security_governance.py 43 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485486487488489490491492493494495496497498499500501502503504505506507508509510511512513514515516517518519520521522523524525526527528529530531532533534535536537538539540541542543544545546547548549550551552553554555556557558559560561562563564565566567568569570571572573574575576577578579580581582583584585586587588589590591592593594595596597598599600601602603604605606607608609610611612613614615616617618619620621622623624625626627628629630631632633634635636637638639640641642643644645646647648649650651652653654655656657658659660661662663664665666667668669670671672673674675676677678679680681682683684685686687688689690691692693694695696697698699700701702703704705706707708709710711712713714715716717718719720721722723724725726727728729730731732733734735736737738739740741742743744745746747748749750751752753754755756757758759760761762763764765766767768769770771772773774775776777778779780781782783784785786787788789790791792793794795796797798799800801802803804805806807808809810811812813814815816817818819820821822823824825826827828829830831832833834835836837838839840841842843844845846847848849850851852853854855856857858859860861862863864865866867868869870871872873874875876877878879880881882
  1. """Cross-domain data security governance and security engineering controls."""
  2. from __future__ import annotations
  3. import copy
  4. import hashlib
  5. import json
  6. import re
  7. import uuid
  8. from collections.abc import Callable
  9. from datetime import UTC, datetime, timedelta
  10. from typing import Any
  11. from urllib.parse import urlsplit
  12. from app.core.common.identifiers import new_governance_uid
  13. from app.core.common.timezone_utils import now_china
  14. CLASSIFICATIONS = ("public", "internal", "sensitive", "highly_sensitive")
  15. CLASSIFICATION_RANK = {value: index for index, value in enumerate(CLASSIFICATIONS)}
  16. FINDING_STATUSES = {"pending_review", "confirmed", "dismissed"}
  17. ACCESS_ACTIONS = {"read", "use"}
  18. ENVIRONMENTS = {"development", "test", "production"}
  19. RETENTION_EVIDENCE_TYPES = {
  20. "classification_evidence", "access_decision", "egress_request",
  21. "audit_event", "siem_delivery", "sbom", "vulnerability",
  22. }
  23. ARCHIVE_MODES = {"hot", "immutable_external"}
  24. DISPOSITION_ACTIONS = {"review", "archive"}
  25. SIEM_CATEGORIES = {
  26. "authentication", "ingestion", "entity_resolution", "publication",
  27. "remediation", "knowledge_query", "authorization", "workflow_task",
  28. "data_product", "agent", "security_governance",
  29. }
  30. SEVERITIES = {"unknown", "low", "medium", "high", "critical"}
  31. RESOLUTION_TYPES = {"patched", "not_affected", "accepted_risk"}
  32. ROLE_PATTERN = re.compile(r"^[a-z][a-z0-9:_-]{1,79}$")
  33. CODE_PATTERN = re.compile(r"^[A-Z][A-Z0-9_]{2,119}$")
  34. FIELD_PATTERN = re.compile(r"^[A-Za-z_][A-Za-z0-9_.-]{0,199}$")
  35. HASH_PATTERN = re.compile(r"^[0-9a-f]{64}$")
  36. PHONE_PATTERN = re.compile(r"(?<!\d)1[3-9]\d{9}(?!\d)")
  37. EMAIL_PATTERN = re.compile(r"(?i)\b[A-Z0-9._%+-]+@[A-Z0-9.-]+\.[A-Z]{2,}\b")
  38. BANK_CARD_PATTERN = re.compile(r"(?<!\d)\d{15,19}(?!\d)")
  39. PRC_ID_PATTERN = re.compile(r"(?<!\d)\d{17}[0-9Xx](?!\d)")
  40. def _closed(value: Any, allowed: set[str], label: str) -> dict[str, Any]:
  41. if not isinstance(value, dict):
  42. raise ValueError(f"{label} must be an object")
  43. unknown = sorted(set(value) - allowed)
  44. if unknown:
  45. raise ValueError(f"{label} contains unsupported fields: {', '.join(unknown)}")
  46. return copy.deepcopy(value)
  47. def _text(value: Any, label: str, maximum: int = 1000) -> str:
  48. if not isinstance(value, str) or not value.strip():
  49. raise ValueError(f"{label} is required")
  50. result = value.strip()
  51. if len(result) > maximum:
  52. raise ValueError(f"{label} exceeds {maximum} characters")
  53. return result
  54. def _optional_text(value: Any, label: str, maximum: int = 1000) -> str | None:
  55. if value in (None, ""):
  56. return None
  57. return _text(value, label, maximum)
  58. def _uid(value: Any, label: str) -> str:
  59. try:
  60. return str(uuid.UUID(str(value)))
  61. except (TypeError, ValueError, AttributeError) as error:
  62. raise ValueError(f"{label} must be a UUID") from error
  63. def _list(value: Any, label: str, minimum: int = 0, maximum: int = 1000) -> list[Any]:
  64. if not isinstance(value, list) or len(value) < minimum or len(value) > maximum:
  65. raise ValueError(f"{label} must contain between {minimum} and {maximum} items")
  66. return copy.deepcopy(value)
  67. def _time(value: Any, label: str) -> datetime:
  68. if isinstance(value, datetime):
  69. result = value
  70. else:
  71. try:
  72. result = datetime.fromisoformat(str(value).replace("Z", "+00:00"))
  73. except (TypeError, ValueError) as error:
  74. raise ValueError(f"{label} must be ISO-8601") from error
  75. if result.tzinfo is None:
  76. raise ValueError(f"{label} must include a timezone")
  77. return result.astimezone(UTC)
  78. def _classification(value: Any, label: str = "classification") -> str:
  79. result = _text(value, label, 40)
  80. if result not in CLASSIFICATION_RANK:
  81. raise ValueError(f"unsupported {label}")
  82. return result
  83. def _canonical(value: Any) -> bytes:
  84. return json.dumps(value, ensure_ascii=False, sort_keys=True, separators=(",", ":")).encode()
  85. def _digest(value: Any) -> str:
  86. return hashlib.sha256(_canonical(value)).hexdigest()
  87. def _mask(value: str, detector: str) -> str:
  88. if detector == "email" and "@" in value:
  89. local, domain = value.split("@", 1)
  90. return f"{local[:1]}***@{domain}"
  91. if detector == "phone" and len(value) >= 7:
  92. return f"{value[:3]}****{value[-4:]}"
  93. if len(value) >= 6:
  94. return f"{value[:2]}****{value[-2:]}"
  95. return "***"
  96. def _normalize_evidence(value: Any) -> list[dict[str, str]]:
  97. result = []
  98. for item in _list(value, "evidence_refs", 1, 50):
  99. body = _closed(item, {"type", "ref", "digest"}, "evidence reference")
  100. digest = _text(body.get("digest"), "evidence digest", 64).lower()
  101. if not HASH_PATTERN.fullmatch(digest):
  102. raise ValueError("evidence digest must be SHA-256")
  103. result.append({
  104. "type": _text(body.get("type"), "evidence type", 60),
  105. "ref": _text(body.get("ref"), "evidence ref", 300),
  106. "digest": digest,
  107. })
  108. return result
  109. class SecurityGovernanceService:
  110. """Own security decisions while leaving data and external security tools authoritative."""
  111. def __init__(
  112. self,
  113. repository,
  114. *,
  115. approval_gateway,
  116. siem_transport,
  117. siem_host_allowlist: set[str] | frozenset[str],
  118. uid_factory: Callable[[], str] = new_governance_uid,
  119. now_factory: Callable[[], datetime] = now_china,
  120. commit: Callable[[], None] = lambda: None,
  121. rollback: Callable[[], None] = lambda: None,
  122. ):
  123. self.repository = repository
  124. self.approval_gateway = approval_gateway
  125. self.siem_transport = siem_transport
  126. self.siem_host_allowlist = {
  127. str(value).strip().lower() for value in siem_host_allowlist if str(value).strip()
  128. }
  129. self.uid_factory = uid_factory
  130. self.now_factory = now_factory
  131. self.commit = commit
  132. self.rollback = rollback
  133. def _actor(self, actor_uid: Any) -> str:
  134. actor = _uid(actor_uid, "actor_uid")
  135. if self.repository.users_available({actor}) != {actor}:
  136. raise ValueError("security actor is unavailable")
  137. return actor
  138. def _save(self, operation):
  139. try:
  140. result = operation()
  141. self.commit()
  142. return result
  143. except Exception:
  144. self.rollback()
  145. raise
  146. def _event(self, resource_type, resource_uid, action, actor_uid, detail=None):
  147. self.repository.add_event(
  148. resource_type, resource_uid, action, actor_uid, copy.deepcopy(detail or {})
  149. )
  150. def create_classification_profile(self, payload: Any, *, actor_uid: str):
  151. body = _closed(
  152. payload,
  153. {"code", "name", "business_domain_uid", "default_classification", "rules"},
  154. "classification profile",
  155. )
  156. actor = self._actor(actor_uid)
  157. code = _text(body.get("code"), "profile code", 120).upper()
  158. if not CODE_PATTERN.fullmatch(code):
  159. raise ValueError("classification profile code is invalid")
  160. rules = []
  161. for raw in _list(body.get("rules"), "classification rules", 1, 100):
  162. rule = _closed(raw, {"field_tokens", "category", "classification"}, "classification rule")
  163. tokens = sorted({
  164. _text(item, "field token", 60).casefold()
  165. for item in _list(rule.get("field_tokens"), "field_tokens", 1, 20)
  166. })
  167. if any(not re.fullmatch(r"[a-z0-9_-]+", item) for item in tokens):
  168. raise ValueError("field tokens must be simple identifiers")
  169. rules.append({
  170. "field_tokens": tokens,
  171. "category": _text(rule.get("category"), "category", 80),
  172. "classification": _classification(rule.get("classification")),
  173. })
  174. now = self.now_factory().isoformat()
  175. record = {
  176. "uid": self.uid_factory(), "code": code,
  177. "name": _text(body.get("name"), "profile name", 300),
  178. "business_domain_uid": _uid(body.get("business_domain_uid"), "business_domain_uid"),
  179. "default_classification": _classification(body.get("default_classification")),
  180. "rules": rules, "status": "active", "current_version": 1,
  181. "created_by": actor, "created_at": now, "updated_at": now,
  182. }
  183. def operation():
  184. result = self.repository.create_profile(record)
  185. self._event("classification_profile", record["uid"], "profile_created", actor, {"code": code})
  186. return result
  187. return self._save(operation)
  188. @staticmethod
  189. def _detect_field(profile: dict[str, Any], field: dict[str, Any]):
  190. field_name = _text(field.get("name"), "field name", 200)
  191. if not FIELD_PATTERN.fullmatch(field_name):
  192. raise ValueError("field name is invalid")
  193. values = [str(value)[:500] for value in _list(field.get("sample_values", []), "sample_values", 0, 20)]
  194. normalized = field_name.casefold().replace("-", "_").replace(".", "_")
  195. matches = []
  196. for rule in profile["rules"]:
  197. if any(token in normalized.split("_") or token in normalized for token in rule["field_tokens"]):
  198. matches.append((rule["classification"], rule["category"], "field_rule", None))
  199. detectors = (
  200. ("prc_id", PRC_ID_PATTERN, "personal_identifier", "highly_sensitive"),
  201. ("bank_card", BANK_CARD_PATTERN, "financial_account", "highly_sensitive"),
  202. ("phone", PHONE_PATTERN, "personal_contact", "sensitive"),
  203. ("email", EMAIL_PATTERN, "personal_contact", "sensitive"),
  204. )
  205. for value in values:
  206. for detector, pattern, category, level in detectors:
  207. found = pattern.search(value)
  208. if found:
  209. sample = found.group(0)
  210. matches.append((level, category, detector, sample))
  211. if not matches:
  212. return None
  213. level = max(matches, key=lambda item: CLASSIFICATION_RANK[item[0]])[0]
  214. categories = sorted({item[1] for item in matches})
  215. detector_codes = sorted({item[2] for item in matches})
  216. raw_matches = [item for item in matches if item[3] is not None]
  217. return {
  218. "field_name": field_name,
  219. "categories": categories,
  220. "proposed_classification": level,
  221. "detector_codes": detector_codes,
  222. "sample_fingerprints": sorted({_digest(item[3]) for item in raw_matches}),
  223. "masked_examples": sorted({_mask(item[3], item[2]) for item in raw_matches})[:3],
  224. }
  225. def scan_sensitive_sample(self, payload: Any, *, actor_uid: str):
  226. body = _closed(
  227. payload,
  228. {"profile_uid", "resource_type", "resource_uid", "business_domain_uid", "fields"},
  229. "sensitive sample scan",
  230. )
  231. actor = self._actor(actor_uid)
  232. profile = self.repository.get_profile(_uid(body.get("profile_uid"), "profile_uid"))
  233. if not profile or profile["status"] != "active":
  234. raise LookupError("active classification profile was not found")
  235. domain_uid = _uid(body.get("business_domain_uid"), "business_domain_uid")
  236. if profile["business_domain_uid"] != domain_uid:
  237. raise PermissionError("classification profile domain does not match")
  238. fields = _list(body.get("fields"), "fields", 1, 100)
  239. now = self.now_factory().isoformat()
  240. scan_uid = self.uid_factory()
  241. findings = []
  242. for field in fields:
  243. normalized = self._detect_field(profile, field)
  244. if not normalized:
  245. continue
  246. findings.append({
  247. "uid": self.uid_factory(), "scan_uid": scan_uid,
  248. **normalized, "status": "pending_review", "final_classification": None,
  249. "review_reason": None, "reviewed_by": None, "reviewed_at": None,
  250. "current_version": 1, "created_at": now,
  251. })
  252. scan = {
  253. "uid": scan_uid, "profile_uid": profile["uid"],
  254. "resource_type": _text(body.get("resource_type"), "resource_type", 80),
  255. "resource_uid": _text(body.get("resource_uid"), "resource_uid", 200),
  256. "business_domain_uid": domain_uid, "field_count": len(fields),
  257. "finding_count": len(findings), "sample_retained": False,
  258. "created_by": actor, "created_at": now,
  259. }
  260. def operation():
  261. result = self.repository.create_scan(scan, findings)
  262. self._event(
  263. "classification_scan", scan_uid, "sample_scanned", actor,
  264. {"finding_count": len(findings), "sample_retained": False},
  265. )
  266. return result
  267. return self._save(operation)
  268. def review_classification_finding(
  269. self, finding_uid: str, payload: Any, *, expected_version: int, actor_uid: str
  270. ):
  271. body = _closed(payload, {"decision", "final_classification", "reason"}, "classification review")
  272. actor = self._actor(actor_uid)
  273. finding = self.repository.get_finding(_uid(finding_uid, "finding_uid"))
  274. if not finding:
  275. raise LookupError("classification finding was not found")
  276. if finding["status"] != "pending_review":
  277. raise RuntimeError("classification finding is not pending review")
  278. scan = self.repository.get_scan(finding["scan_uid"]) if hasattr(self.repository, "get_scan") else None
  279. creator = scan.get("created_by") if scan else None
  280. if actor == creator or (creator is None and actor == finding.get("created_by")):
  281. raise PermissionError("classification requires an independent reviewer")
  282. # Memory repositories keep the creator on the scan rather than the finding.
  283. if (
  284. creator is None
  285. and hasattr(self.repository, "scans")
  286. and actor == self.repository.scans[finding["scan_uid"]]["created_by"]
  287. ):
  288. raise PermissionError("classification requires an independent reviewer")
  289. decision = _text(body.get("decision"), "decision", 20)
  290. if decision not in {"confirm", "dismiss"}:
  291. raise ValueError("unsupported classification review decision")
  292. updated = {
  293. **finding,
  294. "status": "confirmed" if decision == "confirm" else "dismissed",
  295. "final_classification": (
  296. _classification(body.get("final_classification")) if decision == "confirm" else None
  297. ),
  298. "review_reason": _text(body.get("reason"), "review reason", 1000),
  299. "reviewed_by": actor, "reviewed_at": self.now_factory().isoformat(),
  300. }
  301. def operation():
  302. result = self.repository.update_finding(updated, int(expected_version))
  303. self._event("classification_finding", finding["uid"], f"finding_{updated['status']}", actor, {"classification": updated["final_classification"]})
  304. return result
  305. return self._save(operation)
  306. def create_access_policy(self, payload: Any, *, actor_uid: str):
  307. body = _closed(
  308. payload,
  309. {
  310. "code", "name", "business_domain_uid", "subject_user_uids",
  311. "subject_roles", "purposes", "environments", "actions",
  312. "max_classification", "allowed_fields", "expires_at", "review_due_at",
  313. },
  314. "access policy",
  315. )
  316. actor = self._actor(actor_uid)
  317. code = _text(body.get("code"), "policy code", 120).upper()
  318. if not CODE_PATTERN.fullmatch(code):
  319. raise ValueError("access policy code is invalid")
  320. users = sorted({_uid(value, "subject_user_uid") for value in _list(body.get("subject_user_uids", []), "subject_user_uids", 0, 100)})
  321. roles = sorted({_text(value, "subject role", 80) for value in _list(body.get("subject_roles", []), "subject_roles", 0, 50)})
  322. if not users and not roles:
  323. raise ValueError("access policy requires a user or role subject")
  324. if users and self.repository.users_available(users) != set(users):
  325. raise ValueError("access policy contains unavailable users")
  326. if any(not ROLE_PATTERN.fullmatch(role) for role in roles):
  327. raise ValueError("access policy role is invalid")
  328. purposes = sorted({_text(value, "purpose", 100) for value in _list(body.get("purposes"), "purposes", 1, 50)})
  329. environments = sorted({_text(value, "environment", 30) for value in _list(body.get("environments"), "environments", 1, 10)})
  330. actions = sorted({_text(value, "action", 20) for value in _list(body.get("actions"), "actions", 1, 10)})
  331. if not set(environments) <= ENVIRONMENTS or not set(actions) <= ACCESS_ACTIONS:
  332. raise ValueError("unsupported access environment or action")
  333. fields = sorted({_text(value, "allowed field", 200) for value in _list(body.get("allowed_fields"), "allowed_fields", 1, 500)})
  334. if any(value != "*" and not FIELD_PATTERN.fullmatch(value) for value in fields):
  335. raise ValueError("allowed field is invalid")
  336. now = self.now_factory().astimezone(UTC)
  337. expires = _time(body.get("expires_at"), "expires_at")
  338. review = _time(body.get("review_due_at"), "review_due_at")
  339. if review <= now or expires <= review:
  340. raise ValueError("access policy review and expiry dates are invalid")
  341. record = {
  342. "uid": self.uid_factory(), "code": code,
  343. "name": _text(body.get("name"), "policy name", 300),
  344. "business_domain_uid": _uid(body.get("business_domain_uid"), "business_domain_uid"),
  345. "subject_user_uids": users, "subject_roles": roles, "purposes": purposes,
  346. "environments": environments, "actions": actions,
  347. "max_classification": _classification(body.get("max_classification"), "max_classification"),
  348. "allowed_fields": fields, "expires_at": expires.isoformat(),
  349. "review_due_at": review.isoformat(), "status": "active", "current_version": 1,
  350. "created_by": actor, "created_at": now.isoformat(), "updated_at": now.isoformat(),
  351. }
  352. def operation():
  353. result = self.repository.create_access_policy(record)
  354. self._event("access_policy", record["uid"], "access_policy_created", actor, {"code": code})
  355. return result
  356. return self._save(operation)
  357. @staticmethod
  358. def _policy_matches(policy: dict[str, Any], context: dict[str, Any], now: datetime, *, check_fields=True):
  359. subject = (
  360. context["user_uid"] in policy["subject_user_uids"]
  361. or bool(set(context["roles"]) & set(policy["subject_roles"]))
  362. )
  363. if not subject:
  364. return False
  365. if policy["business_domain_uid"] != context["business_domain_uid"]:
  366. return False
  367. if context["purpose"] not in policy["purposes"]:
  368. return False
  369. if context["environment"] not in policy["environments"]:
  370. return False
  371. if context["action"] not in policy["actions"]:
  372. return False
  373. if CLASSIFICATION_RANK[context["classification"]] > CLASSIFICATION_RANK[policy["max_classification"]]:
  374. return False
  375. if _time(policy["expires_at"], "expires_at") <= now or _time(policy["review_due_at"], "review_due_at") <= now:
  376. return False
  377. if check_fields and "*" not in policy["allowed_fields"]:
  378. return set(context["requested_fields"]) <= set(policy["allowed_fields"])
  379. return True
  380. def evaluate_access(self, payload: Any):
  381. body = _closed(
  382. payload,
  383. {
  384. "user_uid", "roles", "business_domain_uid", "purpose", "environment",
  385. "action", "resource_uid", "classification", "requested_fields",
  386. },
  387. "access decision",
  388. )
  389. user_uid = _uid(body.get("user_uid"), "user_uid")
  390. if self.repository.users_available({user_uid}) != {user_uid}:
  391. raise ValueError("access user is unavailable")
  392. roles = sorted({_text(value, "role", 80) for value in _list(body.get("roles"), "roles", 1, 50)})
  393. fields = sorted({_text(value, "requested field", 200) for value in _list(body.get("requested_fields"), "requested_fields", 1, 500)})
  394. context = {
  395. "user_uid": user_uid, "roles": roles,
  396. "business_domain_uid": _uid(body.get("business_domain_uid"), "business_domain_uid"),
  397. "purpose": _text(body.get("purpose"), "purpose", 100),
  398. "environment": _text(body.get("environment"), "environment", 30),
  399. "action": _text(body.get("action"), "action", 20),
  400. "resource_uid": _text(body.get("resource_uid"), "resource_uid", 200),
  401. "classification": _classification(body.get("classification")),
  402. "requested_fields": fields,
  403. }
  404. if context["environment"] not in ENVIRONMENTS or context["action"] not in ACCESS_ACTIONS:
  405. raise ValueError("unsupported access environment or action")
  406. now = self.now_factory().astimezone(UTC)
  407. policies = self.repository.matching_access_policies(**context)
  408. policy = next((item for item in policies if self._policy_matches(item, context, now)), None)
  409. scope_policy = next((item for item in policies if self._policy_matches(item, context, now, check_fields=False)), None)
  410. if policy:
  411. decision, reason = "authorized", "policy_allowed"
  412. elif scope_policy:
  413. decision, reason = "denied", "field_minimization_denied"
  414. else:
  415. decision, reason = "denied", "default_deny"
  416. record = {
  417. "uid": self.uid_factory(), **context,
  418. "policy_uid": policy["uid"] if policy else None,
  419. "decision": decision, "reason_code": reason,
  420. "decided_at": self.now_factory().isoformat(),
  421. }
  422. def operation():
  423. result = self.repository.create_access_decision(record)
  424. self._event("access_decision", record["uid"], f"access_{decision}", user_uid, {"reason_code": reason, "policy_uid": record["policy_uid"]})
  425. return result
  426. return self._save(operation)
  427. def submit_egress_request(self, payload: Any, *, actor_uid: str):
  428. body = _closed(
  429. payload,
  430. {
  431. "business_domain_uid", "resource_uid", "classification", "purpose",
  432. "environment", "requested_fields", "minimized_fields", "masking_applied",
  433. "destination_zone", "expires_at", "workflow_uid",
  434. },
  435. "egress request",
  436. )
  437. actor = self._actor(actor_uid)
  438. requested = sorted({_text(value, "requested field", 200) for value in _list(body.get("requested_fields"), "requested_fields", 1, 500)})
  439. minimized = sorted({_text(value, "minimized field", 200) for value in _list(body.get("minimized_fields"), "minimized_fields", 1, 500)})
  440. if not set(minimized) <= set(requested):
  441. raise ValueError("minimized fields must be a subset of requested fields")
  442. masking = body.get("masking_applied")
  443. if not isinstance(masking, bool):
  444. raise ValueError("masking_applied must be boolean")
  445. classification = _classification(body.get("classification"))
  446. if classification in {"sensitive", "highly_sensitive"} and not masking:
  447. raise ValueError("sensitive egress requires masking")
  448. environment = _text(body.get("environment"), "environment", 30)
  449. if environment not in ENVIRONMENTS:
  450. raise ValueError("unsupported egress environment")
  451. now = self.now_factory().astimezone(UTC)
  452. expires = _time(body.get("expires_at"), "expires_at")
  453. if expires <= now or expires > now + timedelta(days=30):
  454. raise ValueError("egress expiry must be within 30 days")
  455. high = classification == "highly_sensitive"
  456. record = {
  457. "uid": self.uid_factory(),
  458. "business_domain_uid": _uid(body.get("business_domain_uid"), "business_domain_uid"),
  459. "resource_uid": _text(body.get("resource_uid"), "resource_uid", 200),
  460. "classification": classification,
  461. "purpose": _text(body.get("purpose"), "purpose", 300),
  462. "environment": environment, "requested_fields": requested,
  463. "minimized_fields": minimized, "approved_fields": [],
  464. "masking_applied": masking,
  465. "destination_zone": _text(body.get("destination_zone"), "destination_zone", 100),
  466. "expires_at": expires.isoformat(),
  467. "status": "denied" if high else "pending_approval",
  468. "reason_code": "highly_sensitive_egress_disabled" if high else "approval_required",
  469. "approval_task_uid": None, "current_version": 1,
  470. "created_by": actor, "created_at": now.isoformat(), "updated_at": now.isoformat(),
  471. }
  472. if not high:
  473. workflow_uid = _uid(body.get("workflow_uid"), "workflow_uid")
  474. task = self.approval_gateway.create_egress_task(record, workflow_uid, actor)
  475. record["approval_task_uid"] = task["uid"]
  476. def operation():
  477. result = self.repository.create_egress_request(record)
  478. self._event("egress_request", record["uid"], f"egress_{record['status']}", actor, {"classification": classification, "reason_code": record["reason_code"]})
  479. return result
  480. return self._save(operation)
  481. def reconcile_egress_request(self, request_uid: str, *, expected_version: int, actor_uid: str):
  482. actor = self._actor(actor_uid)
  483. record = self.repository.get_egress_request(_uid(request_uid, "request_uid"))
  484. if not record:
  485. raise LookupError("egress request was not found")
  486. if record["status"] != "pending_approval":
  487. raise RuntimeError("egress request is not pending approval")
  488. task = self.approval_gateway.get_task(record["approval_task_uid"])
  489. if not task or task["status"] not in {"approved", "rejected"}:
  490. raise RuntimeError("egress approval has no final decision")
  491. approved = task["status"] == "approved"
  492. updated = {
  493. **record,
  494. "status": "authorized_until_expiry" if approved else "denied",
  495. "reason_code": "approval_granted" if approved else "approval_rejected",
  496. "approved_fields": record["minimized_fields"] if approved else [],
  497. "updated_at": self.now_factory().isoformat(),
  498. }
  499. def operation():
  500. result = self.repository.update_egress_request(updated, int(expected_version))
  501. self._event("egress_request", record["uid"], f"egress_{updated['status']}", actor, {"approved_field_count": len(updated["approved_fields"])})
  502. return result
  503. return self._save(operation)
  504. def create_retention_policy(self, payload: Any, *, actor_uid: str):
  505. body = _closed(
  506. payload,
  507. {"code", "name", "evidence_type", "retention_days", "archive_mode", "disposition_action"},
  508. "retention policy",
  509. )
  510. actor = self._actor(actor_uid)
  511. code = _text(body.get("code"), "retention code", 120).upper()
  512. if not CODE_PATTERN.fullmatch(code):
  513. raise ValueError("retention code is invalid")
  514. evidence_type = _text(body.get("evidence_type"), "evidence_type", 60)
  515. archive_mode = _text(body.get("archive_mode"), "archive_mode", 40)
  516. disposition = _text(body.get("disposition_action"), "disposition_action", 40)
  517. if evidence_type not in RETENTION_EVIDENCE_TYPES or archive_mode not in ARCHIVE_MODES or disposition not in DISPOSITION_ACTIONS:
  518. raise ValueError("unsupported retention policy option")
  519. try:
  520. days = int(body.get("retention_days"))
  521. except (TypeError, ValueError) as error:
  522. raise ValueError("retention_days must be an integer") from error
  523. if days < 30 or days > 36500:
  524. raise ValueError("retention_days must be between 30 and 36500")
  525. now = self.now_factory().isoformat()
  526. record = {
  527. "uid": self.uid_factory(), "code": code,
  528. "name": _text(body.get("name"), "retention name", 300),
  529. "evidence_type": evidence_type, "retention_days": days,
  530. "archive_mode": archive_mode, "disposition_action": disposition,
  531. "automatic_deletion": False, "status": "active", "current_version": 1,
  532. "created_by": actor, "created_at": now, "updated_at": now,
  533. }
  534. def operation():
  535. result = self.repository.create_retention_policy(record)
  536. self._event("retention_policy", record["uid"], "retention_policy_created", actor, {"evidence_type": evidence_type, "automatic_deletion": False})
  537. return result
  538. return self._save(operation)
  539. def retention_candidates(self, *, as_of: datetime, limit: int):
  540. instant = _time(as_of, "as_of")
  541. size = int(limit)
  542. if size < 1 or size > 1000:
  543. raise ValueError("retention candidate limit must be between 1 and 1000")
  544. return self.repository.retention_candidates(instant, size)
  545. def create_siem_sink(self, payload: Any, *, actor_uid: str):
  546. body = _closed(payload, {"name", "sink_type", "endpoint", "categories"}, "SIEM sink")
  547. actor = self._actor(actor_uid)
  548. sink_type = _text(body.get("sink_type"), "sink_type", 30)
  549. endpoint = _text(body.get("endpoint"), "endpoint", 500)
  550. parsed = urlsplit(endpoint)
  551. required_scheme = "https" if sink_type == "webhook" else "tls"
  552. if sink_type not in {"webhook", "syslog_tls"} or parsed.scheme != required_scheme:
  553. raise ValueError("SIEM sink requires HTTPS webhook or TLS syslog")
  554. if parsed.username or parsed.password or parsed.query or parsed.fragment:
  555. raise ValueError("SIEM endpoint must not contain credentials, query or fragment")
  556. host = str(parsed.hostname or "").lower()
  557. if not host or host not in self.siem_host_allowlist:
  558. raise ValueError("SIEM endpoint host is outside the allowlist")
  559. categories = sorted({_text(value, "SIEM category", 40) for value in _list(body.get("categories"), "categories", 1, 20)})
  560. if not set(categories) <= SIEM_CATEGORIES:
  561. raise ValueError("unsupported SIEM audit category")
  562. now = self.now_factory().isoformat()
  563. record = {
  564. "uid": self.uid_factory(), "name": _text(body.get("name"), "sink name", 300),
  565. "sink_type": sink_type, "endpoint": endpoint, "endpoint_host": host,
  566. "categories": categories, "status": "active", "current_version": 1,
  567. "created_by": actor, "created_at": now, "updated_at": now,
  568. }
  569. def operation():
  570. result = self.repository.create_siem_sink(record)
  571. self._event("siem_sink", record["uid"], "siem_sink_created", actor, {"sink_type": sink_type, "endpoint_host": host})
  572. return result
  573. return self._save(operation)
  574. def dispatch_siem_events(self, sink_uid: str, payload: Any, *, actor_uid: str):
  575. body = _closed(payload, {"period_start", "period_end", "limit"}, "SIEM delivery")
  576. actor = self._actor(actor_uid)
  577. sink = self.repository.get_siem_sink(_uid(sink_uid, "sink_uid"))
  578. if not sink or sink["status"] != "active":
  579. raise LookupError("active SIEM sink was not found")
  580. start = _time(body.get("period_start"), "period_start")
  581. end = _time(body.get("period_end"), "period_end")
  582. limit = int(body.get("limit", 100))
  583. if end <= start or limit < 1 or limit > 1000:
  584. raise ValueError("SIEM delivery window or limit is invalid")
  585. events = self.repository.fetch_siem_events(
  586. categories=sink["categories"], period_start=start, period_end=end, limit=limit
  587. )
  588. envelope = {
  589. "schema": "dataops.security.audit.v1",
  590. "period_start": start.isoformat(), "period_end": end.isoformat(),
  591. "events": events,
  592. }
  593. result = self.siem_transport.deliver(sink, envelope)
  594. status = result.get("status")
  595. if status not in {"delivered", "failed"}:
  596. raise RuntimeError("SIEM transport returned an invalid status")
  597. now = self.now_factory().isoformat()
  598. record = {
  599. "uid": self.uid_factory(), "sink_uid": sink["uid"],
  600. "period_start": start.isoformat(), "period_end": end.isoformat(),
  601. "event_count": len(events), "payload_digest": _digest(envelope),
  602. "status": status, "remote_ref": _optional_text(result.get("remote_ref"), "remote_ref", 300),
  603. "error_code": _optional_text(result.get("error_code"), "error_code", 80),
  604. "created_by": actor, "created_at": now,
  605. }
  606. def operation():
  607. saved = self.repository.create_siem_delivery(record)
  608. self._event("siem_delivery", saved["uid"], f"siem_{status}", actor, {"event_count": len(events), "payload_digest": record["payload_digest"]})
  609. return saved
  610. return self._save(operation)
  611. def register_sbom(self, payload: Any, *, actor_uid: str):
  612. body = _closed(
  613. payload,
  614. {"artifact_name", "artifact_version", "artifact_type", "source_ref", "document"},
  615. "SBOM registration",
  616. )
  617. actor = self._actor(actor_uid)
  618. document = body.get("document")
  619. if not isinstance(document, dict):
  620. raise ValueError("SBOM document must be an object")
  621. if document.get("bomFormat") != "CycloneDX" or document.get("specVersion") != "1.5":
  622. raise ValueError("SBOM must use CycloneDX 1.5")
  623. components = []
  624. for raw in _list(document.get("components", []), "SBOM components", 0, 10000):
  625. component = _closed(
  626. raw,
  627. {
  628. "type", "name", "version", "purl", "bom-ref", "licenses",
  629. "externalReferences", "properties", "group", "supplier", "publisher",
  630. "author", "description", "hashes", "scope", "copyright",
  631. },
  632. "SBOM component",
  633. )
  634. components.append({
  635. "type": _text(component.get("type"), "component type", 50),
  636. "name": _text(component.get("name"), "component name", 300),
  637. "version": _optional_text(component.get("version"), "component version", 200),
  638. "purl": _optional_text(component.get("purl"), "component purl", 500),
  639. })
  640. now = self.now_factory().isoformat()
  641. record = {
  642. "uid": self.uid_factory(),
  643. "artifact_name": _text(body.get("artifact_name"), "artifact name", 300),
  644. "artifact_version": _text(body.get("artifact_version"), "artifact version", 120),
  645. "artifact_type": _text(body.get("artifact_type"), "artifact type", 50),
  646. "source_ref": _text(body.get("source_ref"), "source_ref", 500),
  647. "format": "CycloneDX", "spec_version": "1.5",
  648. "document_digest": _digest(document), "component_count": len(components),
  649. "components": sorted(components, key=lambda item: (item["name"], item.get("version") or "")),
  650. "created_by": actor, "created_at": now,
  651. }
  652. def operation():
  653. result = self.repository.create_sbom(record)
  654. self._event("sbom", record["uid"], "sbom_registered", actor, {"artifact_name": record["artifact_name"], "component_count": len(components), "document_digest": record["document_digest"]})
  655. return result
  656. return self._save(operation)
  657. def ingest_vulnerabilities(self, sbom_uid: str, payload: Any, *, actor_uid: str):
  658. body = _closed(payload, {"scanner", "scan_ref", "findings"}, "vulnerability import")
  659. actor = self._actor(actor_uid)
  660. sbom = self.repository.get_sbom(_uid(sbom_uid, "sbom_uid"))
  661. if not sbom:
  662. raise LookupError("SBOM was not found")
  663. scanner = _text(body.get("scanner"), "scanner", 100)
  664. scan_ref = _text(body.get("scan_ref"), "scan_ref", 500)
  665. now = self.now_factory().isoformat()
  666. records = []
  667. for raw in _list(body.get("findings"), "vulnerability findings", 1, 5000):
  668. finding = _closed(
  669. raw,
  670. {"external_id", "severity", "component_name", "installed_version", "fixed_version", "title"},
  671. "vulnerability finding",
  672. )
  673. severity = _text(finding.get("severity"), "severity", 20).lower()
  674. if severity not in SEVERITIES:
  675. raise ValueError("unsupported vulnerability severity")
  676. records.append({
  677. "uid": self.uid_factory(), "sbom_uid": sbom["uid"],
  678. "scanner": scanner, "scan_ref": scan_ref,
  679. "external_id": _text(finding.get("external_id"), "external_id", 120),
  680. "severity": severity,
  681. "component_name": _text(finding.get("component_name"), "component_name", 300),
  682. "installed_version": _text(finding.get("installed_version"), "installed_version", 200),
  683. "fixed_version": _optional_text(finding.get("fixed_version"), "fixed_version", 200),
  684. "title": _text(finding.get("title"), "title", 500),
  685. "status": "open", "assignee_uid": None, "due_at": None,
  686. "resolution_type": None, "resolved_version": None, "resolution": None,
  687. "evidence_refs": [], "resolved_by": None, "resolved_at": None,
  688. "closed_by": None, "closed_at": None, "close_reason": None,
  689. "current_version": 1, "created_by": actor, "created_at": now, "updated_at": now,
  690. })
  691. def operation():
  692. result = self.repository.upsert_vulnerabilities(sbom["uid"], records)
  693. self._event("sbom", sbom["uid"], "vulnerabilities_imported", actor, {"scanner": scanner, "finding_count": len(result)})
  694. return result
  695. return self._save(operation)
  696. def assign_vulnerability(
  697. self, finding_uid: str, payload: Any, *, expected_version: int, actor_uid: str
  698. ):
  699. body = _closed(payload, {"assignee_uid", "due_at"}, "vulnerability assignment")
  700. actor = self._actor(actor_uid)
  701. finding = self.repository.get_vulnerability(_uid(finding_uid, "finding_uid"))
  702. if not finding:
  703. raise LookupError("vulnerability was not found")
  704. if finding["status"] not in {"open", "triaged", "in_progress"}:
  705. raise RuntimeError("vulnerability cannot be assigned")
  706. assignee = _uid(body.get("assignee_uid"), "assignee_uid")
  707. if self.repository.users_available({assignee}) != {assignee}:
  708. raise ValueError("vulnerability assignee is unavailable")
  709. due_at = _time(body.get("due_at"), "due_at")
  710. if due_at <= self.now_factory().astimezone(UTC):
  711. raise ValueError("vulnerability due date must be in the future")
  712. updated = {
  713. **finding, "status": "in_progress", "assignee_uid": assignee,
  714. "due_at": due_at.isoformat(), "updated_at": self.now_factory().isoformat(),
  715. }
  716. return self._update_vulnerability(updated, expected_version, "vulnerability_assigned", actor, {"assignee_uid": assignee})
  717. def resolve_vulnerability(
  718. self, finding_uid: str, payload: Any, *, expected_version: int, actor_uid: str
  719. ):
  720. body = _closed(
  721. payload,
  722. {"resolution_type", "resolved_version", "resolution", "evidence_refs"},
  723. "vulnerability resolution",
  724. )
  725. actor = self._actor(actor_uid)
  726. finding = self.repository.get_vulnerability(_uid(finding_uid, "finding_uid"))
  727. if not finding:
  728. raise LookupError("vulnerability was not found")
  729. if finding["status"] != "in_progress" or actor != finding["assignee_uid"]:
  730. raise PermissionError("only the assigned owner can resolve an in-progress vulnerability")
  731. resolution_type = _text(body.get("resolution_type"), "resolution_type", 30)
  732. if resolution_type not in RESOLUTION_TYPES:
  733. raise ValueError("unsupported vulnerability resolution type")
  734. resolved_version = _optional_text(body.get("resolved_version"), "resolved_version", 200)
  735. if resolution_type == "patched" and not resolved_version:
  736. raise ValueError("patched vulnerabilities require a resolved version")
  737. updated = {
  738. **finding, "status": "resolved", "resolution_type": resolution_type,
  739. "resolved_version": resolved_version,
  740. "resolution": _text(body.get("resolution"), "resolution", 2000),
  741. "evidence_refs": _normalize_evidence(body.get("evidence_refs")),
  742. "resolved_by": actor, "resolved_at": self.now_factory().isoformat(),
  743. "updated_at": self.now_factory().isoformat(),
  744. }
  745. return self._update_vulnerability(updated, expected_version, "vulnerability_resolved", actor, {"resolution_type": resolution_type})
  746. def close_vulnerability(
  747. self, finding_uid: str, payload: Any, *, expected_version: int, actor_uid: str
  748. ):
  749. body = _closed(payload, {"reason"}, "vulnerability closure")
  750. actor = self._actor(actor_uid)
  751. finding = self.repository.get_vulnerability(_uid(finding_uid, "finding_uid"))
  752. if not finding:
  753. raise LookupError("vulnerability was not found")
  754. if finding["status"] != "resolved":
  755. raise RuntimeError("only resolved vulnerabilities can be closed")
  756. if actor in {finding.get("assignee_uid"), finding.get("resolved_by")}:
  757. raise PermissionError("vulnerability closure requires an independent reviewer")
  758. updated = {
  759. **finding, "status": "closed", "closed_by": actor,
  760. "closed_at": self.now_factory().isoformat(),
  761. "close_reason": _text(body.get("reason"), "close reason", 1000),
  762. "updated_at": self.now_factory().isoformat(),
  763. }
  764. return self._update_vulnerability(updated, expected_version, "vulnerability_closed", actor, {"resolution_type": finding["resolution_type"]})
  765. def _update_vulnerability(self, record, expected_version, action, actor, detail):
  766. def operation():
  767. result = self.repository.update_vulnerability(record, int(expected_version), action, actor)
  768. self._event("vulnerability", record["uid"], action, actor, detail)
  769. return result
  770. return self._save(operation)
  771. def list_profiles(self, **filters):
  772. return self.repository.list_profiles(**filters)
  773. def list_classification_scans(self, **filters):
  774. return self.repository.list_scans(**filters)
  775. def list_classification_findings(self, **filters):
  776. return self.repository.list_findings(**filters)
  777. def list_access_policies(self, **filters):
  778. return self.repository.list_access_policies(**filters)
  779. def list_access_decisions(self, **filters):
  780. return self.repository.list_access_decisions(**filters)
  781. def list_egress_requests(self, **filters):
  782. return self.repository.list_egress_requests(**filters)
  783. def list_retention_policies(self):
  784. return self.repository.list_retention_policies()
  785. def list_siem_sinks(self):
  786. return self.repository.list_siem_sinks()
  787. def list_siem_deliveries(self, **filters):
  788. return self.repository.list_siem_deliveries(**filters)
  789. def list_sboms(self):
  790. return self.repository.list_sboms()
  791. def list_vulnerabilities(self, **filters):
  792. return self.repository.list_vulnerabilities(**filters)
  793. def dashboard(self):
  794. return self.repository.dashboard()