unified_responsibilities.py 28 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485486487488489490491492493494495496497498499500501502503504505506507508509510511512513514515516517518519520521522523524525526527528529530531532533534535536537538539540541542543544545546547548549550551552553554555556557558559560561562563564565566567568569570571572573574575576577578579580581582583584585586587588589590591592593594595596597598599600601602603604605606607608609610611612613614615616617618619620621622623624625626627628629630631632633634635636637638639640641642643644645646647648649650651652653654655656657658659660661662663664665666667668669670671672673674675676677678679680681682683684685686687688689690691692693694695696697698699700701702703704705706707708709710711712713714715716717718719720721722723724725726727728729730731732733734
  1. """Hierarchical responsibility resolution, delegation and central policies."""
  2. from __future__ import annotations
  3. import copy
  4. import re
  5. import uuid
  6. from collections.abc import Callable
  7. from datetime import datetime
  8. from typing import Any
  9. from app.core.common.identifiers import new_governance_uid
  10. from app.core.common.timezone_utils import now_china
  11. from app.core.governance.responsibilities import (
  12. RACI_ROLES,
  13. RESPONSIBILITY_ROLES,
  14. )
  15. HIERARCHICAL_RESOURCE_TYPES = frozenset(
  16. {
  17. "organization",
  18. "business_domain",
  19. "data_asset",
  20. "semantic_term",
  21. "data_standard",
  22. "quality_policy",
  23. "data_product",
  24. "agent",
  25. }
  26. )
  27. POLICY_TYPES = frozenset({"central_policy", "joint_review"})
  28. DELEGATION_TYPES = frozenset({"temporary", "departure_transfer"})
  29. DECISIONS = frozenset({"approve", "reject"})
  30. CODE_PATTERN = re.compile(r"^[A-Z][A-Z0-9_]{2,119}$")
  31. def _closed(value: Any, allowed: set[str], label: str) -> dict[str, Any]:
  32. if not isinstance(value, dict):
  33. raise ValueError(f"{label} must be an object")
  34. unknown = sorted(set(value) - allowed)
  35. if unknown:
  36. raise ValueError(
  37. f"{label} contains unsupported fields: {', '.join(unknown)}"
  38. )
  39. return copy.deepcopy(value)
  40. def _string(value: Any, label: str, maximum: int = 500) -> str:
  41. if not isinstance(value, str) or not value.strip():
  42. raise ValueError(f"{label} is required")
  43. result = value.strip()
  44. if len(result) > maximum:
  45. raise ValueError(f"{label} exceeds {maximum} characters")
  46. return result
  47. def _uid(value: Any, label: str) -> str:
  48. try:
  49. return str(uuid.UUID(str(value)))
  50. except (TypeError, ValueError, AttributeError) as error:
  51. raise ValueError(f"{label} must be a UUID") from error
  52. def _time(value: Any, label: str) -> datetime:
  53. if isinstance(value, datetime):
  54. result = value
  55. else:
  56. try:
  57. result = datetime.fromisoformat(str(value).replace("Z", "+00:00"))
  58. except (TypeError, ValueError) as error:
  59. raise ValueError(f"{label} must be an ISO-8601 datetime") from error
  60. if result.tzinfo is None:
  61. raise ValueError(f"{label} must include a timezone")
  62. return result
  63. def _resource_type(value: Any, label: str = "resource_type") -> str:
  64. result = _string(value, label, 40)
  65. if result not in HIERARCHICAL_RESOURCE_TYPES:
  66. raise ValueError(f"unsupported {label}")
  67. return result
  68. def _policy_definition(policy_type: str, value: Any) -> dict[str, Any]:
  69. if policy_type == "central_policy":
  70. definition = _closed(
  71. value,
  72. {
  73. "required_responsibility_roles",
  74. "require_unique_accountable",
  75. "max_delegation_days",
  76. },
  77. "central policy definition",
  78. )
  79. roles = definition.get("required_responsibility_roles")
  80. if not isinstance(roles, list) or not roles:
  81. raise ValueError("central policy required roles are missing")
  82. normalized_roles = sorted({_string(item, "required role", 40) for item in roles})
  83. if not set(normalized_roles) <= RESPONSIBILITY_ROLES:
  84. raise ValueError("central policy contains unsupported role")
  85. if not isinstance(definition.get("require_unique_accountable"), bool):
  86. raise ValueError("require_unique_accountable must be a boolean")
  87. try:
  88. days = int(definition.get("max_delegation_days"))
  89. except (TypeError, ValueError) as error:
  90. raise ValueError("max_delegation_days is invalid") from error
  91. if days < 1 or days > 366:
  92. raise ValueError("max_delegation_days must be between 1 and 366")
  93. return {
  94. "required_responsibility_roles": normalized_roles,
  95. "require_unique_accountable": definition[
  96. "require_unique_accountable"
  97. ],
  98. "max_delegation_days": days,
  99. }
  100. definition = _closed(
  101. value,
  102. {"business_domain_uids", "min_approvals", "require_all_domains"},
  103. "joint review definition",
  104. )
  105. domains = definition.get("business_domain_uids")
  106. if not isinstance(domains, list) or len(domains) < 2:
  107. raise ValueError("joint review requires at least two business domains")
  108. normalized_domains = sorted(
  109. {_string(item, "business domain uid", 120) for item in domains}
  110. )
  111. try:
  112. minimum = int(definition.get("min_approvals"))
  113. except (TypeError, ValueError) as error:
  114. raise ValueError("min_approvals is invalid") from error
  115. if minimum < 1 or minimum > len(normalized_domains):
  116. raise ValueError("min_approvals is outside the domain range")
  117. if not isinstance(definition.get("require_all_domains"), bool):
  118. raise ValueError("require_all_domains must be a boolean")
  119. return {
  120. "business_domain_uids": normalized_domains,
  121. "min_approvals": minimum,
  122. "require_all_domains": definition["require_all_domains"],
  123. }
  124. class UnifiedResponsibilityService:
  125. """Resolve effective accountability without granting platform access."""
  126. def __init__(
  127. self,
  128. repository,
  129. *,
  130. uid_factory: Callable[[], str] = new_governance_uid,
  131. now_factory: Callable[[], datetime] = now_china,
  132. commit: Callable[[], None] = lambda: None,
  133. rollback: Callable[[], None] = lambda: None,
  134. ):
  135. self.repository = repository
  136. self.uid_factory = uid_factory
  137. self.now_factory = now_factory
  138. self.commit = commit
  139. self.rollback = rollback
  140. def _chain(self, resource_type: str, resource_uid: str) -> list[dict[str, Any]]:
  141. kind = _resource_type(resource_type)
  142. uid = _string(resource_uid, "resource_uid", 120)
  143. chain = []
  144. seen = set()
  145. for depth in range(20):
  146. identity = (kind, uid)
  147. if identity in seen:
  148. raise ValueError("responsibility hierarchy contains a cycle")
  149. seen.add(identity)
  150. parent = self.repository.parent(kind, uid)
  151. chain.append(
  152. {
  153. "resource_type": kind,
  154. "resource_uid": uid,
  155. "depth": depth,
  156. "hierarchy_revision": int(parent["revision"]) if parent else 0,
  157. }
  158. )
  159. if parent is None:
  160. return chain
  161. kind = _resource_type(parent["parent_type"], "parent_type")
  162. uid = _string(parent["parent_uid"], "parent_uid", 120)
  163. raise ValueError("responsibility hierarchy exceeds 20 levels")
  164. def resolve(
  165. self,
  166. resource_type: str,
  167. resource_uid: str,
  168. *,
  169. at: datetime | str | None = None,
  170. ) -> dict[str, Any]:
  171. evaluated_at = _time(at, "at") if at is not None else self.now_factory()
  172. chain = self._chain(resource_type, resource_uid)
  173. effective_by_raci: dict[str, list[dict[str, Any]]] = {}
  174. for node in chain:
  175. matrix = self.repository.get(
  176. node["resource_type"], node["resource_uid"]
  177. )
  178. assignments = matrix.get("assignments") or []
  179. for raci_role in RACI_ROLES:
  180. if raci_role in effective_by_raci:
  181. continue
  182. matching = [
  183. item for item in assignments if item["raci_role"] == raci_role
  184. ]
  185. if matching:
  186. effective_by_raci[raci_role] = [
  187. {
  188. **copy.deepcopy(item),
  189. "assigned_user_id": item["user_id"],
  190. "source_resource_type": node["resource_type"],
  191. "source_resource_uid": node["resource_uid"],
  192. "source_revision": matrix.get("revision", 0),
  193. "inheritance_depth": node["depth"],
  194. "inherited": node["depth"] > 0,
  195. }
  196. for item in matching
  197. ]
  198. assignments = []
  199. for raci_role in sorted(effective_by_raci):
  200. for assignment in effective_by_raci[raci_role]:
  201. delegation = self.repository.active_delegation(
  202. assignment["assigned_user_id"],
  203. assignment["responsibility_role"],
  204. chain,
  205. evaluated_at,
  206. )
  207. effective_user = (
  208. delegation["delegate_user_uid"]
  209. if delegation
  210. else assignment["assigned_user_id"]
  211. )
  212. active = effective_user in self.repository.users_available(
  213. [effective_user]
  214. )
  215. assignments.append(
  216. {
  217. **assignment,
  218. "effective_user_id": effective_user,
  219. "delegation_uid": delegation["uid"] if delegation else None,
  220. "delegation_type": (
  221. delegation["delegation_type"] if delegation else None
  222. ),
  223. "effective_user_active": active,
  224. }
  225. )
  226. final_owners = [
  227. item
  228. for item in assignments
  229. if item["raci_role"] == "accountable"
  230. and item["effective_user_active"]
  231. ]
  232. policies = self.repository.policies_for_chain(chain)
  233. central = [
  234. item for item in policies if item["policy_type"] == "central_policy"
  235. ]
  236. required_roles = sorted(
  237. {
  238. role
  239. for policy in central
  240. for role in policy["definition"]["required_responsibility_roles"]
  241. }
  242. )
  243. present_roles = {item["responsibility_role"] for item in assignments}
  244. missing_roles = sorted(set(required_roles) - present_roles)
  245. unique_required = any(
  246. item["definition"]["require_unique_accountable"] for item in central
  247. )
  248. compliant = not missing_roles and (
  249. not unique_required or len(final_owners) == 1
  250. )
  251. joint_policies = [
  252. item for item in policies if item["policy_type"] == "joint_review"
  253. ]
  254. joint = None
  255. if joint_policies:
  256. definition = joint_policies[0]["definition"]
  257. joint = {
  258. "policy_uid": joint_policies[0]["uid"],
  259. "required_domains": definition["business_domain_uids"],
  260. "min_approvals": definition["min_approvals"],
  261. "require_all_domains": definition["require_all_domains"],
  262. }
  263. status = (
  264. "resolved"
  265. if len(final_owners) == 1
  266. else "unresolved"
  267. if not final_owners
  268. else "ambiguous"
  269. )
  270. return {
  271. "resource_type": chain[0]["resource_type"],
  272. "resource_uid": chain[0]["resource_uid"],
  273. "status": status,
  274. "evaluated_at": evaluated_at.isoformat(),
  275. "chain": chain,
  276. "assignments": assignments,
  277. "final_owners": final_owners,
  278. "policy_compliance": {
  279. "status": "compliant" if compliant else "non_compliant",
  280. "required_roles": required_roles,
  281. "missing_roles": missing_roles,
  282. "unique_accountable_required": unique_required,
  283. },
  284. "joint_review": joint,
  285. "grants_data_access": False,
  286. }
  287. def set_parent(
  288. self,
  289. payload: Any,
  290. *,
  291. expected_revision: int,
  292. actor_uid: str,
  293. ) -> dict[str, Any]:
  294. body = _closed(
  295. payload,
  296. {"resource_type", "resource_uid", "parent_type", "parent_uid"},
  297. "responsibility hierarchy",
  298. )
  299. resource_type = _resource_type(body.get("resource_type"))
  300. resource_uid = _string(body.get("resource_uid"), "resource_uid", 120)
  301. parent_type = _resource_type(body.get("parent_type"), "parent_type")
  302. parent_uid = _string(body.get("parent_uid"), "parent_uid", 120)
  303. if (resource_type, resource_uid) == (parent_type, parent_uid):
  304. raise ValueError("responsibility hierarchy contains a cycle")
  305. parent_chain = self._chain(parent_type, parent_uid)
  306. if any(
  307. (item["resource_type"], item["resource_uid"])
  308. == (resource_type, resource_uid)
  309. for item in parent_chain
  310. ):
  311. raise ValueError("responsibility hierarchy contains a cycle")
  312. try:
  313. revision = int(expected_revision)
  314. except (TypeError, ValueError) as error:
  315. raise ValueError("hierarchy revision is invalid") from error
  316. now = self.now_factory().isoformat()
  317. record = {
  318. "resource_type": resource_type,
  319. "resource_uid": resource_uid,
  320. "parent_type": parent_type,
  321. "parent_uid": parent_uid,
  322. "updated_by": _uid(actor_uid, "actor_uid"),
  323. "updated_at": now,
  324. }
  325. try:
  326. result = self.repository.set_parent(record, revision)
  327. self.commit()
  328. return result
  329. except Exception:
  330. self.rollback()
  331. raise
  332. def create_delegation(
  333. self, payload: Any, *, actor_uid: str
  334. ) -> dict[str, Any]:
  335. body = _closed(
  336. payload,
  337. {
  338. "source_user_uid",
  339. "delegate_user_uid",
  340. "scope_type",
  341. "scope_uid",
  342. "responsibility_role",
  343. "delegation_type",
  344. "starts_at",
  345. "ends_at",
  346. "reason",
  347. },
  348. "responsibility delegation",
  349. )
  350. source = _uid(body.get("source_user_uid"), "source_user_uid")
  351. delegate = _uid(body.get("delegate_user_uid"), "delegate_user_uid")
  352. if source == delegate:
  353. raise ValueError("delegation users must be different")
  354. delegation_type = _string(
  355. body.get("delegation_type"), "delegation_type", 30
  356. )
  357. if delegation_type not in DELEGATION_TYPES:
  358. raise ValueError("unsupported delegation type")
  359. scope_type = body.get("scope_type")
  360. scope_uid = body.get("scope_uid")
  361. if (scope_type is None) != (scope_uid is None):
  362. raise ValueError("delegation scope type and uid must be paired")
  363. if scope_type is not None:
  364. scope_type = _resource_type(scope_type, "scope_type")
  365. scope_uid = _string(scope_uid, "scope_uid", 120)
  366. role = body.get("responsibility_role")
  367. if role is not None:
  368. role = _string(role, "responsibility_role", 40)
  369. if role not in RESPONSIBILITY_ROLES:
  370. raise ValueError("unsupported responsibility role")
  371. starts_at = _time(body.get("starts_at"), "starts_at")
  372. ends_at = (
  373. _time(body.get("ends_at"), "ends_at")
  374. if body.get("ends_at") is not None
  375. else None
  376. )
  377. if delegation_type == "temporary":
  378. if ends_at is None or ends_at <= starts_at:
  379. raise ValueError("temporary delegation requires a later end time")
  380. if (ends_at - starts_at).days > 366:
  381. raise ValueError("temporary delegation exceeds 366 days")
  382. if scope_type is not None:
  383. policies = self.repository.policies_for_chain(
  384. self._chain(scope_type, scope_uid)
  385. )
  386. limits = [
  387. int(item["definition"]["max_delegation_days"])
  388. for item in policies
  389. if item["policy_type"] == "central_policy"
  390. ]
  391. if limits and (ends_at - starts_at).total_seconds() > min(limits) * 86400:
  392. raise ValueError("temporary delegation exceeds central policy limit")
  393. elif ends_at is not None:
  394. raise ValueError("departure transfer cannot have an end time")
  395. required_users = [delegate] + ([source] if delegation_type == "temporary" else [])
  396. if self.repository.users_available(required_users) != set(required_users):
  397. raise ValueError("delegation user is unknown or disabled")
  398. now = self.now_factory().isoformat()
  399. record = {
  400. "uid": self.uid_factory(),
  401. "source_user_uid": source,
  402. "delegate_user_uid": delegate,
  403. "scope_type": scope_type,
  404. "scope_uid": scope_uid,
  405. "responsibility_role": role,
  406. "delegation_type": delegation_type,
  407. "starts_at": starts_at.isoformat(),
  408. "ends_at": ends_at.isoformat() if ends_at else None,
  409. "reason": _string(body.get("reason"), "reason", 500),
  410. "status": "active",
  411. "current_version": 1,
  412. "created_by": _uid(actor_uid, "actor_uid"),
  413. "created_at": now,
  414. "updated_by": _uid(actor_uid, "actor_uid"),
  415. "updated_at": now,
  416. }
  417. try:
  418. result = self.repository.create_delegation(record)
  419. self.commit()
  420. return result
  421. except Exception:
  422. self.rollback()
  423. raise
  424. def transfer_departing_user(
  425. self, payload: Any, *, actor_uid: str
  426. ) -> dict[str, Any]:
  427. body = _closed(
  428. payload,
  429. {
  430. "source_user_uid",
  431. "delegate_user_uid",
  432. "scope_type",
  433. "scope_uid",
  434. "responsibility_role",
  435. "reason",
  436. },
  437. "departure transfer",
  438. )
  439. return self.create_delegation(
  440. {
  441. **body,
  442. "responsibility_role": body.get("responsibility_role"),
  443. "delegation_type": "departure_transfer",
  444. "starts_at": self.now_factory().isoformat(),
  445. "ends_at": None,
  446. },
  447. actor_uid=actor_uid,
  448. )
  449. def revoke_delegation(
  450. self, uid: str, *, expected_version: int, actor_uid: str
  451. ) -> dict[str, Any]:
  452. record = self.repository.get_delegation(_uid(uid, "delegation_uid"))
  453. if record is None:
  454. raise LookupError("delegation was not found")
  455. if record["status"] != "active":
  456. raise RuntimeError("delegation is not active")
  457. record.update(
  458. {
  459. "status": "revoked",
  460. "updated_by": _uid(actor_uid, "actor_uid"),
  461. "updated_at": self.now_factory().isoformat(),
  462. }
  463. )
  464. try:
  465. result = self.repository.update_delegation(record, int(expected_version))
  466. self.commit()
  467. return result
  468. except Exception:
  469. self.rollback()
  470. raise
  471. def expire_delegations(
  472. self, *, at: datetime | str | None = None, actor_uid: str
  473. ) -> list[dict[str, Any]]:
  474. timestamp = _time(at, "at") if at is not None else self.now_factory()
  475. try:
  476. result = self.repository.expire_delegations(
  477. timestamp, _uid(actor_uid, "actor_uid")
  478. )
  479. self.commit()
  480. return result
  481. except Exception:
  482. self.rollback()
  483. raise
  484. def list_delegations(self, *, status: str | None = None) -> list[dict[str, Any]]:
  485. normalized = None
  486. if status is not None:
  487. normalized = _string(status, "status", 20)
  488. if normalized not in {"active", "revoked", "expired"}:
  489. raise ValueError("unsupported delegation status")
  490. return self.repository.list_delegations(status=normalized)
  491. def create_policy(self, payload: Any, *, actor_uid: str) -> dict[str, Any]:
  492. body = _closed(
  493. payload,
  494. {
  495. "code",
  496. "name",
  497. "policy_type",
  498. "scope_type",
  499. "scope_uid",
  500. "definition",
  501. },
  502. "responsibility policy",
  503. )
  504. code = _string(body.get("code"), "code", 120).upper()
  505. if not CODE_PATTERN.fullmatch(code):
  506. raise ValueError("responsibility policy code is invalid")
  507. policy_type = _string(body.get("policy_type"), "policy_type", 30)
  508. if policy_type not in POLICY_TYPES:
  509. raise ValueError("unsupported responsibility policy type")
  510. actor = _uid(actor_uid, "actor_uid")
  511. now = self.now_factory().isoformat()
  512. policy_uid = self.uid_factory()
  513. version_uid = self.uid_factory()
  514. policy = {
  515. "uid": policy_uid,
  516. "code": code,
  517. "name": _string(body.get("name"), "name", 300),
  518. "policy_type": policy_type,
  519. "scope_type": _resource_type(body.get("scope_type"), "scope_type"),
  520. "scope_uid": _string(body.get("scope_uid"), "scope_uid", 120),
  521. "status": "draft",
  522. "current_version": 1,
  523. "active_version_uid": None,
  524. "created_by": actor,
  525. "created_at": now,
  526. "updated_at": now,
  527. }
  528. version = {
  529. "uid": version_uid,
  530. "policy_uid": policy_uid,
  531. "version": 1,
  532. "status": "draft",
  533. "definition": _policy_definition(policy_type, body.get("definition")),
  534. "created_by": actor,
  535. "created_at": now,
  536. "published_by": None,
  537. "published_at": None,
  538. }
  539. try:
  540. self.repository.create_policy(policy, version)
  541. self.commit()
  542. return {**policy, "latest_version": version}
  543. except Exception:
  544. self.rollback()
  545. raise
  546. def list_policies(self) -> list[dict[str, Any]]:
  547. return self.repository.list_policies()
  548. def revise_policy(
  549. self,
  550. policy_uid: str,
  551. definition: Any,
  552. *,
  553. expected_version: int,
  554. actor_uid: str,
  555. ) -> dict[str, Any]:
  556. uid = _uid(policy_uid, "policy_uid")
  557. policy = self.repository.get_policy(uid)
  558. if policy is None:
  559. raise LookupError("responsibility policy was not found")
  560. if int(policy["current_version"]) != int(expected_version):
  561. raise RuntimeError("policy version conflict")
  562. now = self.now_factory().isoformat()
  563. version_number = int(expected_version) + 1
  564. version = {
  565. "uid": self.uid_factory(),
  566. "policy_uid": uid,
  567. "version": version_number,
  568. "status": "draft",
  569. "definition": _policy_definition(policy["policy_type"], definition),
  570. "created_by": _uid(actor_uid, "actor_uid"),
  571. "created_at": now,
  572. "published_by": None,
  573. "published_at": None,
  574. }
  575. updated = {
  576. **policy,
  577. "current_version": version_number,
  578. "updated_at": now,
  579. }
  580. try:
  581. self.repository.revise_policy(updated, version, int(expected_version))
  582. self.commit()
  583. return {**updated, "latest_version": version}
  584. except Exception:
  585. self.rollback()
  586. raise
  587. def publish_policy(
  588. self,
  589. policy_uid: str,
  590. *,
  591. expected_version: int,
  592. actor_uid: str,
  593. ) -> dict[str, Any]:
  594. uid = _uid(policy_uid, "policy_uid")
  595. policy = self.repository.get_policy(uid)
  596. if policy is None:
  597. raise LookupError("responsibility policy was not found")
  598. if int(policy["current_version"]) != int(expected_version):
  599. raise RuntimeError("policy version conflict")
  600. version = self.repository.policy_version(uid, int(expected_version))
  601. if version is None or version["status"] != "draft":
  602. raise RuntimeError("responsibility policy version is not publishable")
  603. actor = _uid(actor_uid, "actor_uid")
  604. now = self.now_factory().isoformat()
  605. published_version = {
  606. **version,
  607. "status": "published",
  608. "published_by": actor,
  609. "published_at": now,
  610. }
  611. published = {
  612. **policy,
  613. "status": "published",
  614. "active_version_uid": version["uid"],
  615. "updated_at": now,
  616. }
  617. try:
  618. self.repository.publish_policy(
  619. published, published_version, int(expected_version)
  620. )
  621. self.commit()
  622. return {**published, "active_version": published_version}
  623. except Exception:
  624. self.rollback()
  625. raise
  626. def evaluate_joint_review(
  627. self,
  628. resource_type: str,
  629. resource_uid: str,
  630. decisions: Any,
  631. ) -> dict[str, Any]:
  632. resolved = self.resolve(resource_type, resource_uid)
  633. policy = resolved.get("joint_review")
  634. if policy is None:
  635. raise LookupError("joint review policy was not found")
  636. if not isinstance(decisions, list):
  637. raise ValueError("joint review decisions must be an array")
  638. by_domain = {}
  639. seen_users = set()
  640. for raw in decisions:
  641. item = _closed(
  642. raw,
  643. {"user_uid", "domain_uid", "decision"},
  644. "joint review decision",
  645. )
  646. user_uid = _uid(item.get("user_uid"), "user_uid")
  647. domain_uid = _string(item.get("domain_uid"), "domain_uid", 120)
  648. decision = _string(item.get("decision"), "decision", 20)
  649. if decision not in DECISIONS:
  650. raise ValueError("unsupported joint review decision")
  651. if domain_uid not in policy["required_domains"]:
  652. raise ValueError("review domain is not eligible")
  653. domain = self.resolve("business_domain", domain_uid)
  654. eligible_users = {
  655. owner["effective_user_id"] for owner in domain["final_owners"]
  656. }
  657. if user_uid not in eligible_users:
  658. raise ValueError("reviewer is not the final owner of the domain")
  659. if user_uid in seen_users or domain_uid in by_domain:
  660. raise ValueError("joint review decisions must be independent")
  661. seen_users.add(user_uid)
  662. by_domain[domain_uid] = decision
  663. approved_domains = sorted(
  664. domain for domain, decision in by_domain.items() if decision == "approve"
  665. )
  666. missing = sorted(set(policy["required_domains"]) - set(approved_domains))
  667. if "reject" in by_domain.values():
  668. status = "rejected"
  669. elif len(approved_domains) < policy["min_approvals"] or (
  670. policy["require_all_domains"] and missing
  671. ):
  672. status = "pending"
  673. else:
  674. status = "approved"
  675. return {
  676. "status": status,
  677. "policy_uid": policy["policy_uid"],
  678. "approval_count": len(approved_domains),
  679. "approved_domains": approved_domains,
  680. "missing_domains": missing,
  681. "deterministic": True,
  682. }
  683. def operations(self, owner_uid: str) -> dict[str, Any]:
  684. owner = _uid(owner_uid, "owner_uid")
  685. result = self.repository.responsibility_operations(owner)
  686. tasks = result.get("tasks") or []
  687. metrics = result.get("metrics") or []
  688. return {
  689. "owner_uid": owner,
  690. "tasks": tasks,
  691. "metrics": metrics,
  692. "summary": {
  693. "task_count": len(tasks),
  694. "metric_count": len(metrics),
  695. "overdue_count": sum(bool(item.get("overdue")) for item in tasks),
  696. "recurrent_count": sum(
  697. int(item.get("recurrence_count") or 0) > 1 for item in tasks
  698. ),
  699. },
  700. "grants_data_access": False,
  701. }