| 123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485486487488489490491492493494495496497498499500501502503504505506507508509510511512513514515516517518519520521522523524525526527528529530531532533534535536537538539540541542543544545546547548549550551552553554555556557558559560561562563564565566567568569570571572573574575576577578579580581582583584585586587588589590591592593594595596597598599600601602603604605606607608609610611612613614615616617618619620621622623624625626627628629630631632633634635636637638639640641642643644645646647648649650651652653654655656657658659660661662663664665666667668669670671672673674675676677678679680681682683684685686687688689690691692693694695696697698699700701702703704705706707708709710711712713714715716717718719720721722723724725726727728729730731732733734 |
- """Hierarchical responsibility resolution, delegation and central policies."""
- from __future__ import annotations
- import copy
- import re
- import uuid
- from collections.abc import Callable
- from datetime import datetime
- from typing import Any
- from app.core.common.identifiers import new_governance_uid
- from app.core.common.timezone_utils import now_china
- from app.core.governance.responsibilities import (
- RACI_ROLES,
- RESPONSIBILITY_ROLES,
- )
- HIERARCHICAL_RESOURCE_TYPES = frozenset(
- {
- "organization",
- "business_domain",
- "data_asset",
- "semantic_term",
- "data_standard",
- "quality_policy",
- "data_product",
- "agent",
- }
- )
- POLICY_TYPES = frozenset({"central_policy", "joint_review"})
- DELEGATION_TYPES = frozenset({"temporary", "departure_transfer"})
- DECISIONS = frozenset({"approve", "reject"})
- CODE_PATTERN = re.compile(r"^[A-Z][A-Z0-9_]{2,119}$")
- def _closed(value: Any, allowed: set[str], label: str) -> dict[str, Any]:
- if not isinstance(value, dict):
- raise ValueError(f"{label} must be an object")
- unknown = sorted(set(value) - allowed)
- if unknown:
- raise ValueError(
- f"{label} contains unsupported fields: {', '.join(unknown)}"
- )
- return copy.deepcopy(value)
- def _string(value: Any, label: str, maximum: int = 500) -> str:
- if not isinstance(value, str) or not value.strip():
- raise ValueError(f"{label} is required")
- result = value.strip()
- if len(result) > maximum:
- raise ValueError(f"{label} exceeds {maximum} characters")
- return result
- def _uid(value: Any, label: str) -> str:
- try:
- return str(uuid.UUID(str(value)))
- except (TypeError, ValueError, AttributeError) as error:
- raise ValueError(f"{label} must be a UUID") from error
- def _time(value: Any, label: str) -> datetime:
- if isinstance(value, datetime):
- result = value
- else:
- try:
- result = datetime.fromisoformat(str(value).replace("Z", "+00:00"))
- except (TypeError, ValueError) as error:
- raise ValueError(f"{label} must be an ISO-8601 datetime") from error
- if result.tzinfo is None:
- raise ValueError(f"{label} must include a timezone")
- return result
- def _resource_type(value: Any, label: str = "resource_type") -> str:
- result = _string(value, label, 40)
- if result not in HIERARCHICAL_RESOURCE_TYPES:
- raise ValueError(f"unsupported {label}")
- return result
- def _policy_definition(policy_type: str, value: Any) -> dict[str, Any]:
- if policy_type == "central_policy":
- definition = _closed(
- value,
- {
- "required_responsibility_roles",
- "require_unique_accountable",
- "max_delegation_days",
- },
- "central policy definition",
- )
- roles = definition.get("required_responsibility_roles")
- if not isinstance(roles, list) or not roles:
- raise ValueError("central policy required roles are missing")
- normalized_roles = sorted({_string(item, "required role", 40) for item in roles})
- if not set(normalized_roles) <= RESPONSIBILITY_ROLES:
- raise ValueError("central policy contains unsupported role")
- if not isinstance(definition.get("require_unique_accountable"), bool):
- raise ValueError("require_unique_accountable must be a boolean")
- try:
- days = int(definition.get("max_delegation_days"))
- except (TypeError, ValueError) as error:
- raise ValueError("max_delegation_days is invalid") from error
- if days < 1 or days > 366:
- raise ValueError("max_delegation_days must be between 1 and 366")
- return {
- "required_responsibility_roles": normalized_roles,
- "require_unique_accountable": definition[
- "require_unique_accountable"
- ],
- "max_delegation_days": days,
- }
- definition = _closed(
- value,
- {"business_domain_uids", "min_approvals", "require_all_domains"},
- "joint review definition",
- )
- domains = definition.get("business_domain_uids")
- if not isinstance(domains, list) or len(domains) < 2:
- raise ValueError("joint review requires at least two business domains")
- normalized_domains = sorted(
- {_string(item, "business domain uid", 120) for item in domains}
- )
- try:
- minimum = int(definition.get("min_approvals"))
- except (TypeError, ValueError) as error:
- raise ValueError("min_approvals is invalid") from error
- if minimum < 1 or minimum > len(normalized_domains):
- raise ValueError("min_approvals is outside the domain range")
- if not isinstance(definition.get("require_all_domains"), bool):
- raise ValueError("require_all_domains must be a boolean")
- return {
- "business_domain_uids": normalized_domains,
- "min_approvals": minimum,
- "require_all_domains": definition["require_all_domains"],
- }
- class UnifiedResponsibilityService:
- """Resolve effective accountability without granting platform access."""
- def __init__(
- self,
- repository,
- *,
- uid_factory: Callable[[], str] = new_governance_uid,
- now_factory: Callable[[], datetime] = now_china,
- commit: Callable[[], None] = lambda: None,
- rollback: Callable[[], None] = lambda: None,
- ):
- self.repository = repository
- self.uid_factory = uid_factory
- self.now_factory = now_factory
- self.commit = commit
- self.rollback = rollback
- def _chain(self, resource_type: str, resource_uid: str) -> list[dict[str, Any]]:
- kind = _resource_type(resource_type)
- uid = _string(resource_uid, "resource_uid", 120)
- chain = []
- seen = set()
- for depth in range(20):
- identity = (kind, uid)
- if identity in seen:
- raise ValueError("responsibility hierarchy contains a cycle")
- seen.add(identity)
- parent = self.repository.parent(kind, uid)
- chain.append(
- {
- "resource_type": kind,
- "resource_uid": uid,
- "depth": depth,
- "hierarchy_revision": int(parent["revision"]) if parent else 0,
- }
- )
- if parent is None:
- return chain
- kind = _resource_type(parent["parent_type"], "parent_type")
- uid = _string(parent["parent_uid"], "parent_uid", 120)
- raise ValueError("responsibility hierarchy exceeds 20 levels")
- def resolve(
- self,
- resource_type: str,
- resource_uid: str,
- *,
- at: datetime | str | None = None,
- ) -> dict[str, Any]:
- evaluated_at = _time(at, "at") if at is not None else self.now_factory()
- chain = self._chain(resource_type, resource_uid)
- effective_by_raci: dict[str, list[dict[str, Any]]] = {}
- for node in chain:
- matrix = self.repository.get(
- node["resource_type"], node["resource_uid"]
- )
- assignments = matrix.get("assignments") or []
- for raci_role in RACI_ROLES:
- if raci_role in effective_by_raci:
- continue
- matching = [
- item for item in assignments if item["raci_role"] == raci_role
- ]
- if matching:
- effective_by_raci[raci_role] = [
- {
- **copy.deepcopy(item),
- "assigned_user_id": item["user_id"],
- "source_resource_type": node["resource_type"],
- "source_resource_uid": node["resource_uid"],
- "source_revision": matrix.get("revision", 0),
- "inheritance_depth": node["depth"],
- "inherited": node["depth"] > 0,
- }
- for item in matching
- ]
- assignments = []
- for raci_role in sorted(effective_by_raci):
- for assignment in effective_by_raci[raci_role]:
- delegation = self.repository.active_delegation(
- assignment["assigned_user_id"],
- assignment["responsibility_role"],
- chain,
- evaluated_at,
- )
- effective_user = (
- delegation["delegate_user_uid"]
- if delegation
- else assignment["assigned_user_id"]
- )
- active = effective_user in self.repository.users_available(
- [effective_user]
- )
- assignments.append(
- {
- **assignment,
- "effective_user_id": effective_user,
- "delegation_uid": delegation["uid"] if delegation else None,
- "delegation_type": (
- delegation["delegation_type"] if delegation else None
- ),
- "effective_user_active": active,
- }
- )
- final_owners = [
- item
- for item in assignments
- if item["raci_role"] == "accountable"
- and item["effective_user_active"]
- ]
- policies = self.repository.policies_for_chain(chain)
- central = [
- item for item in policies if item["policy_type"] == "central_policy"
- ]
- required_roles = sorted(
- {
- role
- for policy in central
- for role in policy["definition"]["required_responsibility_roles"]
- }
- )
- present_roles = {item["responsibility_role"] for item in assignments}
- missing_roles = sorted(set(required_roles) - present_roles)
- unique_required = any(
- item["definition"]["require_unique_accountable"] for item in central
- )
- compliant = not missing_roles and (
- not unique_required or len(final_owners) == 1
- )
- joint_policies = [
- item for item in policies if item["policy_type"] == "joint_review"
- ]
- joint = None
- if joint_policies:
- definition = joint_policies[0]["definition"]
- joint = {
- "policy_uid": joint_policies[0]["uid"],
- "required_domains": definition["business_domain_uids"],
- "min_approvals": definition["min_approvals"],
- "require_all_domains": definition["require_all_domains"],
- }
- status = (
- "resolved"
- if len(final_owners) == 1
- else "unresolved"
- if not final_owners
- else "ambiguous"
- )
- return {
- "resource_type": chain[0]["resource_type"],
- "resource_uid": chain[0]["resource_uid"],
- "status": status,
- "evaluated_at": evaluated_at.isoformat(),
- "chain": chain,
- "assignments": assignments,
- "final_owners": final_owners,
- "policy_compliance": {
- "status": "compliant" if compliant else "non_compliant",
- "required_roles": required_roles,
- "missing_roles": missing_roles,
- "unique_accountable_required": unique_required,
- },
- "joint_review": joint,
- "grants_data_access": False,
- }
- def set_parent(
- self,
- payload: Any,
- *,
- expected_revision: int,
- actor_uid: str,
- ) -> dict[str, Any]:
- body = _closed(
- payload,
- {"resource_type", "resource_uid", "parent_type", "parent_uid"},
- "responsibility hierarchy",
- )
- resource_type = _resource_type(body.get("resource_type"))
- resource_uid = _string(body.get("resource_uid"), "resource_uid", 120)
- parent_type = _resource_type(body.get("parent_type"), "parent_type")
- parent_uid = _string(body.get("parent_uid"), "parent_uid", 120)
- if (resource_type, resource_uid) == (parent_type, parent_uid):
- raise ValueError("responsibility hierarchy contains a cycle")
- parent_chain = self._chain(parent_type, parent_uid)
- if any(
- (item["resource_type"], item["resource_uid"])
- == (resource_type, resource_uid)
- for item in parent_chain
- ):
- raise ValueError("responsibility hierarchy contains a cycle")
- try:
- revision = int(expected_revision)
- except (TypeError, ValueError) as error:
- raise ValueError("hierarchy revision is invalid") from error
- now = self.now_factory().isoformat()
- record = {
- "resource_type": resource_type,
- "resource_uid": resource_uid,
- "parent_type": parent_type,
- "parent_uid": parent_uid,
- "updated_by": _uid(actor_uid, "actor_uid"),
- "updated_at": now,
- }
- try:
- result = self.repository.set_parent(record, revision)
- self.commit()
- return result
- except Exception:
- self.rollback()
- raise
- def create_delegation(
- self, payload: Any, *, actor_uid: str
- ) -> dict[str, Any]:
- body = _closed(
- payload,
- {
- "source_user_uid",
- "delegate_user_uid",
- "scope_type",
- "scope_uid",
- "responsibility_role",
- "delegation_type",
- "starts_at",
- "ends_at",
- "reason",
- },
- "responsibility delegation",
- )
- source = _uid(body.get("source_user_uid"), "source_user_uid")
- delegate = _uid(body.get("delegate_user_uid"), "delegate_user_uid")
- if source == delegate:
- raise ValueError("delegation users must be different")
- delegation_type = _string(
- body.get("delegation_type"), "delegation_type", 30
- )
- if delegation_type not in DELEGATION_TYPES:
- raise ValueError("unsupported delegation type")
- scope_type = body.get("scope_type")
- scope_uid = body.get("scope_uid")
- if (scope_type is None) != (scope_uid is None):
- raise ValueError("delegation scope type and uid must be paired")
- if scope_type is not None:
- scope_type = _resource_type(scope_type, "scope_type")
- scope_uid = _string(scope_uid, "scope_uid", 120)
- role = body.get("responsibility_role")
- if role is not None:
- role = _string(role, "responsibility_role", 40)
- if role not in RESPONSIBILITY_ROLES:
- raise ValueError("unsupported responsibility role")
- starts_at = _time(body.get("starts_at"), "starts_at")
- ends_at = (
- _time(body.get("ends_at"), "ends_at")
- if body.get("ends_at") is not None
- else None
- )
- if delegation_type == "temporary":
- if ends_at is None or ends_at <= starts_at:
- raise ValueError("temporary delegation requires a later end time")
- if (ends_at - starts_at).days > 366:
- raise ValueError("temporary delegation exceeds 366 days")
- if scope_type is not None:
- policies = self.repository.policies_for_chain(
- self._chain(scope_type, scope_uid)
- )
- limits = [
- int(item["definition"]["max_delegation_days"])
- for item in policies
- if item["policy_type"] == "central_policy"
- ]
- if limits and (ends_at - starts_at).total_seconds() > min(limits) * 86400:
- raise ValueError("temporary delegation exceeds central policy limit")
- elif ends_at is not None:
- raise ValueError("departure transfer cannot have an end time")
- required_users = [delegate] + ([source] if delegation_type == "temporary" else [])
- if self.repository.users_available(required_users) != set(required_users):
- raise ValueError("delegation user is unknown or disabled")
- now = self.now_factory().isoformat()
- record = {
- "uid": self.uid_factory(),
- "source_user_uid": source,
- "delegate_user_uid": delegate,
- "scope_type": scope_type,
- "scope_uid": scope_uid,
- "responsibility_role": role,
- "delegation_type": delegation_type,
- "starts_at": starts_at.isoformat(),
- "ends_at": ends_at.isoformat() if ends_at else None,
- "reason": _string(body.get("reason"), "reason", 500),
- "status": "active",
- "current_version": 1,
- "created_by": _uid(actor_uid, "actor_uid"),
- "created_at": now,
- "updated_by": _uid(actor_uid, "actor_uid"),
- "updated_at": now,
- }
- try:
- result = self.repository.create_delegation(record)
- self.commit()
- return result
- except Exception:
- self.rollback()
- raise
- def transfer_departing_user(
- self, payload: Any, *, actor_uid: str
- ) -> dict[str, Any]:
- body = _closed(
- payload,
- {
- "source_user_uid",
- "delegate_user_uid",
- "scope_type",
- "scope_uid",
- "responsibility_role",
- "reason",
- },
- "departure transfer",
- )
- return self.create_delegation(
- {
- **body,
- "responsibility_role": body.get("responsibility_role"),
- "delegation_type": "departure_transfer",
- "starts_at": self.now_factory().isoformat(),
- "ends_at": None,
- },
- actor_uid=actor_uid,
- )
- def revoke_delegation(
- self, uid: str, *, expected_version: int, actor_uid: str
- ) -> dict[str, Any]:
- record = self.repository.get_delegation(_uid(uid, "delegation_uid"))
- if record is None:
- raise LookupError("delegation was not found")
- if record["status"] != "active":
- raise RuntimeError("delegation is not active")
- record.update(
- {
- "status": "revoked",
- "updated_by": _uid(actor_uid, "actor_uid"),
- "updated_at": self.now_factory().isoformat(),
- }
- )
- try:
- result = self.repository.update_delegation(record, int(expected_version))
- self.commit()
- return result
- except Exception:
- self.rollback()
- raise
- def expire_delegations(
- self, *, at: datetime | str | None = None, actor_uid: str
- ) -> list[dict[str, Any]]:
- timestamp = _time(at, "at") if at is not None else self.now_factory()
- try:
- result = self.repository.expire_delegations(
- timestamp, _uid(actor_uid, "actor_uid")
- )
- self.commit()
- return result
- except Exception:
- self.rollback()
- raise
- def list_delegations(self, *, status: str | None = None) -> list[dict[str, Any]]:
- normalized = None
- if status is not None:
- normalized = _string(status, "status", 20)
- if normalized not in {"active", "revoked", "expired"}:
- raise ValueError("unsupported delegation status")
- return self.repository.list_delegations(status=normalized)
- def create_policy(self, payload: Any, *, actor_uid: str) -> dict[str, Any]:
- body = _closed(
- payload,
- {
- "code",
- "name",
- "policy_type",
- "scope_type",
- "scope_uid",
- "definition",
- },
- "responsibility policy",
- )
- code = _string(body.get("code"), "code", 120).upper()
- if not CODE_PATTERN.fullmatch(code):
- raise ValueError("responsibility policy code is invalid")
- policy_type = _string(body.get("policy_type"), "policy_type", 30)
- if policy_type not in POLICY_TYPES:
- raise ValueError("unsupported responsibility policy type")
- actor = _uid(actor_uid, "actor_uid")
- now = self.now_factory().isoformat()
- policy_uid = self.uid_factory()
- version_uid = self.uid_factory()
- policy = {
- "uid": policy_uid,
- "code": code,
- "name": _string(body.get("name"), "name", 300),
- "policy_type": policy_type,
- "scope_type": _resource_type(body.get("scope_type"), "scope_type"),
- "scope_uid": _string(body.get("scope_uid"), "scope_uid", 120),
- "status": "draft",
- "current_version": 1,
- "active_version_uid": None,
- "created_by": actor,
- "created_at": now,
- "updated_at": now,
- }
- version = {
- "uid": version_uid,
- "policy_uid": policy_uid,
- "version": 1,
- "status": "draft",
- "definition": _policy_definition(policy_type, body.get("definition")),
- "created_by": actor,
- "created_at": now,
- "published_by": None,
- "published_at": None,
- }
- try:
- self.repository.create_policy(policy, version)
- self.commit()
- return {**policy, "latest_version": version}
- except Exception:
- self.rollback()
- raise
- def list_policies(self) -> list[dict[str, Any]]:
- return self.repository.list_policies()
- def revise_policy(
- self,
- policy_uid: str,
- definition: Any,
- *,
- expected_version: int,
- actor_uid: str,
- ) -> dict[str, Any]:
- uid = _uid(policy_uid, "policy_uid")
- policy = self.repository.get_policy(uid)
- if policy is None:
- raise LookupError("responsibility policy was not found")
- if int(policy["current_version"]) != int(expected_version):
- raise RuntimeError("policy version conflict")
- now = self.now_factory().isoformat()
- version_number = int(expected_version) + 1
- version = {
- "uid": self.uid_factory(),
- "policy_uid": uid,
- "version": version_number,
- "status": "draft",
- "definition": _policy_definition(policy["policy_type"], definition),
- "created_by": _uid(actor_uid, "actor_uid"),
- "created_at": now,
- "published_by": None,
- "published_at": None,
- }
- updated = {
- **policy,
- "current_version": version_number,
- "updated_at": now,
- }
- try:
- self.repository.revise_policy(updated, version, int(expected_version))
- self.commit()
- return {**updated, "latest_version": version}
- except Exception:
- self.rollback()
- raise
- def publish_policy(
- self,
- policy_uid: str,
- *,
- expected_version: int,
- actor_uid: str,
- ) -> dict[str, Any]:
- uid = _uid(policy_uid, "policy_uid")
- policy = self.repository.get_policy(uid)
- if policy is None:
- raise LookupError("responsibility policy was not found")
- if int(policy["current_version"]) != int(expected_version):
- raise RuntimeError("policy version conflict")
- version = self.repository.policy_version(uid, int(expected_version))
- if version is None or version["status"] != "draft":
- raise RuntimeError("responsibility policy version is not publishable")
- actor = _uid(actor_uid, "actor_uid")
- now = self.now_factory().isoformat()
- published_version = {
- **version,
- "status": "published",
- "published_by": actor,
- "published_at": now,
- }
- published = {
- **policy,
- "status": "published",
- "active_version_uid": version["uid"],
- "updated_at": now,
- }
- try:
- self.repository.publish_policy(
- published, published_version, int(expected_version)
- )
- self.commit()
- return {**published, "active_version": published_version}
- except Exception:
- self.rollback()
- raise
- def evaluate_joint_review(
- self,
- resource_type: str,
- resource_uid: str,
- decisions: Any,
- ) -> dict[str, Any]:
- resolved = self.resolve(resource_type, resource_uid)
- policy = resolved.get("joint_review")
- if policy is None:
- raise LookupError("joint review policy was not found")
- if not isinstance(decisions, list):
- raise ValueError("joint review decisions must be an array")
- by_domain = {}
- seen_users = set()
- for raw in decisions:
- item = _closed(
- raw,
- {"user_uid", "domain_uid", "decision"},
- "joint review decision",
- )
- user_uid = _uid(item.get("user_uid"), "user_uid")
- domain_uid = _string(item.get("domain_uid"), "domain_uid", 120)
- decision = _string(item.get("decision"), "decision", 20)
- if decision not in DECISIONS:
- raise ValueError("unsupported joint review decision")
- if domain_uid not in policy["required_domains"]:
- raise ValueError("review domain is not eligible")
- domain = self.resolve("business_domain", domain_uid)
- eligible_users = {
- owner["effective_user_id"] for owner in domain["final_owners"]
- }
- if user_uid not in eligible_users:
- raise ValueError("reviewer is not the final owner of the domain")
- if user_uid in seen_users or domain_uid in by_domain:
- raise ValueError("joint review decisions must be independent")
- seen_users.add(user_uid)
- by_domain[domain_uid] = decision
- approved_domains = sorted(
- domain for domain, decision in by_domain.items() if decision == "approve"
- )
- missing = sorted(set(policy["required_domains"]) - set(approved_domains))
- if "reject" in by_domain.values():
- status = "rejected"
- elif len(approved_domains) < policy["min_approvals"] or (
- policy["require_all_domains"] and missing
- ):
- status = "pending"
- else:
- status = "approved"
- return {
- "status": status,
- "policy_uid": policy["policy_uid"],
- "approval_count": len(approved_domains),
- "approved_domains": approved_domains,
- "missing_domains": missing,
- "deterministic": True,
- }
- def operations(self, owner_uid: str) -> dict[str, Any]:
- owner = _uid(owner_uid, "owner_uid")
- result = self.repository.responsibility_operations(owner)
- tasks = result.get("tasks") or []
- metrics = result.get("metrics") or []
- return {
- "owner_uid": owner,
- "tasks": tasks,
- "metrics": metrics,
- "summary": {
- "task_count": len(tasks),
- "metric_count": len(metrics),
- "overdue_count": sum(bool(item.get("overdue")) for item in tasks),
- "recurrent_count": sum(
- int(item.get("recurrence_count") or 0) > 1 for item in tasks
- ),
- },
- "grants_data_access": False,
- }
|