from __future__ import annotations from collections.abc import Mapping, Sequence from dataclasses import dataclass from datetime import datetime, timezone from sqlalchemy import or_ from sqlalchemy.orm import Session from govoplan_cases.backend.db.models import ( CaseAccessGrant, CaseIdentity, CaseRecordRevision, CaseTimelineEntry, ) from govoplan_cases.backend.domain import CaseRecord from govoplan_core.core.dsar import ( DsarErasureActionRef, DsarExecutionResultRef, DsarRecordRef, DsarSubjectRef, dsar_capability_name, ) CASES_DSAR_CAPABILITY = dsar_capability_name("cases") _MAX_RECORDS = 5_000 @dataclass(frozen=True, slots=True) class _SubjectSelectors: account_id: str | None identity_id: str | None membership_id: str | None case_id: str | None revision_id: str | None access_grant_id: str | None timeline_id: str | None @property def actor_ids(self) -> tuple[str, ...]: return tuple( value for value in (self.account_id, self.identity_id, self.membership_id) if value ) @property def has_canonical_selector(self) -> bool: return bool(self.actor_ids) @property def has_direct_selector(self) -> bool: return bool( self.case_id or self.revision_id or self.access_grant_id or self.timeline_id ) @property def has_recognized_selector(self) -> bool: return self.has_canonical_selector or self.has_direct_selector class CasesDsarProvider: provider_id = "cases" module_id = "cases" def search_subject( self, session: object, *, tenant_id: str, subject: DsarSubjectRef, ) -> Sequence[DsarRecordRef]: db = _session(session) selectors = _subject_selectors(subject) if selectors is None or not selectors.has_recognized_selector: return () identities = _matching_identities(db, tenant_id, selectors) revisions = _matching_revisions(db, tenant_id, selectors) grants = _matching_grants(db, tenant_id, selectors) timeline = _matching_timeline(db, tenant_id, selectors) if _direct_reference_conflicts( selectors, identities=identities, revisions=revisions, grants=grants, timeline=timeline, ): return () direct_case_ids = _direct_case_ids( selectors, identities=identities, revisions=revisions, grants=grants, timeline=timeline, ) records: list[DsarRecordRef] = [] seen: set[tuple[str, str]] = set() def append(record: DsarRecordRef) -> None: key = (record.resource_type, record.resource_id) if key in seen: return if len(records) >= _MAX_RECORDS: raise ValueError( "Cases DSAR result limit exceeded; narrow the subject selectors." ) seen.add(key) records.append(record) for identity in identities: if identity.case_id in direct_case_ids: append(_case_identity_record(identity)) if identity.created_by in selectors.actor_ids: append( _operator_record( resource_id=f"identity:{identity.id}", case_id=identity.case_id, activity="created_case_identity", observed_at=identity.created_at, ) ) for revision in revisions: if _revision_is_direct(revision, selectors, direct_case_ids): record = CaseRecord.from_mapping(revision.snapshot) append(_revision_record(revision, record)) if revision.superseded_at is None and revision.closed_at is None: append(_current_fact_record(revision, record)) if revision.changed_by in selectors.actor_ids: append( _operator_record( resource_id=f"revision:{revision.id}", case_id=revision.case_id, activity="changed_case_revision", observed_at=revision.recorded_at, case_revision=revision.revision, ) ) for grant in grants: if _grant_subject_matches(grant, selectors) or ( _grant_is_direct(grant, selectors) and not selectors.has_canonical_selector ): append(_access_grant_record(grant)) if grant.created_by in selectors.actor_ids: append( _operator_record( resource_id=f"grant:{grant.id}", case_id=grant.case_id, activity="changed_case_access_grant", observed_at=grant.updated_at, ) ) for entry in timeline: direct = _timeline_is_direct(entry, selectors, direct_case_ids) actor_match = entry.actor_id in selectors.actor_ids if direct or actor_match: append( _timeline_record( entry, expose_actor=actor_match, match_fields=( (["reference"] if direct else []) + (["actor_id"] if actor_match else []) ), ) ) return tuple(records) def plan_erasure( self, session: object, *, tenant_id: str, subject: DsarSubjectRef, records: Sequence[DsarRecordRef], ) -> Sequence[DsarErasureActionRef]: del tenant_id _session(session) if _subject_selectors(subject) is None: raise ValueError("Cases DSAR subject selectors conflict.") actions: list[DsarErasureActionRef] = [] for record in records: _validate_record(record) if record.immutable_evidence: kind = "retain" title = f"Retain {record.title}" rationale = record.retention_reason or ( "Case history is retained as institutional evidence." ) else: kind = "manual_review" title = f"Review {record.title}" rationale = ( "An authorized case operator must amend, close, supersede, or " "deactivate the current fact through the governed case lifecycle " "after reviewing legal, procedural, access, and third-party effects." ) actions.append( DsarErasureActionRef( action_id=f"cases:{kind}:{record.resource_type}:{record.resource_id}", provider_id=self.provider_id, module_id=self.module_id, kind=kind, resource_type=record.resource_type, resource_id=record.resource_id, title=title, rationale=rationale, executable=False, ) ) return tuple(actions) def execute_erasure( self, session: object, *, tenant_id: str, subject: DsarSubjectRef, actions: Sequence[DsarErasureActionRef], request_id: str, ) -> Sequence[DsarExecutionResultRef]: del tenant_id _session(session) if _subject_selectors(subject) is None: raise ValueError("Cases DSAR subject selectors conflict.") results: list[DsarExecutionResultRef] = [] for action in actions: _validate_action(action) if action.executable: raise ValueError( "Cases DSAR does not publish executable erasure actions." ) results.append( DsarExecutionResultRef( action_id=action.action_id, status="blocked", summary=( "Use the governed case or access lifecycle after legal, " "institutional-evidence, and third-party review." ), evidence={"request_id": request_id}, ) ) return tuple(results) def _matching_identities( session: Session, tenant_id: str, selectors: _SubjectSelectors, ) -> list[CaseIdentity]: conditions = [] if selectors.actor_ids: conditions.append(CaseIdentity.created_by.in_(selectors.actor_ids)) if selectors.case_id: conditions.append(CaseIdentity.case_id == selectors.case_id) return _query_conditions(session, CaseIdentity, tenant_id, conditions) def _matching_revisions( session: Session, tenant_id: str, selectors: _SubjectSelectors, ) -> list[CaseRecordRevision]: conditions = [] if selectors.actor_ids: conditions.append(CaseRecordRevision.changed_by.in_(selectors.actor_ids)) if selectors.case_id: conditions.append(CaseRecordRevision.case_id == selectors.case_id) if selectors.revision_id: conditions.append(CaseRecordRevision.id == selectors.revision_id) return _query_conditions(session, CaseRecordRevision, tenant_id, conditions) def _matching_grants( session: Session, tenant_id: str, selectors: _SubjectSelectors, ) -> list[CaseAccessGrant]: conditions = [] for kind, value in ( ("account", selectors.account_id), ("identity", selectors.identity_id), ("membership", selectors.membership_id), ): if value: conditions.append( (CaseAccessGrant.subject_kind == kind) & (CaseAccessGrant.subject_id == value) ) if selectors.actor_ids: conditions.append(CaseAccessGrant.created_by.in_(selectors.actor_ids)) if selectors.case_id: conditions.append(CaseAccessGrant.case_id == selectors.case_id) if selectors.access_grant_id: conditions.append(CaseAccessGrant.id == selectors.access_grant_id) return _query_conditions(session, CaseAccessGrant, tenant_id, conditions) def _matching_timeline( session: Session, tenant_id: str, selectors: _SubjectSelectors, ) -> list[CaseTimelineEntry]: conditions = [] if selectors.actor_ids: conditions.append(CaseTimelineEntry.actor_id.in_(selectors.actor_ids)) if selectors.case_id: conditions.append(CaseTimelineEntry.case_id == selectors.case_id) if selectors.timeline_id: conditions.append(CaseTimelineEntry.id == selectors.timeline_id) return _query_conditions(session, CaseTimelineEntry, tenant_id, conditions) def _direct_reference_conflicts( selectors: _SubjectSelectors, *, identities: Sequence[CaseIdentity], revisions: Sequence[CaseRecordRevision], grants: Sequence[CaseAccessGrant], timeline: Sequence[CaseTimelineEntry], ) -> bool: referenced_case_ids: list[str] = [] if selectors.case_id: if not any(row.case_id == selectors.case_id for row in identities): return True referenced_case_ids.append(selectors.case_id) if selectors.revision_id: row = next( (item for item in revisions if item.id == selectors.revision_id), None ) if row is None: return True referenced_case_ids.append(row.case_id) if selectors.access_grant_id: row = next( (item for item in grants if item.id == selectors.access_grant_id), None ) if row is None: return True if selectors.has_canonical_selector and not _grant_subject_matches( row, selectors ): return True referenced_case_ids.append(row.case_id) if selectors.timeline_id: row = next( (item for item in timeline if item.id == selectors.timeline_id), None ) if row is None: return True if selectors.has_canonical_selector and row.actor_id not in selectors.actor_ids: return True referenced_case_ids.append(row.case_id) if len(set(referenced_case_ids)) > 1: return True if not selectors.has_canonical_selector or not referenced_case_ids: return False related_case_ids = { row.case_id for row in identities if row.created_by in selectors.actor_ids } related_case_ids.update( row.case_id for row in revisions if row.changed_by in selectors.actor_ids ) related_case_ids.update( row.case_id for row in grants if _grant_subject_matches(row, selectors) ) related_case_ids.update( row.case_id for row in timeline if row.actor_id in selectors.actor_ids ) return referenced_case_ids[0] not in related_case_ids def _direct_case_ids( selectors: _SubjectSelectors, *, identities: Sequence[CaseIdentity], revisions: Sequence[CaseRecordRevision], grants: Sequence[CaseAccessGrant], timeline: Sequence[CaseTimelineEntry], ) -> set[str]: values = {selectors.case_id} if selectors.case_id else set() values.update(row.case_id for row in revisions if row.id == selectors.revision_id) values.update(row.case_id for row in grants if row.id == selectors.access_grant_id) values.update(row.case_id for row in timeline if row.id == selectors.timeline_id) values.update( row.case_id for row in identities if selectors.case_id and row.case_id == selectors.case_id ) return {value for value in values if value} def _revision_is_direct( row: CaseRecordRevision, selectors: _SubjectSelectors, direct_case_ids: set[str], ) -> bool: if selectors.revision_id: return row.id == selectors.revision_id return bool(selectors.case_id and row.case_id in direct_case_ids) def _grant_is_direct( row: CaseAccessGrant, selectors: _SubjectSelectors, ) -> bool: return bool(selectors.access_grant_id and row.id == selectors.access_grant_id) def _timeline_is_direct( row: CaseTimelineEntry, selectors: _SubjectSelectors, direct_case_ids: set[str], ) -> bool: if selectors.timeline_id: return row.id == selectors.timeline_id return bool(selectors.case_id and row.case_id in direct_case_ids) def _grant_subject_matches( row: CaseAccessGrant, selectors: _SubjectSelectors, ) -> bool: return bool( (row.subject_kind == "account" and row.subject_id == selectors.account_id) or (row.subject_kind == "identity" and row.subject_id == selectors.identity_id) or ( row.subject_kind == "membership" and row.subject_id == selectors.membership_id ) ) def _case_identity_record(row: CaseIdentity) -> DsarRecordRef: return _record( "cases_case_identity", row.id, "case_identity", "Case identity", { "match_fields": ["reference"], "case_id": row.case_id, "case_number": _bounded_text(row.case_number, 255), "created_at": _iso(row.created_at), }, observed_at=row.created_at, immutable=True, retention_reason=( "The stable case identifier and number are retained so immutable case and " "record evidence remains reconstructable." ), ) def _revision_record( row: CaseRecordRevision, record: CaseRecord, ) -> DsarRecordRef: return _record( "cases_case_revision", row.id, "case_history", "Case revision", { "match_fields": ["reference"], "case_id": row.case_id, "case_number": _bounded_text(record.case_number, 255), "revision": row.revision, "previous_revision_id": row.previous_revision_id, "case_type_key": row.case_type_key, "status_key": row.status_key, "title": _bounded_text(row.title, 500), "access_mode": row.access_mode, "opened_at": _iso(row.opened_at), "deadline_at": _iso(row.deadline_at), "closed_at": _iso(row.closed_at), "recorded_at": _iso(row.recorded_at), "superseded_at": _iso(row.superseded_at), "service_ref": _institutional_reference(record.service_ref), "party_reference_count": len(record.party_refs), "assignment_reference_count": len(record.assignment_refs), "evidence_reference_count": len(record.evidence_refs), "decision_reference_count": len(record.decision_refs), "record_reference_count": len(record.record_refs), "access_grant_count": len(record.access_grants), }, observed_at=row.recorded_at, immutable=True, retention_reason=( "Case revisions are immutable procedure and accountability evidence; " "corrections append a new governed revision." ), ) def _current_fact_record( row: CaseRecordRevision, record: CaseRecord, ) -> DsarRecordRef: return _record( "cases_current_case_fact", row.case_id, "current_case_fact", "Current case fact", { "match_fields": ["reference"], "revision_record_id": row.id, "revision": row.revision, "case_number": _bounded_text(record.case_number, 255), "status_key": row.status_key, "title": _bounded_text(row.title, 500), "deadline_at": _iso(row.deadline_at), }, observed_at=row.recorded_at, ) def _access_grant_record(row: CaseAccessGrant) -> DsarRecordRef: immutable = not row.active return _record( "cases_access_grant", row.id, "case_access_fact", "Case access grant", { "match_fields": ["subject"], "case_id": row.case_id, "subject_kind": row.subject_kind, "subject_id": row.subject_id, "permissions": [str(value)[:40] for value in row.permissions[:20]], "allowed_purposes": [ str(value)[:255] for value in row.allowed_purposes[:100] ], "source": _bounded_text(row.source, 30), "active": row.active, "source_revision": row.source_revision, "created_at": _iso(row.created_at), "updated_at": _iso(row.updated_at), }, observed_at=row.updated_at, immutable=immutable, retention_reason=( "Inactive access-grant state is retained to explain historical case access." if immutable else None ), ) def _timeline_record( row: CaseTimelineEntry, *, expose_actor: bool, match_fields: Sequence[str], ) -> DsarRecordRef: return _record( "cases_timeline_event", row.id, "case_lifecycle_evidence", "Case lifecycle event", { "match_fields": list(match_fields), "case_id": row.case_id, "event_type": _bounded_text(row.event_type, 120), "purpose": _bounded_text(row.purpose, 255), "case_revision": row.case_revision, "occurred_at": _iso(row.occurred_at), "actor_id": row.actor_id if expose_actor else None, }, observed_at=row.occurred_at, immutable=True, retention_reason=( "Case timeline events are immutable lifecycle and accountability evidence." ), ) def _operator_record( *, resource_id: str, case_id: str, activity: str, observed_at: datetime | None, case_revision: int | None = None, ) -> DsarRecordRef: return _record( "cases_operator_attribution", resource_id, "operator_accountability_evidence", "Case operator attribution", { "match_fields": ["actor_id"], "case_id": case_id, "activity": activity, "case_revision": case_revision, "observed_at": _iso(observed_at), }, observed_at=observed_at, immutable=True, retention_reason=( "Operator attribution is retained as accountability evidence; raw case, " "party, evidence, change-reason, and event payload content is excluded." ), ) def _subject_selectors(subject: DsarSubjectRef) -> _SubjectSelectors | None: groups = { "account_id": ( subject.account_id, subject.external_references.get("cases.account"), subject.external_references.get("access.account"), ), "identity_id": ( subject.identity_id, subject.external_references.get("cases.identity"), subject.external_references.get("identity.id"), ), "membership_id": ( subject.membership_id, subject.external_references.get("cases.membership"), subject.external_references.get("tenancy.membership"), ), "case_id": (subject.external_references.get("cases.case"),), "revision_id": (subject.external_references.get("cases.revision"),), "access_grant_id": (subject.external_references.get("cases.access_grant"),), "timeline_id": (subject.external_references.get("cases.timeline"),), } normalized: dict[str, str | None] = {} for key, values in groups.items(): distinct = {value for item in values if (value := _normalized_id(item))} if len(distinct) > 1: return None normalized[key] = next(iter(distinct), None) return _SubjectSelectors(**normalized) def _query_conditions( session: Session, model: type, tenant_id: str, conditions: Sequence[object], ) -> list[object]: if not conditions: return [] rows = ( session.query(model) .filter(model.tenant_id == tenant_id, or_(*conditions)) .order_by(model.id.asc()) .limit(_MAX_RECORDS + 1) .all() ) if len(rows) > _MAX_RECORDS: raise ValueError("Cases DSAR match limit exceeded; narrow the selectors.") return rows def _institutional_reference(value: object | None) -> dict[str, object] | None: if value is None: return None return { "kind": str(getattr(value, "kind", "")), "owner_module": str(getattr(value, "owner_module", "")), "object_id": str(getattr(value, "object_id", "")), "version": _bounded_text(getattr(value, "version", None), 120), } def _validate_record(record: DsarRecordRef) -> None: if record.provider_id != "cases" or record.module_id != "cases": raise ValueError("Cases DSAR received a foreign provider record.") def _validate_action(action: DsarErasureActionRef) -> None: if action.provider_id != "cases" or action.module_id != "cases": raise ValueError("Cases DSAR received a foreign provider action.") def _record( resource_type: str, resource_id: str, category: str, title: str, data: Mapping[str, object], *, observed_at: datetime | None, immutable: bool = False, retention_reason: str | None = None, ) -> DsarRecordRef: return DsarRecordRef( provider_id="cases", module_id="cases", resource_type=resource_type, resource_id=resource_id, category=category, title=title, data=data, observed_at=observed_at, immutable_evidence=immutable, retention_reason=retention_reason, source_path="/cases", ) def _session(value: object) -> Session: if not isinstance(value, Session): raise TypeError("Cases DSAR provider requires a SQLAlchemy session.") return value def _bounded_text(value: str | None, limit: int) -> str | None: return value[:limit] if value else None def _normalized_id(value: object) -> str | None: if value is None: return None normalized = str(value).strip() return normalized or None def _iso(value: datetime | None) -> str | None: if value is None: return None if value.tzinfo is None: value = value.replace(tzinfo=timezone.utc) return value.isoformat() __all__ = ["CASES_DSAR_CAPABILITY", "CasesDsarProvider"]