"""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, }