From df51bb1787bb761163dc110ee11a7c31938b9fca Mon Sep 17 00:00:00 2001 From: Albrecht Degering Date: Thu, 6 Aug 2026 05:36:18 +0200 Subject: [PATCH] feat(records): govern lifecycle recovery and transfer evidence --- docs/EAKTE_ARCHITECTURE.md | 51 +- docs/RECORDS_DOMAIN_BOUNDARY.md | 20 +- src/govoplan_records/backend/archive.py | 91 + src/govoplan_records/backend/db/__init__.py | 6 + src/govoplan_records/backend/db/models.py | 203 ++ src/govoplan_records/backend/manifest.py | 145 +- .../7f5b3d9a2c1e_records_lifecycle.py | 291 +++ src/govoplan_records/backend/recovery.py | 203 ++ src/govoplan_records/backend/router.py | 424 +++- src/govoplan_records/backend/schemas.py | 124 ++ src/govoplan_records/backend/service.py | 1753 ++++++++++++++++- .../fixtures/service_to_decision_journey.json | 18 + tests/test_migrations.py | 5 +- tests/test_records.py | 421 +++- tests/test_recovery.py | 181 ++ webui/src/api/records.ts | 185 ++ .../features/records/RecordCatalogDialog.tsx | 174 ++ .../features/records/RecordLifecyclePanel.tsx | 383 ++++ webui/src/features/records/RecordsPage.tsx | 85 +- webui/src/i18n/generatedTranslations.ts | 150 +- webui/src/module.ts | 1 + webui/src/styles/records.css | 106 + 22 files changed, 4954 insertions(+), 66 deletions(-) create mode 100644 src/govoplan_records/backend/archive.py create mode 100644 src/govoplan_records/backend/migrations/versions/7f5b3d9a2c1e_records_lifecycle.py create mode 100644 src/govoplan_records/backend/recovery.py create mode 100644 tests/fixtures/service_to_decision_journey.json create mode 100644 tests/test_recovery.py create mode 100644 webui/src/features/records/RecordCatalogDialog.tsx create mode 100644 webui/src/features/records/RecordLifecyclePanel.tsx diff --git a/docs/EAKTE_ARCHITECTURE.md b/docs/EAKTE_ARCHITECTURE.md index cbd815d..96c616d 100644 --- a/docs/EAKTE_ARCHITECTURE.md +++ b/docs/EAKTE_ARCHITECTURE.md @@ -17,8 +17,8 @@ Implementation is tracked in ## Implementation Status -The native foundation is implemented through the work packages tracked by -Records #2-#4: +The native foundation and governed lifecycle are implemented through the work +packages tracked by Records #2-#6: - versioned file plans and record classes; - stable records, immutable revisions, volumes, exact record items, and @@ -29,13 +29,24 @@ Records #2-#4: temporal status, create/edit, and filing actions; - a provider-neutral Core filing contract with exact Files-version and Cases-revision providers. +- governed close/reopen transitions, class-bound retention calculation, + effective-dated holds, appraisal, evidence-bound disposition proposals, and + mandatory independent approval before finalization; +- archive-neutral transfer manifests and bounded receipts, plus an explicitly + non-conformant simulation provider that never claims custody; +- durable recovery-ledger fences for every API write, atomic domain/checkpoint + commits, Audit events, source-revision revalidation, transfer-manifest + checksum checks, and an operator recovery-evidence view; +- tenant administration for versioned file-plan nodes and record classes, + explicit volume management, and lifecycle controls in the eAkte workspace. -The remaining delivery order is intentionally visible rather than implied: -Records #5 owns closure, retention, holds, appraisal, and disposition; #6 owns -recovery and signed evidence; #7 requires selection and target testing of an -archive/xdomea provider; and #8 proves the reference journey. Restricted -per-record access grants also remain a dedicated access-policy slice. No -archive or destructive effect is currently claimed. +The remaining boundary is intentionally visible rather than implied: Records +#7 requires selection and target testing of an archive/xdomea endpoint and +conformance profile; #8 completes the cross-module reference journey and +production evidence. Restricted per-record access grants also remain a +dedicated access-policy slice. Destruction is represented only as an approved +pending state; no content deletion or real archive effect is currently +claimed. ## Ownership Boundary @@ -103,10 +114,15 @@ planned -> open -> closed -> retention_running -> appraisal_due destroyed ``` -Reopening creates a governed transition and does not reset elapsed retention -without an explicit rule. A hold preserves the reason and affected scope. A -destruction action requires exact manifest, current authority, policy, -approval, preflight, idempotency, outcome-unknown recovery, and evidence. +Reopening creates a governed transition and preserves the preceding retention +schedule. A later closure only restarts retention under an explicit +administrator action. A hold preserves reason, authority, scope, effective +interval, policy references, and release evidence. Holds block proposal, +approval finalization, packaging, and dispatch. A destruction approval changes +the record to `destruction_pending`; no source object or content is deleted. +An unapproved disposition can be withdrawn through a new immutable revision so +a corrected proposal can supersede it. An approved disposition cannot use this +correction path. ## Filing Semantics @@ -189,10 +205,17 @@ available in an evidence/details view. storage or an external provider, never node-local paths. - Every external filing, transfer, or destruction uses intent-before-effect, idempotency, durable receipts, outcome-unknown state, and reconciliation. +- Native Records API writes use a per-resource distributed lease and the Core + recovery ledger. Immutable revision, chronology, Audit projection, and the + terminal checkpoint commit atomically. - Backup evidence binds record rows, object manifests, provider mappings, policy/configuration versions, and key references. - Restore verifies content digests, missing keys/objects, provider reachability, and disposition holds before reopening effects. +- Automated evidence tests prove that terminal Records operations and their + hash chain remain verifiable after a database backup/restore round trip. An + outcome-unknown transfer is persisted as non-retryable and cannot advance + custody or record disposition state. - Search indexes are rebuildable projections and cannot become record authority. @@ -203,8 +226,10 @@ available in an evidence/details view. 2. Integrate filing from Cases, Forms Runtime, Decisions, Campaign/Postbox, Files, and Reporting. 3. Add closure, retention calculation, holds, appraisal, and reviewed - disposition without destructive provider effects. + disposition without destructive provider effects. **Implemented.** 4. Add native transfer packages and one target-tested xdomea/archive provider. + Native packaging and simulation are implemented; target selection/testing + remains external. 5. Add destruction/recovery, TR-ESOR provider integration, migration, and signed reference-journey evidence. diff --git a/docs/RECORDS_DOMAIN_BOUNDARY.md b/docs/RECORDS_DOMAIN_BOUNDARY.md index 240d78d..61314ec 100644 --- a/docs/RECORDS_DOMAIN_BOUNDARY.md +++ b/docs/RECORDS_DOMAIN_BOUNDARY.md @@ -42,16 +42,24 @@ The native kernel currently provides: file versions and Cases revisions; - tenant APIs, search projection, uninstall/retirement guards, and a Records workspace using shared WebUI controls. +- immutable close/reopen, retention, hold, appraisal, and independently + approved disposition transitions; +- archive-neutral manifests, bounded receipts, and a clearly marked transfer + simulation that does not claim custody; +- Core recovery-ledger adoption, atomic terminal evidence, Audit projection, + source and package restore diagnostics, catalog administration, volumes, and + lifecycle UI. -Restricted object grants and the lifecycle after the planned/open stages remain -separate governed slices. Archive transfer and destructive effects are not -implemented by this kernel. +Restricted object grants remain a separate governed slice. A target-tested +archive adapter and any destructive effect remain deliberately unimplemented; +approved destruction is only a pending lifecycle state. ## First Implementation Slice -Complete restricted access, closure and retention calculation, holds, -appraisal, disposition, transfer, recovery evidence, and one target-tested -archive provider without moving source-module ownership into Records. +Complete restricted access and one target-tested archive provider without +moving source-module ownership into Records. Extend the executable reference +journey across assisted service, decision, filing, hold, restore, search, and +archive simulation before claiming production maturity. The complete native/external boundary, temporal and purpose-aware record model, disposition lifecycle, German public-sector provider profiles, and staged diff --git a/src/govoplan_records/backend/archive.py b/src/govoplan_records/backend/archive.py new file mode 100644 index 0000000..8261c43 --- /dev/null +++ b/src/govoplan_records/backend/archive.py @@ -0,0 +1,91 @@ +from __future__ import annotations + +from datetime import UTC, datetime +import hashlib +import json + +from govoplan_core.core.records import ( + RecordArchiveProviderState, + RecordArchiveReceipt, + RecordArchiveTransferRequest, + RecordContractError, +) + + +SIMULATION_PROVIDER_ID = "simulation" +SIMULATION_PROFILE = "govoplan-simulation-v1" + + +class SimulatedRecordArchiveProvider: + """Exercise the transfer boundary without claiming archival custody.""" + + provider_id = SIMULATION_PROVIDER_ID + + def state(self) -> RecordArchiveProviderState: + checked_at = datetime.now(UTC) + return RecordArchiveProviderState( + provider_id=self.provider_id, + label="Records transfer simulation", + profiles=(SIMULATION_PROFILE,), + authority_modes=("linked_reference",), + healthy=True, + checked_at=checked_at, + last_success_at=checked_at, + freshness_seconds=0, + limitations=( + "Simulation validates package and receipt handling but does not transfer custody.", + "It is not an xDOMEA or archival conformance profile.", + ), + simulated=True, + ) + + def dispatch( + self, + session: object, + principal: object, + *, + request: RecordArchiveTransferRequest, + ) -> RecordArchiveReceipt: + del session, principal + if request.package.profile != SIMULATION_PROFILE: + raise RecordContractError( + "The simulation provider only accepts its declared profile." + ) + observed_at = datetime.now(UTC) + receipt_payload = { + "provider_id": self.provider_id, + "package_id": request.package.package_id, + "manifest_sha256": request.package.manifest_sha256, + "profile": request.package.profile, + "outcome": "accepted", + "simulated": True, + } + receipt_sha256 = hashlib.sha256( + json.dumps( + receipt_payload, + sort_keys=True, + separators=(",", ":"), + ).encode("utf-8") + ).hexdigest() + return RecordArchiveReceipt( + provider_id=self.provider_id, + package_id=request.package.package_id, + outcome="accepted", + observed_at=observed_at, + receipt_sha256=receipt_sha256, + external_reference=f"simulation:{request.package.package_id}", + retry_safe=True, + simulated=True, + metadata={ + "manifest_sha256": request.package.manifest_sha256, + "profile": request.package.profile, + "custody_transferred": False, + }, + ) + + +__all__ = [ + "SIMULATION_PROFILE", + "SIMULATION_PROVIDER_ID", + "SimulatedRecordArchiveProvider", +] diff --git a/src/govoplan_records/backend/db/__init__.py b/src/govoplan_records/backend/db/__init__.py index 963deb6..9278306 100644 --- a/src/govoplan_records/backend/db/__init__.py +++ b/src/govoplan_records/backend/db/__init__.py @@ -1,19 +1,25 @@ from govoplan_records.backend.db.models import ( RecordChronologyEntry, RecordClassRevision, + RecordDispositionRevision, RecordFilePlanRevision, + RecordHoldRevision, RecordIdentity, RecordItem, RecordRevision, + RecordTransferPackageRevision, RecordVolumeRevision, ) __all__ = [ "RecordChronologyEntry", "RecordClassRevision", + "RecordDispositionRevision", "RecordFilePlanRevision", + "RecordHoldRevision", "RecordIdentity", "RecordItem", "RecordRevision", + "RecordTransferPackageRevision", "RecordVolumeRevision", ] diff --git a/src/govoplan_records/backend/db/models.py b/src/govoplan_records/backend/db/models.py index aad4f9e..298c3e5 100644 --- a/src/govoplan_records/backend/db/models.py +++ b/src/govoplan_records/backend/db/models.py @@ -240,6 +240,24 @@ class RecordRevision(Base, TimestampMixin): changed_by: Mapped[str | None] = mapped_column( String(255), nullable=True, index=True ) + closed_at: Mapped[datetime | None] = mapped_column( + DateTime(timezone=True), nullable=True, index=True + ) + retention_started_at: Mapped[datetime | None] = mapped_column( + DateTime(timezone=True), nullable=True, index=True + ) + retention_due_at: Mapped[datetime | None] = mapped_column( + DateTime(timezone=True), nullable=True, index=True + ) + retention_rule: Mapped[dict[str, Any]] = mapped_column( + JSON, default=dict, nullable=False + ) + appraisal_state: Mapped[str | None] = mapped_column( + String(40), nullable=True, index=True + ) + appraisal: Mapped[dict[str, Any]] = mapped_column( + JSON, default=dict, nullable=False + ) snapshot: Mapped[dict[str, Any]] = mapped_column(JSON, default=dict, nullable=False) @@ -387,12 +405,197 @@ class RecordChronologyEntry(Base, TimestampMixin): payload: Mapped[dict[str, Any]] = mapped_column(JSON, default=dict, nullable=False) +class RecordHoldRevision(Base, TimestampMixin): + __tablename__ = "record_hold_revisions" + __table_args__ = ( + UniqueConstraint( + "tenant_id", "hold_id", "revision", name="uq_record_hold_revision" + ), + UniqueConstraint( + "tenant_id", "idempotency_key", name="uq_record_hold_idempotency" + ), + Index("ix_record_hold_current", "tenant_id", "hold_id", "superseded_at"), + Index("ix_record_hold_record", "tenant_id", "record_id", "status"), + ) + + id: Mapped[str] = mapped_column(String(36), primary_key=True, default=new_uuid) + tenant_id: Mapped[str] = mapped_column(String(36), nullable=False, index=True) + hold_id: Mapped[str] = mapped_column(String(255), nullable=False, index=True) + record_id: Mapped[str] = mapped_column(String(255), nullable=False, index=True) + revision: Mapped[int] = mapped_column(Integer, nullable=False) + previous_revision_id: Mapped[str | None] = mapped_column( + ForeignKey("record_hold_revisions.id", ondelete="RESTRICT"), + nullable=True, + index=True, + ) + status: Mapped[str] = mapped_column(String(30), nullable=False, index=True) + reason: Mapped[str] = mapped_column(Text, nullable=False) + authority: Mapped[str] = mapped_column(String(500), nullable=False) + scope: Mapped[dict[str, Any]] = mapped_column(JSON, default=dict, nullable=False) + effective_from: Mapped[datetime] = mapped_column( + DateTime(timezone=True), nullable=False, index=True + ) + effective_to: Mapped[datetime | None] = mapped_column( + DateTime(timezone=True), nullable=True, index=True + ) + released_at: Mapped[datetime | None] = mapped_column( + DateTime(timezone=True), nullable=True + ) + policy_refs: Mapped[list[str]] = mapped_column(JSON, default=list, nullable=False) + institutional_context: Mapped[dict[str, Any]] = mapped_column( + JSON, default=dict, nullable=False + ) + recorded_at: Mapped[datetime] = mapped_column( + DateTime(timezone=True), nullable=False, index=True + ) + superseded_at: Mapped[datetime | None] = mapped_column( + DateTime(timezone=True), nullable=True, index=True + ) + changed_by: Mapped[str | None] = mapped_column( + String(255), nullable=True, index=True + ) + idempotency_key: Mapped[str] = mapped_column(String(255), nullable=False) + request_sha256: Mapped[str] = mapped_column(String(64), nullable=False) + + +class RecordDispositionRevision(Base, TimestampMixin): + __tablename__ = "record_disposition_revisions" + __table_args__ = ( + UniqueConstraint( + "tenant_id", + "disposition_id", + "revision", + name="uq_record_disposition_revision", + ), + UniqueConstraint( + "tenant_id", "idempotency_key", name="uq_record_disposition_idempotency" + ), + Index( + "ix_record_disposition_current", + "tenant_id", + "disposition_id", + "superseded_at", + ), + Index("ix_record_disposition_record", "tenant_id", "record_id", "status"), + ) + + id: Mapped[str] = mapped_column(String(36), primary_key=True, default=new_uuid) + tenant_id: Mapped[str] = mapped_column(String(36), nullable=False, index=True) + disposition_id: Mapped[str] = mapped_column(String(255), nullable=False, index=True) + record_id: Mapped[str] = mapped_column(String(255), nullable=False, index=True) + revision: Mapped[int] = mapped_column(Integer, nullable=False) + previous_revision_id: Mapped[str | None] = mapped_column( + ForeignKey("record_disposition_revisions.id", ondelete="RESTRICT"), + nullable=True, + index=True, + ) + action: Mapped[str] = mapped_column(String(30), nullable=False, index=True) + status: Mapped[str] = mapped_column(String(40), nullable=False, index=True) + reason: Mapped[str] = mapped_column(Text, nullable=False) + subject_revision: Mapped[int] = mapped_column(Integer, nullable=False) + subject_sha256: Mapped[str] = mapped_column(String(64), nullable=False) + consequence_preview: Mapped[dict[str, Any]] = mapped_column( + JSON, default=dict, nullable=False + ) + policy_refs: Mapped[list[str]] = mapped_column(JSON, default=list, nullable=False) + approval_request_id: Mapped[str | None] = mapped_column( + String(255), nullable=True, index=True + ) + proposed_by: Mapped[str | None] = mapped_column( + String(255), nullable=True, index=True + ) + reviewed_by: Mapped[str | None] = mapped_column( + String(255), nullable=True, index=True + ) + reviewed_at: Mapped[datetime | None] = mapped_column( + DateTime(timezone=True), nullable=True + ) + institutional_context: Mapped[dict[str, Any]] = mapped_column( + JSON, default=dict, nullable=False + ) + recorded_at: Mapped[datetime] = mapped_column( + DateTime(timezone=True), nullable=False, index=True + ) + superseded_at: Mapped[datetime | None] = mapped_column( + DateTime(timezone=True), nullable=True, index=True + ) + idempotency_key: Mapped[str] = mapped_column(String(255), nullable=False) + request_sha256: Mapped[str] = mapped_column(String(64), nullable=False) + + +class RecordTransferPackageRevision(Base, TimestampMixin): + __tablename__ = "record_transfer_package_revisions" + __table_args__ = ( + UniqueConstraint( + "tenant_id", + "package_id", + "revision", + name="uq_record_transfer_package_revision", + ), + UniqueConstraint( + "tenant_id", + "idempotency_key", + name="uq_record_transfer_package_idempotency", + ), + Index( + "ix_record_transfer_package_current", + "tenant_id", + "package_id", + "superseded_at", + ), + Index("ix_record_transfer_package_record", "tenant_id", "record_id", "status"), + ) + + id: Mapped[str] = mapped_column(String(36), primary_key=True, default=new_uuid) + tenant_id: Mapped[str] = mapped_column(String(36), nullable=False, index=True) + package_id: Mapped[str] = mapped_column(String(255), nullable=False, index=True) + record_id: Mapped[str] = mapped_column(String(255), nullable=False, index=True) + disposition_id: Mapped[str] = mapped_column(String(255), nullable=False, index=True) + revision: Mapped[int] = mapped_column(Integer, nullable=False) + previous_revision_id: Mapped[str | None] = mapped_column( + ForeignKey("record_transfer_package_revisions.id", ondelete="RESTRICT"), + nullable=True, + index=True, + ) + record_revision: Mapped[int] = mapped_column(Integer, nullable=False) + provider_id: Mapped[str] = mapped_column(String(100), nullable=False, index=True) + profile: Mapped[str] = mapped_column(String(255), nullable=False) + status: Mapped[str] = mapped_column(String(40), nullable=False, index=True) + authority_mode: Mapped[str] = mapped_column(String(40), nullable=False) + manifest: Mapped[dict[str, Any]] = mapped_column(JSON, default=dict, nullable=False) + manifest_sha256: Mapped[str] = mapped_column(String(64), nullable=False) + receipt: Mapped[dict[str, Any]] = mapped_column(JSON, default=dict, nullable=False) + receipt_sha256: Mapped[str | None] = mapped_column(String(64), nullable=True) + external_reference: Mapped[str | None] = mapped_column(String(1500), nullable=True) + recovery_operation_id: Mapped[str | None] = mapped_column( + String(36), nullable=True, index=True + ) + simulated: Mapped[bool] = mapped_column(Boolean, default=False, nullable=False) + institutional_context: Mapped[dict[str, Any]] = mapped_column( + JSON, default=dict, nullable=False + ) + recorded_at: Mapped[datetime] = mapped_column( + DateTime(timezone=True), nullable=False, index=True + ) + superseded_at: Mapped[datetime | None] = mapped_column( + DateTime(timezone=True), nullable=True, index=True + ) + changed_by: Mapped[str | None] = mapped_column( + String(255), nullable=True, index=True + ) + idempotency_key: Mapped[str] = mapped_column(String(255), nullable=False) + request_sha256: Mapped[str] = mapped_column(String(64), nullable=False) + + __all__ = [ "RecordChronologyEntry", "RecordClassRevision", "RecordFilePlanRevision", + "RecordHoldRevision", "RecordIdentity", "RecordItem", + "RecordDispositionRevision", "RecordRevision", + "RecordTransferPackageRevision", "RecordVolumeRevision", ] diff --git a/src/govoplan_records/backend/manifest.py b/src/govoplan_records/backend/manifest.py index dd65f18..13052fe 100644 --- a/src/govoplan_records/backend/manifest.py +++ b/src/govoplan_records/backend/manifest.py @@ -29,13 +29,20 @@ from govoplan_core.core.modules import ( RoleTemplate, ) from govoplan_core.core.provider_governance import declared_module_architecture -from govoplan_core.core.records import CAPABILITY_RECORDS_FILING +from govoplan_core.core.records import ( + CAPABILITY_RECORDS_FILING, + record_archive_capability, +) from govoplan_core.core.search import SearchSourceProviderRegistration from govoplan_core.core.views import ViewSurface from govoplan_core.db.base import Base from govoplan_records.backend.db import models as record_models from govoplan_records.backend.search_source import create_records_search_source from govoplan_records.backend.service import SqlRecordRegistry +from govoplan_records.backend.archive import ( + SIMULATION_PROVIDER_ID, + SimulatedRecordArchiveProvider, +) MODULE_ID = "records" @@ -55,6 +62,7 @@ OPTIONAL_DEPENDENCIES = ( "dms", "docs", "policy", + "approvals", "audit", "transparency", "search", @@ -85,6 +93,13 @@ def _records_registry(context: ModuleContext) -> SqlRecordRegistry: return SqlRecordRegistry(context.registry) +def _simulated_archive_provider( + context: ModuleContext, +) -> SimulatedRecordArchiveProvider: + del context + return SimulatedRecordArchiveProvider() + + def _tenant_summary(session, tenant_id: str) -> dict[str, int]: records = ( session.query(record_models.RecordIdentity) @@ -100,7 +115,32 @@ def _tenant_summary(session, tenant_id: str) -> dict[str, int]: ) .count() ) - return {"records": records, "open_records": open_records} + active_holds = ( + session.query(record_models.RecordHoldRevision) + .filter( + record_models.RecordHoldRevision.tenant_id == tenant_id, + record_models.RecordHoldRevision.superseded_at.is_(None), + record_models.RecordHoldRevision.status == "active", + ) + .count() + ) + pending_dispositions = ( + session.query(record_models.RecordDispositionRevision) + .filter( + record_models.RecordDispositionRevision.tenant_id == tenant_id, + record_models.RecordDispositionRevision.superseded_at.is_(None), + record_models.RecordDispositionRevision.status.in_( + ("review_pending", "review_unavailable") + ), + ) + .count() + ) + return { + "records": records, + "open_records": open_records, + "active_holds": active_holds, + "pending_dispositions": pending_dispositions, + } PERMISSIONS = ( @@ -247,14 +287,16 @@ DOCUMENTATION = ( }, ), DocumentationTopic( - id="records.lifecycle-limitations", - title="Records lifecycle limitations", - summary="Identifies lifecycle controls intentionally deferred beyond the native kernel.", + id="records.lifecycle", + title="Governed records lifecycle", + summary="Close, retain, hold, appraise, approve, and package records with durable evidence.", body=( - "The current vertical supports planned and open records. Restricted object grants, closure, " - "retention calculation, holds, appraisal, disposition, transfer, destruction, and external " - "archive effects are separate governed work packages. No destructive effect is implied by " - "enabling Records." + "Closing a record applies the versioned record-class retention rule. Reopening preserves the " + "previous schedule unless an administrator deliberately restarts it. Effective-dated holds block " + "disposition and transfer. Appraisal selects retain, transfer, destroy, or reclassify; a disposition " + "binds the exact evidence digest and requires independent approval through Approvals. Approval only " + "changes lifecycle state. Destruction remains pending and transfer simulation never claims custody. " + "Every API mutation is fenced and recorded in the Core recovery ledger." ), layer="configured", documentation_types=("admin", "user"), @@ -263,17 +305,44 @@ DOCUMENTATION = ( related_modules=("policy", "approvals", "audit", "dms"), translations={ "de": { - "title": "Grenzen des Aktenlebenszyklus", - "summary": "Kennzeichnet bewusst nach dem nativen Kern umzusetzende Lebenszyklussteuerungen.", + "title": "Gesteuerter Aktenlebenszyklus", + "summary": "Akten mit dauerhaftem Nachweis abschließen, aufbewahren, sperren, bewerten, freigeben und paketieren.", "body": ( - "Der aktuelle Stand unterstützt geplante und offene Akten. Objektbezogene Freigaben, " - "Abschluss, Aufbewahrungsberechnung, Sperren, Bewertung, Aussonderung, Übergabe, Vernichtung " - "und externe Archiveffekte sind getrennte gesteuerte Arbeitspakete. Die Aktivierung von " - "Records löst keine vernichtende Wirkung aus." + "Beim Abschluss einer Akte wird die versionierte Aufbewahrungsregel der Aktenklasse angewendet. " + "Eine Wiedereröffnung bewahrt den bisherigen Zeitplan, sofern ein Administrator ihn nicht " + "bewusst neu startet. Gültigkeitsbezogene Sperren blockieren Aussonderung und Übergabe. Die " + "Bewertung wählt Aufbewahrung, Übergabe, Vernichtung oder Neuklassifikation; die Aussonderung " + "bindet den exakten Nachweis und benötigt eine unabhängige Freigabe durch Approvals. Eine " + "Freigabe ändert nur den Lebenszyklusstatus. Vernichtung bleibt vorgemerkt, und eine Simulation " + "behauptet keine Archivverwahrung. Jede API-Änderung wird im Recovery-Ledger abgesichert." ), } }, - metadata={"known_limit": True}, + metadata={ + "help_contexts": [ + "records.lifecycle", + "records.lifecycle.volume", + "records.lifecycle.close", + "records.lifecycle.reopen", + "records.lifecycle.appraise", + "records.lifecycle.hold", + "records.lifecycle.release-hold", + "records.lifecycle.disposition", + "records.lifecycle.finalize", + "records.lifecycle.prepare-transfer", + "records.lifecycle.dispatch-transfer", + "records.lifecycle.recovery", + "records.catalog.admin", + "records.field.volume", + "records.field.volume-label", + "records.field.retention-trigger", + "records.field.disposition", + "records.field.hold-authority", + "records.field.lifecycle-reason", + "records.field.archive-provider", + "records.field.archive-profile", + ], + }, ), ) @@ -361,14 +430,25 @@ manifest = ModuleManifest( provides_interfaces=( ModuleInterfaceProvider(name="records.registry", version="1.0.0"), ModuleInterfaceProvider(name="records.filing", version="1.0.0"), + ModuleInterfaceProvider(name="records.archive", version="1.0.0"), ), - capability_factories={CAPABILITY_RECORDS_FILING: _records_registry}, + capability_factories={ + CAPABILITY_RECORDS_FILING: _records_registry, + record_archive_capability(SIMULATION_PROVIDER_ID): _simulated_archive_provider, + }, capability_documentation={ CAPABILITY_RECORDS_FILING: CapabilityDocumentation( label="Record filing", summary="Resolves authorized exact source revisions and files immutable record items.", contract_version="1.0.0", ), + record_archive_capability(SIMULATION_PROVIDER_ID): CapabilityDocumentation( + label="Record archive simulation", + summary=( + "Validates archive-neutral package and receipt handling without transferring custody." + ), + contract_version="1.0.0", + ), }, migration_spec=MigrationSpec( module_id=MODULE_ID, @@ -377,6 +457,9 @@ manifest = ModuleManifest( retirement_supported=True, retirement_provider=drop_table_retirement_provider( record_models.RecordChronologyEntry, + record_models.RecordTransferPackageRevision, + record_models.RecordDispositionRevision, + record_models.RecordHoldRevision, record_models.RecordItem, record_models.RecordVolumeRevision, record_models.RecordRevision, @@ -396,6 +479,9 @@ manifest = ModuleManifest( record_models.RecordRevision, record_models.RecordItem, record_models.RecordChronologyEntry, + record_models.RecordHoldRevision, + record_models.RecordDispositionRevision, + record_models.RecordTransferPackageRevision, record_models.RecordClassRevision, record_models.RecordFilePlanRevision, label="Records", @@ -417,6 +503,9 @@ manifest = ModuleManifest( "record_item", "record_class", "file_plan_node", + "record_hold", + "record_disposition", + "record_transfer_package", ), evidence=( "src/govoplan_records/backend/service.py", @@ -433,10 +522,16 @@ manifest = ModuleManifest( ), ), retention=InformationGovernanceDimension( - adoption="contract_only", - limitation=( - "Record classes preserve retention inputs; closure, holds, calculation, appraisal, " - "and disposition are tracked in Records #5." + adoption="enforced", + object_types=( + "record", + "record_class", + "record_hold", + "record_disposition", + ), + evidence=( + "src/govoplan_records/backend/service.py", + "tests/test_records.py", ), ), institutional_context=InformationGovernanceDimension( @@ -462,8 +557,8 @@ manifest = ModuleManifest( documentation_ref="docs/EAKTE_ARCHITECTURE.md", test_ref="tests/test_records.py", known_limits=( - "Restricted object grants and lifecycle stages after open are tracked separately.", - "Archive transfer and destructive effects are not part of the native kernel.", + "Restricted object grants are tracked separately.", + "Archive packaging is native, but real target conformance and destructive effects remain external work.", ), supported_authority_modes=( "native_authoritative", @@ -481,6 +576,10 @@ manifest = ModuleManifest( "record item", "filing decision", "record chronology", + "record hold", + "record appraisal", + "record disposition", + "record transfer package", ), non_owned_concepts=( "file content", diff --git a/src/govoplan_records/backend/migrations/versions/7f5b3d9a2c1e_records_lifecycle.py b/src/govoplan_records/backend/migrations/versions/7f5b3d9a2c1e_records_lifecycle.py new file mode 100644 index 0000000..154d84b --- /dev/null +++ b/src/govoplan_records/backend/migrations/versions/7f5b3d9a2c1e_records_lifecycle.py @@ -0,0 +1,291 @@ +"""Add governed Records lifecycle evidence. + +Revision ID: 7f5b3d9a2c1e +Revises: 6e4a2c8f1d9b +""" + +from __future__ import annotations + +from alembic import op +import sqlalchemy as sa + + +revision = "7f5b3d9a2c1e" +down_revision = "6e4a2c8f1d9b" +branch_labels = None +depends_on = None + + +def upgrade() -> None: + with op.batch_alter_table("record_revisions") as batch_op: + batch_op.add_column( + sa.Column("closed_at", sa.DateTime(timezone=True), nullable=True) + ) + batch_op.add_column( + sa.Column("retention_started_at", sa.DateTime(timezone=True), nullable=True) + ) + batch_op.add_column( + sa.Column("retention_due_at", sa.DateTime(timezone=True), nullable=True) + ) + batch_op.add_column( + sa.Column( + "retention_rule", + sa.JSON(), + nullable=False, + server_default=sa.text("'{}'"), + ) + ) + batch_op.add_column( + sa.Column("appraisal_state", sa.String(length=40), nullable=True) + ) + batch_op.add_column( + sa.Column( + "appraisal", + sa.JSON(), + nullable=False, + server_default=sa.text("'{}'"), + ) + ) + batch_op.create_index( + op.f("ix_record_revisions_closed_at"), ["closed_at"], unique=False + ) + batch_op.create_index( + op.f("ix_record_revisions_retention_started_at"), + ["retention_started_at"], + unique=False, + ) + batch_op.create_index( + op.f("ix_record_revisions_retention_due_at"), + ["retention_due_at"], + unique=False, + ) + batch_op.create_index( + op.f("ix_record_revisions_appraisal_state"), + ["appraisal_state"], + unique=False, + ) + + op.create_table( + "record_hold_revisions", + sa.Column("id", sa.String(length=36), nullable=False), + sa.Column("tenant_id", sa.String(length=36), nullable=False), + sa.Column("hold_id", sa.String(length=255), nullable=False), + sa.Column("record_id", sa.String(length=255), nullable=False), + sa.Column("revision", sa.Integer(), nullable=False), + sa.Column("previous_revision_id", sa.String(length=36), nullable=True), + sa.Column("status", sa.String(length=30), nullable=False), + sa.Column("reason", sa.Text(), nullable=False), + sa.Column("authority", sa.String(length=500), nullable=False), + sa.Column("scope", sa.JSON(), nullable=False), + sa.Column("effective_from", sa.DateTime(timezone=True), nullable=False), + sa.Column("effective_to", sa.DateTime(timezone=True), nullable=True), + sa.Column("released_at", sa.DateTime(timezone=True), nullable=True), + sa.Column("policy_refs", sa.JSON(), nullable=False), + sa.Column("institutional_context", sa.JSON(), nullable=False), + sa.Column("recorded_at", sa.DateTime(timezone=True), nullable=False), + sa.Column("superseded_at", sa.DateTime(timezone=True), nullable=True), + sa.Column("changed_by", sa.String(length=255), nullable=True), + sa.Column("idempotency_key", sa.String(length=255), nullable=False), + sa.Column("request_sha256", sa.String(length=64), nullable=False), + sa.Column("created_at", sa.DateTime(timezone=True), nullable=False), + sa.Column("updated_at", sa.DateTime(timezone=True), nullable=False), + sa.ForeignKeyConstraint( + ["previous_revision_id"], + ["record_hold_revisions.id"], + ondelete="RESTRICT", + ), + sa.PrimaryKeyConstraint("id"), + sa.UniqueConstraint( + "tenant_id", "hold_id", "revision", name="uq_record_hold_revision" + ), + sa.UniqueConstraint( + "tenant_id", "idempotency_key", name="uq_record_hold_idempotency" + ), + ) + _indexes( + "record_hold_revisions", + "tenant_id", + "hold_id", + "record_id", + "previous_revision_id", + "status", + "effective_from", + "effective_to", + "recorded_at", + "superseded_at", + "changed_by", + ) + op.create_index( + "ix_record_hold_current", + "record_hold_revisions", + ["tenant_id", "hold_id", "superseded_at"], + ) + op.create_index( + "ix_record_hold_record", + "record_hold_revisions", + ["tenant_id", "record_id", "status"], + ) + + op.create_table( + "record_disposition_revisions", + sa.Column("id", sa.String(length=36), nullable=False), + sa.Column("tenant_id", sa.String(length=36), nullable=False), + sa.Column("disposition_id", sa.String(length=255), nullable=False), + sa.Column("record_id", sa.String(length=255), nullable=False), + sa.Column("revision", sa.Integer(), nullable=False), + sa.Column("previous_revision_id", sa.String(length=36), nullable=True), + sa.Column("action", sa.String(length=30), nullable=False), + sa.Column("status", sa.String(length=40), nullable=False), + sa.Column("reason", sa.Text(), nullable=False), + sa.Column("subject_revision", sa.Integer(), nullable=False), + sa.Column("subject_sha256", sa.String(length=64), nullable=False), + sa.Column("consequence_preview", sa.JSON(), nullable=False), + sa.Column("policy_refs", sa.JSON(), nullable=False), + sa.Column("approval_request_id", sa.String(length=255), nullable=True), + sa.Column("proposed_by", sa.String(length=255), nullable=True), + sa.Column("reviewed_by", sa.String(length=255), nullable=True), + sa.Column("reviewed_at", sa.DateTime(timezone=True), nullable=True), + sa.Column("institutional_context", sa.JSON(), nullable=False), + sa.Column("recorded_at", sa.DateTime(timezone=True), nullable=False), + sa.Column("superseded_at", sa.DateTime(timezone=True), nullable=True), + sa.Column("idempotency_key", sa.String(length=255), nullable=False), + sa.Column("request_sha256", sa.String(length=64), nullable=False), + sa.Column("created_at", sa.DateTime(timezone=True), nullable=False), + sa.Column("updated_at", sa.DateTime(timezone=True), nullable=False), + sa.ForeignKeyConstraint( + ["previous_revision_id"], + ["record_disposition_revisions.id"], + ondelete="RESTRICT", + ), + sa.PrimaryKeyConstraint("id"), + sa.UniqueConstraint( + "tenant_id", + "disposition_id", + "revision", + name="uq_record_disposition_revision", + ), + sa.UniqueConstraint( + "tenant_id", + "idempotency_key", + name="uq_record_disposition_idempotency", + ), + ) + _indexes( + "record_disposition_revisions", + "tenant_id", + "disposition_id", + "record_id", + "previous_revision_id", + "action", + "status", + "approval_request_id", + "proposed_by", + "reviewed_by", + "recorded_at", + "superseded_at", + ) + op.create_index( + "ix_record_disposition_current", + "record_disposition_revisions", + ["tenant_id", "disposition_id", "superseded_at"], + ) + op.create_index( + "ix_record_disposition_record", + "record_disposition_revisions", + ["tenant_id", "record_id", "status"], + ) + + op.create_table( + "record_transfer_package_revisions", + sa.Column("id", sa.String(length=36), nullable=False), + sa.Column("tenant_id", sa.String(length=36), nullable=False), + sa.Column("package_id", sa.String(length=255), nullable=False), + sa.Column("record_id", sa.String(length=255), nullable=False), + sa.Column("disposition_id", sa.String(length=255), nullable=False), + sa.Column("revision", sa.Integer(), nullable=False), + sa.Column("previous_revision_id", sa.String(length=36), nullable=True), + sa.Column("record_revision", sa.Integer(), nullable=False), + sa.Column("provider_id", sa.String(length=100), nullable=False), + sa.Column("profile", sa.String(length=255), nullable=False), + sa.Column("status", sa.String(length=40), nullable=False), + sa.Column("authority_mode", sa.String(length=40), nullable=False), + sa.Column("manifest", sa.JSON(), nullable=False), + sa.Column("manifest_sha256", sa.String(length=64), nullable=False), + sa.Column("receipt", sa.JSON(), nullable=False), + sa.Column("receipt_sha256", sa.String(length=64), nullable=True), + sa.Column("external_reference", sa.String(length=1500), nullable=True), + sa.Column("recovery_operation_id", sa.String(length=36), nullable=True), + sa.Column("simulated", sa.Boolean(), nullable=False), + sa.Column("institutional_context", sa.JSON(), nullable=False), + sa.Column("recorded_at", sa.DateTime(timezone=True), nullable=False), + sa.Column("superseded_at", sa.DateTime(timezone=True), nullable=True), + sa.Column("changed_by", sa.String(length=255), nullable=True), + sa.Column("idempotency_key", sa.String(length=255), nullable=False), + sa.Column("request_sha256", sa.String(length=64), nullable=False), + sa.Column("created_at", sa.DateTime(timezone=True), nullable=False), + sa.Column("updated_at", sa.DateTime(timezone=True), nullable=False), + sa.ForeignKeyConstraint( + ["previous_revision_id"], + ["record_transfer_package_revisions.id"], + ondelete="RESTRICT", + ), + sa.PrimaryKeyConstraint("id"), + sa.UniqueConstraint( + "tenant_id", + "package_id", + "revision", + name="uq_record_transfer_package_revision", + ), + sa.UniqueConstraint( + "tenant_id", + "idempotency_key", + name="uq_record_transfer_package_idempotency", + ), + ) + _indexes( + "record_transfer_package_revisions", + "tenant_id", + "package_id", + "record_id", + "disposition_id", + "previous_revision_id", + "provider_id", + "status", + "recovery_operation_id", + "recorded_at", + "superseded_at", + "changed_by", + ) + op.create_index( + "ix_record_transfer_package_current", + "record_transfer_package_revisions", + ["tenant_id", "package_id", "superseded_at"], + ) + op.create_index( + "ix_record_transfer_package_record", + "record_transfer_package_revisions", + ["tenant_id", "record_id", "status"], + ) + + +def downgrade() -> None: + op.drop_table("record_transfer_package_revisions") + op.drop_table("record_disposition_revisions") + op.drop_table("record_hold_revisions") + + with op.batch_alter_table("record_revisions") as batch_op: + batch_op.drop_index(op.f("ix_record_revisions_appraisal_state")) + batch_op.drop_index(op.f("ix_record_revisions_retention_due_at")) + batch_op.drop_index(op.f("ix_record_revisions_retention_started_at")) + batch_op.drop_index(op.f("ix_record_revisions_closed_at")) + batch_op.drop_column("appraisal") + batch_op.drop_column("appraisal_state") + batch_op.drop_column("retention_rule") + batch_op.drop_column("retention_due_at") + batch_op.drop_column("retention_started_at") + batch_op.drop_column("closed_at") + + +def _indexes(table: str, *columns: str) -> None: + for column in columns: + op.create_index(op.f(f"ix_{table}_{column}"), table, [column], unique=False) diff --git a/src/govoplan_records/backend/recovery.py b/src/govoplan_records/backend/recovery.py new file mode 100644 index 0000000..868b196 --- /dev/null +++ b/src/govoplan_records/backend/recovery.py @@ -0,0 +1,203 @@ +from __future__ import annotations + +from dataclasses import dataclass +from contextvars import ContextVar, Token +import hashlib +import json +from typing import Any + +from sqlalchemy.orm import Session, sessionmaker + +from govoplan_core.core.recovery import ( + RecoveryGuaranteeError, + RecoveryMode, + RecoveryPlan, +) +from govoplan_core.core.recovery_runtime import ( + DurableRecoveryOperation, + RecoveryOperationBusy, + RecoveryOperationStateConflict, + begin_durable_recovery_operation, +) +from govoplan_core.core.runtime_coordination import process_runtime_identity + + +class RecordRecoveryError(RuntimeError): + pass + + +_recovery_operation_id: ContextVar[str | None] = ContextVar( + "records_recovery_operation_id", default=None +) + + +def bind_record_recovery_operation(operation_id: str) -> Token[str | None]: + return _recovery_operation_id.set(operation_id) + + +def current_record_recovery_operation() -> str | None: + return _recovery_operation_id.get() + + +def reset_record_recovery_operation(token: Token[str | None]) -> None: + _recovery_operation_id.reset(token) + + +def record_session_factory(session: Session) -> sessionmaker[Session]: + bind = session.get_bind() + if bind is None: + raise RecordRecoveryError("Records recovery requires a bound database session.") + return sessionmaker(bind=bind, expire_on_commit=False) + + +@dataclass(slots=True) +class RecordAtomicRecovery: + operation: DurableRecoveryOperation | None + operation_id: str + replayed: bool + + def commit_success( + self, + session: Session, + *, + result: object, + resource_id: str, + ) -> None: + evidence = { + "verified": True, + "resource_id": resource_id, + "result_sha256": _canonical_sha256(result), + "domain_and_checkpoint_atomic": True, + "checks": { + "domain_result_digest_recorded": True, + "domain_and_checkpoint_atomic": True, + }, + } + if self.operation is None: + session.commit() + return + self.operation.commit_atomic_success(session, evidence=evidence) + + def reject(self, *, summary: str, error_type: str) -> None: + if self.operation is None: + return + self.operation.reject( + summary=summary, + evidence={ + "verified": True, + "domain_mutation_committed": False, + "error_type": error_type, + "checks": {"definitive_domain_rejection": True}, + }, + ) + + def fail(self, *, summary: str, error_type: str) -> None: + if self.operation is None: + return + self.operation.fail( + summary=summary, + evidence={ + "verified": True, + "domain_mutation_committed": False, + "error_type": error_type, + }, + ) + + +def begin_record_atomic_recovery( + session: Session, + *, + tenant_id: str, + operation_type: str, + idempotency_key: str, + request: dict[str, Any], + resource_type: str, + resource_id: str, +) -> RecordAtomicRecovery: + operation_key = hashlib.sha256( + f"{tenant_id}:{operation_type}:{idempotency_key}".encode("utf-8") + ).hexdigest() + lease_id = hashlib.sha256( + f"{tenant_id}:{resource_type}:{resource_id}".encode("utf-8") + ).hexdigest() + try: + started = begin_durable_recovery_operation( + record_session_factory(session), + identity=process_runtime_identity(), + module_id="records", + operation_type=operation_type, + idempotency_key=f"records:{operation_key}", + request={"tenant_id": tenant_id, **request}, + recovery_plan=RecoveryPlan( + mode=RecoveryMode.ATOMIC, + preconditions=( + "the current actor is authorized for the Records mutation", + "the expected revision and idempotency key are present", + "the target resource has no unresolved recovery operation", + ), + verification_steps=( + "commit the immutable domain revision and chronology entry", + "commit the terminal recovery checkpoint in the same transaction", + "verify the result digest and recovery evidence chain", + ), + ), + precondition_evidence={ + "tenant_id": tenant_id, + "resource_type": resource_type, + "resource_id_sha256": hashlib.sha256( + resource_id.encode("utf-8") + ).hexdigest(), + "request_sha256": _canonical_sha256(request), + "external_effect": False, + }, + lease_resource_key=f"records:{tenant_id}:{lease_id}", + lease_ttl_seconds=5 * 60, + resource_type=resource_type, + resource_id=resource_id, + metadata={ + "resources": ["postgresql"], + "external_effect": False, + "recovery_declaration": "records-atomic-revision", + }, + block_unresolved_resource=True, + ) + except RecoveryOperationBusy as exc: + raise RecordRecoveryError( + "Another runtime is changing this Records resource." + ) from exc + except RecoveryOperationStateConflict as exc: + raise RecordRecoveryError( + "This Records resource has an active or unresolved recovery operation." + ) from exc + except (RecoveryGuaranteeError, RuntimeError) as exc: + raise RecordRecoveryError( + "The recovery ledger is unavailable; the Records mutation was not started." + ) from exc + return RecordAtomicRecovery( + operation=started.operation, + operation_id=started.operation_id, + replayed=started.replayed, + ) + + +def _canonical_sha256(value: object) -> str: + return hashlib.sha256( + json.dumps( + value, + sort_keys=True, + separators=(",", ":"), + ensure_ascii=False, + default=str, + ).encode("utf-8") + ).hexdigest() + + +__all__ = [ + "RecordAtomicRecovery", + "RecordRecoveryError", + "begin_record_atomic_recovery", + "bind_record_recovery_operation", + "current_record_recovery_operation", + "record_session_factory", + "reset_record_recovery_operation", +] diff --git a/src/govoplan_records/backend/router.py b/src/govoplan_records/backend/router.py index 1748d89..391afd2 100644 --- a/src/govoplan_records/backend/router.py +++ b/src/govoplan_records/backend/router.py @@ -1,5 +1,6 @@ from __future__ import annotations +import hashlib from typing import Any from fastapi import APIRouter, Depends, HTTPException, Query, status @@ -7,18 +8,31 @@ from sqlalchemy.exc import IntegrityError from sqlalchemy.orm import Session from govoplan_core.auth import ApiPrincipal, get_api_principal, has_scope +from govoplan_core.audit.logging import audit_from_principal from govoplan_core.core.records import RecordFilingRequest, RecordSourceLocator from govoplan_core.db.session import get_session from govoplan_records.backend.manifest import ADMIN_SCOPE, READ_SCOPE, WRITE_SCOPE from govoplan_records.backend.schemas import ( FilePlanNodeWriteRequest, + RecordAppraisalRequest, + RecordArchiveProviderResponse, RecordCatalogResponse, RecordClassWriteRequest, + RecordCloseRequest, RecordCreateRequest, RecordDetailResponse, + RecordDispositionCreateRequest, + RecordDispositionFinalizeRequest, + RecordDispositionWithdrawRequest, + RecordHoldCreateRequest, + RecordHoldReleaseRequest, RecordItemCreateRequest, + RecordLifecycleActionRequest, RecordListResponse, + RecordRecoveryStatusResponse, RecordSourceProviderResponse, + RecordTransferDispatchRequest, + RecordTransferPackageCreateRequest, RecordUpdateRequest, RecordVolumeCreateRequest, ) @@ -29,6 +43,12 @@ from govoplan_records.backend.service import ( RecordStoreError, SqlRecordRegistry, ) +from govoplan_records.backend.recovery import ( + RecordRecoveryError, + begin_record_atomic_recovery, + bind_record_recovery_operation, + reset_record_recovery_operation, +) def create_router(registry: object | None = None) -> APIRouter: @@ -59,6 +79,12 @@ def create_router(registry: object | None = None) -> APIRouter: lambda: records.write_file_plan_node( session, principal, payload=payload.model_dump(mode="python") ), + principal=principal, + operation_type="catalog.file_plan.write", + idempotency_key=payload.idempotency_key, + request=payload.model_dump(mode="json"), + resource_type="record_file_plan_node", + resource_id=payload.node_id, ) @router.post( @@ -77,6 +103,12 @@ def create_router(registry: object | None = None) -> APIRouter: lambda: records.write_record_class( session, principal, payload=payload.model_dump(mode="python") ), + principal=principal, + operation_type="catalog.class.write", + idempotency_key=payload.idempotency_key, + request=payload.model_dump(mode="json"), + resource_type="record_class", + resource_id=payload.class_id, ) @router.get("/sources", response_model=RecordSourceProviderResponse) @@ -89,6 +121,16 @@ def create_router(registry: object | None = None) -> APIRouter: providers=records.source_providers(session, principal) ) + @router.get("/archive-providers", response_model=RecordArchiveProviderResponse) + def api_archive_providers( + session: Session = Depends(get_session), + principal: ApiPrincipal = Depends(get_api_principal), + ) -> RecordArchiveProviderResponse: + _require(principal, WRITE_SCOPE) + return RecordArchiveProviderResponse( + providers=records.archive_providers(session, principal) + ) + @router.get("", response_model=RecordListResponse) def api_list_records( query: str | None = Query(default=None, max_length=500), @@ -127,6 +169,12 @@ def create_router(registry: object | None = None) -> APIRouter: lambda: records.create_record( session, principal, payload=payload.model_dump(mode="python") ), + principal=principal, + operation_type="record.create", + idempotency_key=payload.idempotency_key, + request=payload.model_dump(mode="json"), + resource_type="record", + resource_id=payload.record_id or payload.record_number, ) @router.get("/{record_id}", response_model=RecordDetailResponse) @@ -146,6 +194,20 @@ def create_router(registry: object | None = None) -> APIRouter: except RecordStoreError as exc: raise _http_error(exc) from exc + @router.get("/{record_id}/recovery", response_model=RecordRecoveryStatusResponse) + def api_record_recovery_status( + record_id: str, + session: Session = Depends(get_session), + principal: ApiPrincipal = Depends(get_api_principal), + ) -> RecordRecoveryStatusResponse: + _require(principal, ADMIN_SCOPE) + try: + return RecordRecoveryStatusResponse( + **records.recovery_status(session, principal, record_id=record_id) + ) + except RecordStoreError as exc: + raise _http_error(exc) from exc + @router.patch("/{record_id}", response_model=dict[str, Any]) def api_update_record( record_id: str, @@ -162,6 +224,294 @@ def create_router(registry: object | None = None) -> APIRouter: record_id=record_id, payload=payload.model_dump(mode="python", exclude_unset=True), ), + principal=principal, + operation_type="record.revise", + idempotency_key=payload.idempotency_key, + request=payload.model_dump(mode="json", exclude_unset=True), + resource_type="record", + resource_id=record_id, + ) + + @router.post("/{record_id}/close", response_model=dict[str, Any]) + def api_close_record( + record_id: str, + payload: RecordCloseRequest, + session: Session = Depends(get_session), + principal: ApiPrincipal = Depends(get_api_principal), + ) -> dict[str, Any]: + _require(principal, WRITE_SCOPE) + if payload.restart_retention: + _require(principal, ADMIN_SCOPE) + return _write( + session, + lambda: records.close_record( + session, + principal, + record_id=record_id, + payload=payload.model_dump(mode="python"), + ), + principal=principal, + operation_type="record.close", + idempotency_key=payload.idempotency_key, + request=payload.model_dump(mode="json"), + resource_type="record", + resource_id=record_id, + ) + + @router.post("/{record_id}/reopen", response_model=dict[str, Any]) + def api_reopen_record( + record_id: str, + payload: RecordLifecycleActionRequest, + session: Session = Depends(get_session), + principal: ApiPrincipal = Depends(get_api_principal), + ) -> dict[str, Any]: + _require(principal, WRITE_SCOPE) + return _write( + session, + lambda: records.reopen_record( + session, + principal, + record_id=record_id, + payload=payload.model_dump(mode="python"), + ), + principal=principal, + operation_type="record.reopen", + idempotency_key=payload.idempotency_key, + request=payload.model_dump(mode="json"), + resource_type="record", + resource_id=record_id, + ) + + @router.post("/{record_id}/appraise", response_model=dict[str, Any]) + def api_appraise_record( + record_id: str, + payload: RecordAppraisalRequest, + session: Session = Depends(get_session), + principal: ApiPrincipal = Depends(get_api_principal), + ) -> dict[str, Any]: + _require(principal, WRITE_SCOPE) + if payload.override_retention_not_due: + _require(principal, ADMIN_SCOPE) + return _write( + session, + lambda: records.appraise_record( + session, + principal, + record_id=record_id, + payload=payload.model_dump(mode="python"), + ), + principal=principal, + operation_type="record.appraise", + idempotency_key=payload.idempotency_key, + request=payload.model_dump(mode="json"), + resource_type="record", + resource_id=record_id, + ) + + @router.post( + "/{record_id}/holds", + response_model=dict[str, Any], + status_code=status.HTTP_201_CREATED, + ) + def api_apply_hold( + record_id: str, + payload: RecordHoldCreateRequest, + session: Session = Depends(get_session), + principal: ApiPrincipal = Depends(get_api_principal), + ) -> dict[str, Any]: + _require(principal, WRITE_SCOPE) + return _write( + session, + lambda: records.apply_hold( + session, + principal, + record_id=record_id, + payload=payload.model_dump(mode="python"), + ), + principal=principal, + operation_type="hold.apply", + idempotency_key=payload.idempotency_key, + request=payload.model_dump(mode="json"), + resource_type="record", + resource_id=record_id, + ) + + @router.post( + "/{record_id}/holds/{hold_id}/release", + response_model=dict[str, Any], + ) + def api_release_hold( + record_id: str, + hold_id: str, + payload: RecordHoldReleaseRequest, + session: Session = Depends(get_session), + principal: ApiPrincipal = Depends(get_api_principal), + ) -> dict[str, Any]: + _require(principal, WRITE_SCOPE) + return _write( + session, + lambda: records.release_hold( + session, + principal, + record_id=record_id, + hold_id=hold_id, + payload=payload.model_dump(mode="python"), + ), + principal=principal, + operation_type="hold.release", + idempotency_key=payload.idempotency_key, + request={"hold_id": hold_id, **payload.model_dump(mode="json")}, + resource_type="record", + resource_id=record_id, + ) + + @router.post( + "/{record_id}/dispositions", + response_model=dict[str, Any], + status_code=status.HTTP_201_CREATED, + ) + def api_propose_disposition( + record_id: str, + payload: RecordDispositionCreateRequest, + session: Session = Depends(get_session), + principal: ApiPrincipal = Depends(get_api_principal), + ) -> dict[str, Any]: + _require(principal, WRITE_SCOPE) + return _write( + session, + lambda: records.propose_disposition( + session, + principal, + record_id=record_id, + payload=payload.model_dump(mode="python"), + ), + principal=principal, + operation_type="disposition.propose", + idempotency_key=payload.idempotency_key, + request=payload.model_dump(mode="json"), + resource_type="record", + resource_id=record_id, + ) + + @router.post( + "/{record_id}/dispositions/{disposition_id}/finalize", + response_model=dict[str, Any], + ) + def api_finalize_disposition( + record_id: str, + disposition_id: str, + payload: RecordDispositionFinalizeRequest, + session: Session = Depends(get_session), + principal: ApiPrincipal = Depends(get_api_principal), + ) -> dict[str, Any]: + _require(principal, WRITE_SCOPE) + return _write( + session, + lambda: records.finalize_disposition( + session, + principal, + record_id=record_id, + disposition_id=disposition_id, + payload=payload.model_dump(mode="python"), + ), + principal=principal, + operation_type="disposition.finalize", + idempotency_key=payload.idempotency_key, + request={ + "disposition_id": disposition_id, + **payload.model_dump(mode="json"), + }, + resource_type="record", + resource_id=record_id, + ) + + @router.post( + "/{record_id}/dispositions/{disposition_id}/withdraw", + response_model=dict[str, Any], + ) + def api_withdraw_disposition( + record_id: str, + disposition_id: str, + payload: RecordDispositionWithdrawRequest, + session: Session = Depends(get_session), + principal: ApiPrincipal = Depends(get_api_principal), + ) -> dict[str, Any]: + _require(principal, WRITE_SCOPE) + return _write( + session, + lambda: records.withdraw_disposition( + session, + principal, + record_id=record_id, + disposition_id=disposition_id, + payload=payload.model_dump(mode="python"), + ), + principal=principal, + operation_type="disposition.withdraw", + idempotency_key=payload.idempotency_key, + request={ + "disposition_id": disposition_id, + **payload.model_dump(mode="json"), + }, + resource_type="record", + resource_id=record_id, + ) + + @router.post( + "/{record_id}/transfer-packages", + response_model=dict[str, Any], + status_code=status.HTTP_201_CREATED, + ) + def api_prepare_transfer_package( + record_id: str, + payload: RecordTransferPackageCreateRequest, + session: Session = Depends(get_session), + principal: ApiPrincipal = Depends(get_api_principal), + ) -> dict[str, Any]: + _require(principal, WRITE_SCOPE) + return _write( + session, + lambda: records.prepare_transfer_package( + session, + principal, + record_id=record_id, + payload=payload.model_dump(mode="python"), + ), + principal=principal, + operation_type="transfer.prepare", + idempotency_key=payload.idempotency_key, + request=payload.model_dump(mode="json"), + resource_type="record", + resource_id=record_id, + ) + + @router.post( + "/{record_id}/transfer-packages/{package_id}/dispatch", + response_model=dict[str, Any], + ) + def api_dispatch_transfer_package( + record_id: str, + package_id: str, + payload: RecordTransferDispatchRequest, + session: Session = Depends(get_session), + principal: ApiPrincipal = Depends(get_api_principal), + ) -> dict[str, Any]: + _require(principal, WRITE_SCOPE) + return _write( + session, + lambda: records.dispatch_transfer_package( + session, + principal, + record_id=record_id, + package_id=package_id, + payload=payload.model_dump(mode="python"), + ), + principal=principal, + operation_type="transfer.simulate", + idempotency_key=payload.idempotency_key, + request={"package_id": package_id, **payload.model_dump(mode="json")}, + resource_type="record", + resource_id=record_id, ) @router.post( @@ -184,6 +534,12 @@ def create_router(registry: object | None = None) -> APIRouter: record_id=record_id, payload=payload.model_dump(mode="python"), ), + principal=principal, + operation_type="volume.create", + idempotency_key=payload.idempotency_key, + request=payload.model_dump(mode="json"), + resource_type="record", + resource_id=record_id, ) @router.post( @@ -235,7 +591,16 @@ def create_router(registry: object | None = None) -> APIRouter: }, } - return _write(session, operation) + return _write( + session, + operation, + principal=principal, + operation_type="item.file", + idempotency_key=payload.idempotency_key, + request=payload.model_dump(mode="json"), + resource_type="record", + resource_id=record_id, + ) return router @@ -245,18 +610,69 @@ def _require(principal: ApiPrincipal, scope: str) -> None: raise HTTPException(status_code=403, detail=f"Missing scope: {scope}") -def _write(session: Session, operation): +def _write( + session: Session, + operation, + *, + principal: ApiPrincipal, + operation_type: str, + idempotency_key: str, + request: dict[str, Any], + resource_type: str, + resource_id: str, +): + recovery = None try: - result = operation() - session.commit() + recovery = begin_record_atomic_recovery( + session, + tenant_id=principal.tenant_id, + operation_type=operation_type, + idempotency_key=idempotency_key, + request=request, + resource_type=resource_type, + resource_id=resource_id, + ) + recovery_token = bind_record_recovery_operation(recovery.operation_id) + try: + result = operation() + finally: + reset_record_recovery_operation(recovery_token) + if not recovery.replayed: + audit_from_principal( + session, + principal, + action=f"records.{operation_type}", + object_type=resource_type, + object_id=resource_id, + details={ + "recovery_operation_id": recovery.operation_id, + "idempotency_key_sha256": _sha256(idempotency_key), + }, + commit=False, + ) + recovery.commit_success(session, result=result, resource_id=resource_id) return result except (RecordStoreError, IntegrityError) as exc: session.rollback() + if recovery is not None: + recovery.reject(summary=str(exc), error_type=type(exc).__name__) if isinstance(exc, IntegrityError): raise HTTPException( status_code=409, detail="The record write conflicts with existing data." ) from exc raise _http_error(exc) from exc + except RecordRecoveryError as exc: + session.rollback() + raise HTTPException(status_code=503, detail=str(exc)) from exc + except Exception as exc: + session.rollback() + if recovery is not None: + recovery.fail(summary=str(exc), error_type=type(exc).__name__) + raise + + +def _sha256(value: str) -> str: + return hashlib.sha256(value.encode("utf-8")).hexdigest() def _http_error(exc: RecordStoreError) -> HTTPException: diff --git a/src/govoplan_records/backend/schemas.py b/src/govoplan_records/backend/schemas.py index 95e36a7..b105eec 100644 --- a/src/govoplan_records/backend/schemas.py +++ b/src/govoplan_records/backend/schemas.py @@ -148,6 +148,99 @@ class RecordItemCreateRequest(StrictModel): metadata: dict[str, Any] = Field(default_factory=dict) +class RecordLifecycleActionRequest(StrictModel): + expected_revision: int = Field(ge=1) + purpose: str = Field(min_length=1, max_length=255) + reason: str = Field(min_length=1, max_length=2_000) + recorded_at: datetime + idempotency_key: str = Field(min_length=1, max_length=255) + + +class RecordCloseRequest(RecordLifecycleActionRequest): + retention_trigger_at: datetime | None = None + restart_retention: bool = False + + +class RecordAppraisalRequest(RecordLifecycleActionRequest): + outcome: Literal["retain", "transfer", "destroy", "reclassify"] + policy_refs: list[str] = Field(default_factory=list, max_length=100) + override_retention_not_due: bool = False + + +class RecordHoldCreateRequest(StrictModel): + hold_id: str | None = Field(default=None, max_length=255) + expected_record_revision: int = Field(ge=1) + reason: str = Field(min_length=1, max_length=10_000) + authority: str = Field(min_length=1, max_length=500) + purpose: str = Field(min_length=1, max_length=255) + scope: dict[str, Any] = Field(default_factory=dict) + effective_from: datetime | None = None + effective_to: datetime | None = None + policy_refs: list[str] = Field(default_factory=list, max_length=100) + institutional_context: dict[str, Any] = Field(default_factory=dict) + recorded_at: datetime + idempotency_key: str = Field(min_length=1, max_length=255) + + @model_validator(mode="after") + def validate_interval(self): + if ( + self.effective_from + and self.effective_to + and self.effective_to <= self.effective_from + ): + raise ValueError("effective_to must be after effective_from") + return self + + +class RecordHoldReleaseRequest(StrictModel): + expected_hold_revision: int = Field(ge=1) + reason: str = Field(min_length=1, max_length=2_000) + purpose: str = Field(min_length=1, max_length=255) + recorded_at: datetime + idempotency_key: str = Field(min_length=1, max_length=255) + + +class RecordDispositionCreateRequest(StrictModel): + disposition_id: str | None = Field(default=None, max_length=255) + expected_record_revision: int = Field(ge=1) + action: Literal["retain", "transfer", "destroy", "reclassify"] + reason: str = Field(min_length=1, max_length=10_000) + purpose: str = Field(min_length=1, max_length=255) + policy_refs: list[str] = Field(default_factory=list, max_length=100) + institutional_context: dict[str, Any] = Field(default_factory=dict) + recorded_at: datetime + idempotency_key: str = Field(min_length=1, max_length=255) + + +class RecordDispositionFinalizeRequest(StrictModel): + expected_disposition_revision: int = Field(ge=1) + purpose: str = Field(min_length=1, max_length=255) + recorded_at: datetime + idempotency_key: str = Field(min_length=1, max_length=255) + + +class RecordDispositionWithdrawRequest(RecordDispositionFinalizeRequest): + reason: str = Field(min_length=1, max_length=2_000) + + +class RecordTransferPackageCreateRequest(StrictModel): + package_id: str | None = Field(default=None, max_length=255) + disposition_id: str = Field(min_length=1, max_length=255) + expected_record_revision: int = Field(ge=1) + provider_id: str = Field(min_length=1, max_length=100) + profile: str = Field(min_length=1, max_length=255) + purpose: str = Field(min_length=1, max_length=255) + recorded_at: datetime + idempotency_key: str = Field(min_length=1, max_length=255) + + +class RecordTransferDispatchRequest(StrictModel): + expected_package_revision: int = Field(ge=1) + purpose: str = Field(min_length=1, max_length=255) + recorded_at: datetime + idempotency_key: str = Field(min_length=1, max_length=255) + + class RecordListResponse(StrictModel): records: list[dict[str, Any]] total: int @@ -165,6 +258,9 @@ class RecordDetailResponse(StrictModel): volumes: list[dict[str, Any]] items: list[dict[str, Any]] chronology: list[dict[str, Any]] + holds: list[dict[str, Any]] + dispositions: list[dict[str, Any]] + transfer_packages: list[dict[str, Any]] access_explanation: dict[str, Any] @@ -172,15 +268,43 @@ class RecordSourceProviderResponse(StrictModel): providers: list[dict[str, Any]] +class RecordArchiveProviderResponse(StrictModel): + providers: list[dict[str, Any]] + + +class RecordRecoveryStatusResponse(StrictModel): + record_id: str + record_revision: int + evidence_sha256: str + healthy: bool + source_checks: list[dict[str, Any]] + package_checks: list[dict[str, Any]] + recovery_operations: list[dict[str, Any]] + failure_count: int + limitations: list[str] + + __all__ = [ "FilePlanNodeWriteRequest", + "RecordAppraisalRequest", + "RecordArchiveProviderResponse", "RecordCatalogResponse", "RecordClassWriteRequest", + "RecordCloseRequest", "RecordCreateRequest", "RecordDetailResponse", + "RecordDispositionCreateRequest", + "RecordDispositionFinalizeRequest", + "RecordDispositionWithdrawRequest", + "RecordHoldCreateRequest", + "RecordHoldReleaseRequest", "RecordItemCreateRequest", + "RecordLifecycleActionRequest", "RecordListResponse", + "RecordRecoveryStatusResponse", "RecordSourceProviderResponse", + "RecordTransferDispatchRequest", + "RecordTransferPackageCreateRequest", "RecordUpdateRequest", "RecordVolumeCreateRequest", ] diff --git a/src/govoplan_records/backend/service.py b/src/govoplan_records/backend/service.py index c7bb204..2824b82 100644 --- a/src/govoplan_records/backend/service.py +++ b/src/govoplan_records/backend/service.py @@ -1,7 +1,7 @@ from __future__ import annotations from collections.abc import Mapping, Sequence -from datetime import UTC, datetime +from datetime import UTC, datetime, timedelta import hashlib import json from typing import Any @@ -11,26 +11,46 @@ from sqlalchemy import func, or_ from sqlalchemy.orm import Session from govoplan_core.core.records import ( + RecordArchiveProvider, + RecordArchiveTransferRequest, RecordContractError, RecordFilingRequest, RecordFilingResult, RecordSourceLocator, RecordSourceProvider, RecordSourceReference, + RecordTransferPackage, + record_archive_capabilities, + record_archive_capability, record_source_capabilities, record_source_capability, ) +from govoplan_core.core.approvals import ( + CAPABILITY_APPROVAL_REQUESTS, + ApprovalActorSelector, + ApprovalRequestCreateCommand, + ApprovalRequestProvider, + ApprovalStepDefinition, +) +from govoplan_core.core.recovery import ( + RecoveryOperation, + verify_recovery_evidence_chain, +) from govoplan_core.core.temporal import current_temporal_data_context from govoplan_core.db.temporal import apply_temporal_revision_filter from govoplan_records.backend.db.models import ( RecordChronologyEntry, RecordClassRevision, + RecordDispositionRevision, RecordFilePlanRevision, + RecordHoldRevision, RecordIdentity, RecordItem, RecordRevision, + RecordTransferPackageRevision, RecordVolumeRevision, ) +from govoplan_records.backend.recovery import current_record_recovery_operation class RecordStoreError(ValueError): @@ -370,6 +390,14 @@ class SqlRecordRegistry: identity = session.get(RecordIdentity, current.identity_id) if identity is None: raise RecordStoreError("Record identity is missing.") + if ( + current.state not in {"planned", "open"} + and "state" in payload + and str(payload["state"]) != current.state + ): + raise RecordStoreError( + "Lifecycle state can only change through a governed lifecycle action." + ) values = _revision_values(current, payload) if values["access_mode"] != "tenant": raise RecordStoreError( @@ -497,11 +525,23 @@ class SqlRecordRegistry: chronology = _record_chronology( session, tenant_id=tenant_id, record_id=record_id ) + holds = _record_holds(session, tenant_id=tenant_id, record_id=record_id) + dispositions = _record_dispositions( + session, tenant_id=tenant_id, record_id=record_id + ) + transfer_packages = _record_transfer_packages( + session, tenant_id=tenant_id, record_id=record_id + ) return { "record": _record_dict(row, identity), "volumes": [_volume_dict(item) for item in volumes], "items": [_item_dict(item) for item in items], "chronology": [_chronology_dict(item) for item in chronology], + "holds": [_hold_dict(item) for item in holds], + "dispositions": [_disposition_dict(item) for item in dispositions], + "transfer_packages": [ + _transfer_package_dict(item) for item in transfer_packages + ], "access_explanation": { "decision": "allowed", "reason": "Current tenant and Records permission were evaluated for this read.", @@ -516,6 +556,1233 @@ class SqlRecordRegistry: }, } + def close_record( + self, + session: Session, + principal: object, + *, + record_id: str, + payload: Mapping[str, object], + ) -> dict[str, Any]: + tenant_id = _tenant(principal) + request_hash = _request_hash({"record_id": record_id, **payload}) + replay = _replay_event( + session, + tenant_id=tenant_id, + idempotency_key=_text(payload, "idempotency_key"), + request_hash=request_hash, + ) + if replay is not None: + return self.get_record( + session, + principal, + record_id=record_id, + revision=replay.record_revision, + )["record"] + current = _required_current_record( + session, tenant_id=tenant_id, record_id=record_id, lock=True + ) + _validate_expected( + current.revision, + payload.get("expected_revision"), + label="Record", + ) + if current.state != "open": + raise RecordStoreError("Only open records can be closed.") + recorded_at = _timestamp(payload.get("recorded_at"), "recorded_at") + _validate_later_revision(current.recorded_at, recorded_at) + record_class = _current_class(session, tenant_id, current.class_id) + if record_class is None: + raise RecordStoreError("The current record class is unavailable.") + restart_retention = bool(payload.get("restart_retention", False)) + if current.retention_started_at is not None and not restart_retention: + retention_started_at = current.retention_started_at + retention_due_at = current.retention_due_at + retention_rule = dict(current.retention_rule) + else: + retention_started_at, retention_due_at, retention_rule = ( + _calculate_retention( + record_class, + closed_at=recorded_at, + explicit_trigger_at=_optional_timestamp( + payload.get("retention_trigger_at"), + "retention_trigger_at", + ), + ) + ) + state = "retention_running" if retention_due_at is not None else "closed" + row, identity = _revise_record_state( + session, + principal, + current=current, + recorded_at=recorded_at, + values={ + "state": state, + "closed_at": recorded_at, + "retention_started_at": retention_started_at, + "retention_due_at": retention_due_at, + "retention_rule": retention_rule, + "appraisal_state": None, + "appraisal": {}, + }, + snapshot_patch={ + "lifecycle": { + "event": "closed", + "reason": _text(payload, "reason"), + "retention_rule": retention_rule, + } + }, + ) + _append_event( + session, + principal, + row=row, + event_type="record.closed", + summary=f"Record {identity.record_number} closed", + occurred_at=recorded_at, + purpose=_text(payload, "purpose"), + idempotency_key=_text(payload, "idempotency_key"), + request_hash=request_hash, + institutional_context=row.institutional_context, + payload={ + "reason": _text(payload, "reason"), + "retention_started_at": _datetime_text(retention_started_at), + "retention_due_at": _datetime_text(retention_due_at), + "restart_retention": restart_retention, + }, + ) + session.flush() + return _record_dict(row, identity) + + def reopen_record( + self, + session: Session, + principal: object, + *, + record_id: str, + payload: Mapping[str, object], + ) -> dict[str, Any]: + tenant_id = _tenant(principal) + request_hash = _request_hash({"record_id": record_id, **payload}) + replay = _replay_event( + session, + tenant_id=tenant_id, + idempotency_key=_text(payload, "idempotency_key"), + request_hash=request_hash, + ) + if replay is not None: + return self.get_record( + session, + principal, + record_id=record_id, + revision=replay.record_revision, + )["record"] + current = _required_current_record( + session, tenant_id=tenant_id, record_id=record_id, lock=True + ) + _validate_expected( + current.revision, + payload.get("expected_revision"), + label="Record", + ) + if current.state not in {"closed", "retention_running", "appraised"}: + raise RecordStoreError( + "Only closed, retention-running, or appraised records can be reopened." + ) + if _current_disposition_for_record(session, tenant_id, record_id) is not None: + raise RecordStoreError( + "A record with a disposition proposal cannot be reopened." + ) + recorded_at = _timestamp(payload.get("recorded_at"), "recorded_at") + _validate_later_revision(current.recorded_at, recorded_at) + row, identity = _revise_record_state( + session, + principal, + current=current, + recorded_at=recorded_at, + values={"state": "open", "appraisal_state": None, "appraisal": {}}, + snapshot_patch={ + "lifecycle": { + "event": "reopened", + "reason": _text(payload, "reason"), + } + }, + ) + _append_event( + session, + principal, + row=row, + event_type="record.reopened", + summary=f"Record {identity.record_number} reopened", + occurred_at=recorded_at, + purpose=_text(payload, "purpose"), + idempotency_key=_text(payload, "idempotency_key"), + request_hash=request_hash, + institutional_context=row.institutional_context, + payload={ + "reason": _text(payload, "reason"), + "retention_schedule_preserved": True, + }, + ) + session.flush() + return _record_dict(row, identity) + + def appraise_record( + self, + session: Session, + principal: object, + *, + record_id: str, + payload: Mapping[str, object], + ) -> dict[str, Any]: + tenant_id = _tenant(principal) + request_hash = _request_hash({"record_id": record_id, **payload}) + replay = _replay_event( + session, + tenant_id=tenant_id, + idempotency_key=_text(payload, "idempotency_key"), + request_hash=request_hash, + ) + if replay is not None: + return self.get_record( + session, + principal, + record_id=record_id, + revision=replay.record_revision, + )["record"] + current = _required_current_record( + session, tenant_id=tenant_id, record_id=record_id, lock=True + ) + _validate_expected( + current.revision, + payload.get("expected_revision"), + label="Record", + ) + if current.state not in {"closed", "retention_running"}: + raise RecordStoreError("Only closed records can be appraised.") + recorded_at = _timestamp(payload.get("recorded_at"), "recorded_at") + _validate_later_revision(current.recorded_at, recorded_at) + override = bool(payload.get("override_retention_not_due", False)) + if ( + current.retention_due_at is not None + and _aware(recorded_at) < _aware(current.retention_due_at) + and not override + ): + raise RecordStoreError( + "The retention period has not elapsed; an administrator override is required." + ) + outcome = _text(payload, "outcome") + if outcome not in {"retain", "transfer", "destroy", "reclassify"}: + raise RecordStoreError("Unsupported appraisal outcome.") + appraisal = { + "outcome": outcome, + "reason": _text(payload, "reason"), + "policy_refs": _text_list(payload.get("policy_refs")), + "retention_not_due_override": override, + "appraised_at": _datetime_text(recorded_at), + "appraised_by": _actor(principal), + } + row, identity = _revise_record_state( + session, + principal, + current=current, + recorded_at=recorded_at, + values={ + "state": "appraised", + "appraisal_state": outcome, + "appraisal": appraisal, + }, + snapshot_patch={"appraisal": appraisal}, + ) + _append_event( + session, + principal, + row=row, + event_type="record.appraised", + summary=f"Record {identity.record_number} appraised for {outcome}", + occurred_at=recorded_at, + purpose=_text(payload, "purpose"), + idempotency_key=_text(payload, "idempotency_key"), + request_hash=request_hash, + institutional_context=row.institutional_context, + payload=appraisal, + ) + session.flush() + return _record_dict(row, identity) + + def apply_hold( + self, + session: Session, + principal: object, + *, + record_id: str, + payload: Mapping[str, object], + ) -> dict[str, Any]: + tenant_id = _tenant(principal) + request_hash = _request_hash({"record_id": record_id, **payload}) + replay = _replay_by_key( + session, + RecordHoldRevision, + tenant_id=tenant_id, + idempotency_key=_text(payload, "idempotency_key"), + request_hash=request_hash, + ) + if replay is not None: + return _hold_dict(replay) + record = _required_current_record( + session, tenant_id=tenant_id, record_id=record_id, lock=True + ) + _validate_expected( + record.revision, + payload.get("expected_record_revision"), + label="Record", + ) + hold_id = _optional_text(payload.get("hold_id")) or str(uuid.uuid4()) + if _current_hold(session, tenant_id, hold_id) is not None: + raise RecordConflictError("A current hold already uses this identifier.") + recorded_at = _timestamp(payload.get("recorded_at"), "recorded_at") + effective_from = ( + _optional_timestamp(payload.get("effective_from"), "effective_from") + or recorded_at + ) + effective_to = _optional_timestamp(payload.get("effective_to"), "effective_to") + _validate_interval(effective_from, effective_to) + row = RecordHoldRevision( + tenant_id=tenant_id, + hold_id=hold_id, + record_id=record_id, + revision=1, + status="active", + reason=_text(payload, "reason"), + authority=_text(payload, "authority"), + scope=_mapping(payload.get("scope")), + effective_from=effective_from, + effective_to=effective_to, + policy_refs=_text_list(payload.get("policy_refs")), + institutional_context={ + **dict(record.institutional_context), + **_mapping(payload.get("institutional_context")), + }, + recorded_at=recorded_at, + changed_by=_actor(principal), + idempotency_key=_text(payload, "idempotency_key"), + request_sha256=request_hash, + ) + session.add(row) + _append_event( + session, + principal, + row=record, + event_type="record.hold_applied", + summary=f"Hold applied: {row.authority}", + occurred_at=recorded_at, + purpose=_text(payload, "purpose"), + idempotency_key=row.idempotency_key, + request_hash=request_hash, + institutional_context=row.institutional_context, + payload={ + "hold_id": hold_id, + "reason": row.reason, + "authority": row.authority, + "effective_from": _datetime_text(effective_from), + "effective_to": _datetime_text(effective_to), + }, + ) + session.flush() + return _hold_dict(row) + + def release_hold( + self, + session: Session, + principal: object, + *, + record_id: str, + hold_id: str, + payload: Mapping[str, object], + ) -> dict[str, Any]: + tenant_id = _tenant(principal) + request_hash = _request_hash( + {"record_id": record_id, "hold_id": hold_id, **payload} + ) + replay = _replay_by_key( + session, + RecordHoldRevision, + tenant_id=tenant_id, + idempotency_key=_text(payload, "idempotency_key"), + request_hash=request_hash, + ) + if replay is not None: + return _hold_dict(replay) + record = _required_current_record( + session, tenant_id=tenant_id, record_id=record_id, lock=True + ) + current = _current_hold(session, tenant_id, hold_id, lock=True) + if current is None or current.record_id != record_id: + raise RecordNotFoundError("Record hold not found.") + _validate_expected( + current.revision, + payload.get("expected_hold_revision"), + label="Record hold", + ) + if current.status != "active": + raise RecordStoreError("Only an active hold can be released.") + recorded_at = _timestamp(payload.get("recorded_at"), "recorded_at") + _validate_later_revision(current.recorded_at, recorded_at) + current.superseded_at = recorded_at + row = RecordHoldRevision( + tenant_id=tenant_id, + hold_id=hold_id, + record_id=record_id, + revision=current.revision + 1, + previous_revision_id=current.id, + status="released", + reason=current.reason, + authority=current.authority, + scope={ + **dict(current.scope), + "release_reason": _text(payload, "reason"), + }, + effective_from=current.effective_from, + effective_to=current.effective_to or recorded_at, + released_at=recorded_at, + policy_refs=list(current.policy_refs), + institutional_context=dict(current.institutional_context), + recorded_at=recorded_at, + changed_by=_actor(principal), + idempotency_key=_text(payload, "idempotency_key"), + request_sha256=request_hash, + ) + session.add(row) + _append_event( + session, + principal, + row=record, + event_type="record.hold_released", + summary=f"Hold released: {row.authority}", + occurred_at=recorded_at, + purpose=_text(payload, "purpose"), + idempotency_key=row.idempotency_key, + request_hash=request_hash, + institutional_context=row.institutional_context, + payload={ + "hold_id": hold_id, + "reason": _text(payload, "reason"), + "released_revision": row.revision, + }, + ) + session.flush() + return _hold_dict(row) + + def propose_disposition( + self, + session: Session, + principal: object, + *, + record_id: str, + payload: Mapping[str, object], + ) -> dict[str, Any]: + tenant_id = _tenant(principal) + request_hash = _request_hash({"record_id": record_id, **payload}) + replay = _replay_by_key( + session, + RecordDispositionRevision, + tenant_id=tenant_id, + idempotency_key=_text(payload, "idempotency_key"), + request_hash=request_hash, + ) + if replay is not None: + return _disposition_dict(replay) + record = _required_current_record( + session, tenant_id=tenant_id, record_id=record_id, lock=True + ) + _validate_expected( + record.revision, + payload.get("expected_record_revision"), + label="Record", + ) + if record.state != "appraised": + raise RecordStoreError( + "A disposition proposal requires an appraised record." + ) + recorded_at = _timestamp(payload.get("recorded_at"), "recorded_at") + _require_no_active_hold(session, tenant_id, record_id, at=recorded_at) + if _current_disposition_for_record(session, tenant_id, record_id) is not None: + raise RecordConflictError( + "The record already has a current disposition proposal." + ) + action = _text(payload, "action") + if action not in {"retain", "transfer", "destroy", "reclassify"}: + raise RecordStoreError("Unsupported disposition action.") + if record.appraisal_state != action: + raise RecordStoreError( + "The disposition action must match the current appraisal outcome." + ) + disposition_id = _optional_text(payload.get("disposition_id")) or str( + uuid.uuid4() + ) + evidence = _record_evidence_manifest( + session, tenant_id=tenant_id, record=record + ) + subject_sha256 = _request_hash(evidence) + policy_refs = _text_list(payload.get("policy_refs")) + approval_request_id: str | None = None + status = "review_unavailable" + approval_provider = self._tenant_capability( + CAPABILITY_APPROVAL_REQUESTS, session, tenant_id + ) + if isinstance(approval_provider, ApprovalRequestProvider): + approval = approval_provider.create_request( + session, + principal, + command=ApprovalRequestCreateCommand( + title=f"Disposition of record {evidence['record_number']}", + description=_text(payload, "reason"), + subject_module="records", + subject_type="record_disposition", + subject_id=disposition_id, + subject_version=str(record.revision), + subject_digest=subject_sha256, + steps=( + ApprovalStepDefinition( + key="records-disposition-review", + label="Independent disposition review", + selectors=( + ApprovalActorSelector( + kind="any_account", + value="*", + label="Authorized Records reviewer", + ), + ), + required_approvals=1, + forbidden_evidence_roles=("proposer",), + ), + ), + separation_of_duties=True, + policy_refs=tuple(policy_refs), + evidence_actors={"proposer": ((_actor(principal) or "unknown"),)}, + metadata={ + "record_id": record_id, + "record_revision": record.revision, + "action": action, + }, + ), + idempotency_key=f"records-disposition:{_text(payload, 'idempotency_key')}", + ) + approval_request_id = approval.id + status = "review_pending" + consequence_preview = _disposition_consequence_preview( + action=action, + record=record, + evidence=evidence, + approvals_available=approval_request_id is not None, + ) + row = RecordDispositionRevision( + tenant_id=tenant_id, + disposition_id=disposition_id, + record_id=record_id, + revision=1, + action=action, + status=status, + reason=_text(payload, "reason"), + subject_revision=record.revision, + subject_sha256=subject_sha256, + consequence_preview=consequence_preview, + policy_refs=policy_refs, + approval_request_id=approval_request_id, + proposed_by=_actor(principal), + institutional_context={ + **dict(record.institutional_context), + **_mapping(payload.get("institutional_context")), + }, + recorded_at=recorded_at, + idempotency_key=_text(payload, "idempotency_key"), + request_sha256=request_hash, + ) + session.add(row) + _append_event( + session, + principal, + row=record, + event_type="record.disposition_proposed", + summary=f"Disposition proposed: {action}", + occurred_at=recorded_at, + purpose=_text(payload, "purpose"), + idempotency_key=row.idempotency_key, + request_hash=request_hash, + institutional_context=row.institutional_context, + payload={ + "disposition_id": disposition_id, + "action": action, + "subject_sha256": subject_sha256, + "approval_request_id": approval_request_id, + "consequence_preview": consequence_preview, + }, + ) + session.flush() + return _disposition_dict(row) + + def finalize_disposition( + self, + session: Session, + principal: object, + *, + record_id: str, + disposition_id: str, + payload: Mapping[str, object], + ) -> dict[str, Any]: + tenant_id = _tenant(principal) + request_hash = _request_hash( + {"record_id": record_id, "disposition_id": disposition_id, **payload} + ) + replay = _replay_by_key( + session, + RecordDispositionRevision, + tenant_id=tenant_id, + idempotency_key=_text(payload, "idempotency_key"), + request_hash=request_hash, + ) + if replay is not None: + record = _required_current_record( + session, tenant_id=tenant_id, record_id=record_id + ) + identity = session.get(RecordIdentity, record.identity_id) + if identity is None: + raise RecordStoreError("Record identity is missing.") + return { + "disposition": _disposition_dict(replay), + "record": _record_dict(record, identity), + } + record = _required_current_record( + session, tenant_id=tenant_id, record_id=record_id, lock=True + ) + current = _current_disposition(session, tenant_id, disposition_id, lock=True) + if current is None or current.record_id != record_id: + raise RecordNotFoundError("Record disposition not found.") + _validate_expected( + current.revision, + payload.get("expected_disposition_revision"), + label="Record disposition", + ) + if current.status != "review_pending": + raise RecordStoreError( + "The disposition is not awaiting an available approval." + ) + recorded_at = _timestamp(payload.get("recorded_at"), "recorded_at") + _require_no_active_hold(session, tenant_id, record_id, at=recorded_at) + evidence = _record_evidence_manifest( + session, tenant_id=tenant_id, record=record + ) + if _request_hash(evidence) != current.subject_sha256: + raise RecordConflictError( + "The record evidence changed after the disposition was proposed." + ) + approval_provider = self._tenant_capability( + CAPABILITY_APPROVAL_REQUESTS, session, tenant_id + ) + if not isinstance(approval_provider, ApprovalRequestProvider): + raise RecordSourceUnavailableError( + "The Approvals module is required to finalize a disposition." + ) + if not current.approval_request_id: + raise RecordStoreError("The disposition has no approval request.") + approval = approval_provider.check_approved( + session, + principal, + request_id=current.approval_request_id, + subject_module="records", + subject_type="record_disposition", + subject_id=disposition_id, + subject_version=str(current.subject_revision), + subject_digest=current.subject_sha256, + ) + if not approval.approved: + raise RecordStoreError("The disposition has not been approved.") + current.superseded_at = recorded_at + disposition = RecordDispositionRevision( + tenant_id=tenant_id, + disposition_id=disposition_id, + record_id=record_id, + revision=current.revision + 1, + previous_revision_id=current.id, + action=current.action, + status="approved", + reason=current.reason, + subject_revision=current.subject_revision, + subject_sha256=current.subject_sha256, + consequence_preview=dict(current.consequence_preview), + policy_refs=list(current.policy_refs), + approval_request_id=current.approval_request_id, + proposed_by=current.proposed_by, + reviewed_by=_actor(principal), + reviewed_at=recorded_at, + institutional_context=dict(current.institutional_context), + recorded_at=recorded_at, + idempotency_key=_text(payload, "idempotency_key"), + request_sha256=request_hash, + ) + session.add(disposition) + state = { + "retain": "closed", + "transfer": "transfer_pending", + "destroy": "destruction_pending", + "reclassify": "open", + }[current.action] + revised, identity = _revise_record_state( + session, + principal, + current=record, + recorded_at=recorded_at, + values={"state": state}, + snapshot_patch={ + "disposition": { + "id": disposition_id, + "action": current.action, + "status": "approved", + "subject_sha256": current.subject_sha256, + } + }, + ) + _append_event( + session, + principal, + row=revised, + event_type="record.disposition_approved", + summary=f"Disposition approved: {current.action}", + occurred_at=recorded_at, + purpose=_text(payload, "purpose"), + idempotency_key=disposition.idempotency_key, + request_hash=request_hash, + institutional_context=disposition.institutional_context, + payload={ + "disposition_id": disposition_id, + "approval_request_id": current.approval_request_id, + "subject_sha256": current.subject_sha256, + "record_state": state, + }, + ) + session.flush() + return { + "disposition": _disposition_dict(disposition), + "record": _record_dict(revised, identity), + } + + def withdraw_disposition( + self, + session: Session, + principal: object, + *, + record_id: str, + disposition_id: str, + payload: Mapping[str, object], + ) -> dict[str, Any]: + tenant_id = _tenant(principal) + request_hash = _request_hash( + {"record_id": record_id, "disposition_id": disposition_id, **payload} + ) + replay = _replay_by_key( + session, + RecordDispositionRevision, + tenant_id=tenant_id, + idempotency_key=_text(payload, "idempotency_key"), + request_hash=request_hash, + ) + if replay is not None: + return _disposition_dict(replay) + record = _required_current_record( + session, tenant_id=tenant_id, record_id=record_id, lock=True + ) + current = _current_disposition(session, tenant_id, disposition_id, lock=True) + if current is None or current.record_id != record_id: + raise RecordNotFoundError("Record disposition not found.") + _validate_expected( + current.revision, + payload.get("expected_disposition_revision"), + label="Record disposition", + ) + if current.status not in {"review_pending", "review_unavailable"}: + raise RecordStoreError( + "Only an unapproved disposition proposal can be withdrawn." + ) + recorded_at = _timestamp(payload.get("recorded_at"), "recorded_at") + _validate_later_revision(current.recorded_at, recorded_at) + current.superseded_at = recorded_at + row = RecordDispositionRevision( + tenant_id=tenant_id, + disposition_id=disposition_id, + record_id=record_id, + revision=current.revision + 1, + previous_revision_id=current.id, + action=current.action, + status="withdrawn", + reason=current.reason, + subject_revision=current.subject_revision, + subject_sha256=current.subject_sha256, + consequence_preview={ + **dict(current.consequence_preview), + "withdrawal_reason": _text(payload, "reason"), + }, + policy_refs=list(current.policy_refs), + approval_request_id=current.approval_request_id, + proposed_by=current.proposed_by, + reviewed_by=_actor(principal), + reviewed_at=recorded_at, + institutional_context=dict(current.institutional_context), + recorded_at=recorded_at, + idempotency_key=_text(payload, "idempotency_key"), + request_sha256=request_hash, + ) + session.add(row) + _append_event( + session, + principal, + row=record, + event_type="record.disposition_withdrawn", + summary=f"Disposition withdrawn: {current.action}", + occurred_at=recorded_at, + purpose=_text(payload, "purpose"), + idempotency_key=row.idempotency_key, + request_hash=request_hash, + institutional_context=row.institutional_context, + payload={ + "disposition_id": disposition_id, + "reason": _text(payload, "reason"), + "superseded_revision": current.revision, + }, + ) + session.flush() + return _disposition_dict(row) + + def archive_providers( + self, session: Session, principal: object + ) -> list[dict[str, Any]]: + tenant_id = _tenant(principal) + providers: list[dict[str, Any]] = [] + for capability_name in record_archive_capabilities(self.registry): + provider = self._tenant_capability(capability_name, session, tenant_id) + if not isinstance(provider, RecordArchiveProvider): + continue + state = provider.state() + providers.append( + { + "id": state.provider_id, + "label": state.label, + "profiles": list(state.profiles), + "authority_modes": list(state.authority_modes), + "healthy": state.healthy, + "checked_at": _datetime_text(state.checked_at), + "last_success_at": _datetime_text(state.last_success_at), + "freshness_seconds": state.freshness_seconds, + "limitations": list(state.limitations), + "simulated": state.simulated, + } + ) + return providers + + def recovery_status( + self, + session: Session, + principal: object, + *, + record_id: str, + ) -> dict[str, Any]: + tenant_id = _tenant(principal) + record = _required_current_record( + session, tenant_id=tenant_id, record_id=record_id + ) + identity = session.get(RecordIdentity, record.identity_id) + if identity is None: + raise RecordStoreError("Record identity is missing.") + source_checks: list[dict[str, Any]] = [] + for item in ( + session.query(RecordItem) + .filter( + RecordItem.tenant_id == tenant_id, + RecordItem.record_id == record_id, + ) + .order_by(RecordItem.sequence) + .all() + ): + check = { + "item_id": item.id, + "source_module": item.source_module, + "resource_type": item.resource_type, + "resource_id": item.resource_id, + "source_revision": item.source_revision, + "stored_content_sha256": item.content_sha256, + "status": "unavailable", + "error": None, + } + provider = self._tenant_capability( + record_source_capability(item.source_module), session, tenant_id + ) + if not isinstance(provider, RecordSourceProvider): + check["error"] = "The source provider is not enabled." + source_checks.append(check) + continue + try: + reference = provider.resolve( + session, + principal, + locator=RecordSourceLocator( + tenant_id=tenant_id, + source_module=item.source_module, + resource_type=item.resource_type, + resource_id=item.resource_id, + source_revision=item.source_revision, + metadata=dict(item.source_metadata), + ), + purpose=item.purpose, + ) + except Exception as exc: + check["error"] = f"Source revalidation failed ({type(exc).__name__})." + source_checks.append(check) + continue + resolved_digest = (reference.content_sha256 or "").removeprefix( + "sha256:" + ) or None + check["resolved_content_sha256"] = resolved_digest + if item.content_sha256 and resolved_digest != item.content_sha256: + check["status"] = "digest_mismatch" + check["error"] = ( + "The source digest no longer matches the filed evidence." + ) + else: + check["status"] = "verified" + source_checks.append(check) + + package_checks = [] + for package in _record_transfer_packages( + session, tenant_id=tenant_id, record_id=record_id + ): + calculated = _request_hash(dict(package.manifest)) + package_checks.append( + { + "package_id": package.package_id, + "revision": package.revision, + "status": package.status, + "simulated": package.simulated, + "stored_manifest_sha256": package.manifest_sha256, + "calculated_manifest_sha256": calculated, + "manifest_verified": calculated == package.manifest_sha256, + "receipt_sha256": package.receipt_sha256, + "recovery_operation_id": package.recovery_operation_id, + } + ) + + operations = ( + session.query(RecoveryOperation) + .filter( + RecoveryOperation.module_id == "records", + RecoveryOperation.resource_type == "record", + RecoveryOperation.resource_id.in_((record_id, identity.record_number)), + ) + .order_by(RecoveryOperation.created_at) + .all() + ) + operation_checks = [ + { + "operation_id": operation.id, + "operation_type": operation.operation_type, + "status": operation.status, + "mode": operation.mode, + "checkpoint_count": operation.checkpoint_count, + "evidence_head_sha256": operation.evidence_head_sha256, + "evidence_chain_verified": verify_recovery_evidence_chain( + session, operation.id + ), + "completed_at": _datetime_text(operation.completed_at), + } + for operation in operations + ] + evidence = _record_evidence_manifest( + session, tenant_id=tenant_id, record=record + ) + failures = [ + check for check in source_checks if check["status"] not in {"verified"} + ] + failures.extend( + check for check in package_checks if not check["manifest_verified"] + ) + failures.extend( + check for check in operation_checks if not check["evidence_chain_verified"] + ) + return { + "record_id": record_id, + "record_revision": record.revision, + "evidence_sha256": _request_hash(evidence), + "healthy": not failures, + "source_checks": source_checks, + "package_checks": package_checks, + "recovery_operations": operation_checks, + "failure_count": len(failures), + "limitations": [ + "Source checks use current authorization and never bypass the owning module.", + "A verified simulation receipt does not prove archival custody.", + ], + } + + def prepare_transfer_package( + self, + session: Session, + principal: object, + *, + record_id: str, + payload: Mapping[str, object], + ) -> dict[str, Any]: + tenant_id = _tenant(principal) + request_hash = _request_hash({"record_id": record_id, **payload}) + replay = _replay_by_key( + session, + RecordTransferPackageRevision, + tenant_id=tenant_id, + idempotency_key=_text(payload, "idempotency_key"), + request_hash=request_hash, + ) + if replay is not None: + return _transfer_package_dict(replay) + record = _required_current_record( + session, tenant_id=tenant_id, record_id=record_id, lock=True + ) + _validate_expected( + record.revision, + payload.get("expected_record_revision"), + label="Record", + ) + recorded_at = _timestamp(payload.get("recorded_at"), "recorded_at") + _require_no_active_hold(session, tenant_id, record_id, at=recorded_at) + disposition_id = _text(payload, "disposition_id") + disposition = _current_disposition(session, tenant_id, disposition_id) + if ( + disposition is None + or disposition.record_id != record_id + or disposition.action != "transfer" + or disposition.status != "approved" + ): + raise RecordStoreError( + "An approved transfer disposition is required before packaging." + ) + provider_id = _text(payload, "provider_id") + provider = self._tenant_capability( + record_archive_capability(provider_id), session, tenant_id + ) + if not isinstance(provider, RecordArchiveProvider): + raise RecordSourceUnavailableError( + "The selected record archive provider is unavailable." + ) + provider_state = provider.state() + profile = _text(payload, "profile") + if not provider_state.healthy or profile not in provider_state.profiles: + raise RecordSourceUnavailableError( + "The selected archive profile is unavailable or unhealthy." + ) + package_id = _optional_text(payload.get("package_id")) or str(uuid.uuid4()) + manifest = { + "format": "govoplan-record-transfer-manifest", + "version": 1, + "package_id": package_id, + "provider_id": provider_id, + "profile": profile, + "prepared_at": _datetime_text(recorded_at), + "record": _record_evidence_manifest( + session, tenant_id=tenant_id, record=record + ), + "disposition": _disposition_dict(disposition), + "institutional_context": dict(record.institutional_context), + } + manifest_sha256 = _request_hash(manifest) + row = RecordTransferPackageRevision( + tenant_id=tenant_id, + package_id=package_id, + record_id=record_id, + disposition_id=disposition_id, + revision=1, + record_revision=record.revision, + provider_id=provider_id, + profile=profile, + status="prepared", + authority_mode=provider_state.authority_modes[0], + manifest=manifest, + manifest_sha256=manifest_sha256, + receipt={}, + recovery_operation_id=current_record_recovery_operation(), + simulated=provider_state.simulated, + institutional_context=dict(record.institutional_context), + recorded_at=recorded_at, + changed_by=_actor(principal), + idempotency_key=_text(payload, "idempotency_key"), + request_sha256=request_hash, + ) + session.add(row) + _append_event( + session, + principal, + row=record, + event_type="record.transfer_package_prepared", + summary=f"Transfer package prepared for {provider_state.label}", + occurred_at=recorded_at, + purpose=_text(payload, "purpose"), + idempotency_key=row.idempotency_key, + request_hash=request_hash, + institutional_context=row.institutional_context, + payload={ + "package_id": package_id, + "provider_id": provider_id, + "profile": profile, + "manifest_sha256": manifest_sha256, + "simulated": provider_state.simulated, + }, + ) + session.flush() + return _transfer_package_dict(row) + + def dispatch_transfer_package( + self, + session: Session, + principal: object, + *, + record_id: str, + package_id: str, + payload: Mapping[str, object], + ) -> dict[str, Any]: + tenant_id = _tenant(principal) + request_hash = _request_hash( + {"record_id": record_id, "package_id": package_id, **payload} + ) + replay = _replay_by_key( + session, + RecordTransferPackageRevision, + tenant_id=tenant_id, + idempotency_key=_text(payload, "idempotency_key"), + request_hash=request_hash, + ) + if replay is not None: + return _transfer_package_dict(replay) + record = _required_current_record( + session, tenant_id=tenant_id, record_id=record_id, lock=True + ) + current = _current_transfer_package(session, tenant_id, package_id, lock=True) + if current is None or current.record_id != record_id: + raise RecordNotFoundError("Record transfer package not found.") + _validate_expected( + current.revision, + payload.get("expected_package_revision"), + label="Record transfer package", + ) + if current.status != "prepared": + raise RecordStoreError( + "Only a prepared package can be dispatched; unknown outcomes require reconciliation." + ) + recorded_at = _timestamp(payload.get("recorded_at"), "recorded_at") + _require_no_active_hold(session, tenant_id, record_id, at=recorded_at) + provider = self._tenant_capability( + record_archive_capability(current.provider_id), session, tenant_id + ) + if not isinstance(provider, RecordArchiveProvider): + raise RecordSourceUnavailableError( + "The record archive provider is unavailable." + ) + if not provider.state().simulated: + raise RecordSourceUnavailableError( + "Real archive dispatch requires a configured target-specific recovery profile." + ) + try: + receipt = provider.dispatch( + session, + principal, + request=RecordArchiveTransferRequest( + package=RecordTransferPackage( + tenant_id=tenant_id, + package_id=package_id, + record_id=record_id, + record_revision=current.record_revision, + profile=current.profile, + manifest_sha256=current.manifest_sha256, + manifest=dict(current.manifest), + ), + purpose=_text(payload, "purpose"), + idempotency_key=_text(payload, "idempotency_key"), + institutional_context=dict(current.institutional_context), + ), + ) + except RecordContractError as exc: + raise RecordStoreError(str(exc)) from exc + if ( + receipt.package_id != package_id + or receipt.provider_id != current.provider_id + ): + raise RecordStoreError( + "The archive provider returned a receipt for another package." + ) + if not receipt.simulated: + raise RecordStoreError( + "The simulation provider returned a receipt that could be mistaken for a real custody transfer." + ) + current.superseded_at = recorded_at + status = { + "accepted": ("simulated_accepted" if receipt.simulated else "accepted"), + "rejected": "rejected", + "outcome_unknown": "outcome_unknown", + }[receipt.outcome] + receipt_data = { + "provider_id": receipt.provider_id, + "package_id": receipt.package_id, + "outcome": receipt.outcome, + "observed_at": _datetime_text(receipt.observed_at), + "external_reference": receipt.external_reference, + "retry_safe": receipt.retry_safe, + "simulated": receipt.simulated, + "metadata": dict(receipt.metadata), + } + row = RecordTransferPackageRevision( + tenant_id=tenant_id, + package_id=package_id, + record_id=record_id, + disposition_id=current.disposition_id, + revision=current.revision + 1, + previous_revision_id=current.id, + record_revision=current.record_revision, + provider_id=current.provider_id, + profile=current.profile, + status=status, + authority_mode=current.authority_mode, + manifest=dict(current.manifest), + manifest_sha256=current.manifest_sha256, + receipt=receipt_data, + receipt_sha256=receipt.receipt_sha256.removeprefix("sha256:"), + external_reference=receipt.external_reference, + recovery_operation_id=current_record_recovery_operation() + or current.recovery_operation_id, + simulated=receipt.simulated, + institutional_context=dict(current.institutional_context), + recorded_at=recorded_at, + changed_by=_actor(principal), + idempotency_key=_text(payload, "idempotency_key"), + request_sha256=request_hash, + ) + session.add(row) + _append_event( + session, + principal, + row=record, + event_type="record.transfer_package_dispatched", + summary=f"Transfer package result: {status}", + occurred_at=recorded_at, + purpose=_text(payload, "purpose"), + idempotency_key=row.idempotency_key, + request_hash=request_hash, + institutional_context=row.institutional_context, + payload={ + "package_id": package_id, + "provider_id": current.provider_id, + "status": status, + "receipt_sha256": row.receipt_sha256, + "simulated": row.simulated, + "external_reference": row.external_reference, + }, + ) + session.flush() + return _transfer_package_dict(row) + def create_volume( self, session: Session, @@ -770,7 +2037,7 @@ class SqlRecordRegistry: return self.registry.capability(name) except Exception as exc: raise RecordSourceUnavailableError( - f"Record source provider {name} is unavailable: {exc}" + f"Record capability {name} is unavailable: {exc}" ) from exc return None @@ -827,6 +2094,79 @@ def _current_record( return query.with_for_update().one_or_none() if lock else query.one_or_none() +def _required_current_record( + session: Session, + *, + tenant_id: str, + record_id: str, + lock: bool = False, +) -> RecordRevision: + row = _current_record(session, tenant_id, record_id, lock=lock) + if row is None: + raise RecordNotFoundError("Record not found.") + return row + + +def _current_hold( + session: Session, + tenant_id: str, + hold_id: str, + *, + lock: bool = False, +) -> RecordHoldRevision | None: + query = session.query(RecordHoldRevision).filter( + RecordHoldRevision.tenant_id == tenant_id, + RecordHoldRevision.hold_id == hold_id, + RecordHoldRevision.superseded_at.is_(None), + ) + return query.with_for_update().one_or_none() if lock else query.one_or_none() + + +def _current_disposition( + session: Session, + tenant_id: str, + disposition_id: str, + *, + lock: bool = False, +) -> RecordDispositionRevision | None: + query = session.query(RecordDispositionRevision).filter( + RecordDispositionRevision.tenant_id == tenant_id, + RecordDispositionRevision.disposition_id == disposition_id, + RecordDispositionRevision.superseded_at.is_(None), + ) + return query.with_for_update().one_or_none() if lock else query.one_or_none() + + +def _current_disposition_for_record( + session: Session, tenant_id: str, record_id: str +) -> RecordDispositionRevision | None: + return ( + session.query(RecordDispositionRevision) + .filter( + RecordDispositionRevision.tenant_id == tenant_id, + RecordDispositionRevision.record_id == record_id, + RecordDispositionRevision.superseded_at.is_(None), + RecordDispositionRevision.status != "withdrawn", + ) + .one_or_none() + ) + + +def _current_transfer_package( + session: Session, + tenant_id: str, + package_id: str, + *, + lock: bool = False, +) -> RecordTransferPackageRevision | None: + query = session.query(RecordTransferPackageRevision).filter( + RecordTransferPackageRevision.tenant_id == tenant_id, + RecordTransferPackageRevision.package_id == package_id, + RecordTransferPackageRevision.superseded_at.is_(None), + ) + return query.with_for_update().one_or_none() if lock else query.one_or_none() + + def _current_volume( session: Session, tenant_id: str, volume_id: str ) -> RecordVolumeRevision | None: @@ -893,6 +2233,98 @@ def _record_chronology( return query.order_by(RecordChronologyEntry.occurred_at.desc()).all() +def _record_holds( + session: Session, *, tenant_id: str, record_id: str +) -> Sequence[RecordHoldRevision]: + query = session.query(RecordHoldRevision).filter( + RecordHoldRevision.tenant_id == tenant_id, + RecordHoldRevision.record_id == record_id, + ) + return ( + apply_temporal_revision_filter( + query, + RecordHoldRevision, + valid_from=None, + valid_to=None, + ) + .order_by(RecordHoldRevision.recorded_at.desc()) + .all() + ) + + +def _record_dispositions( + session: Session, *, tenant_id: str, record_id: str +) -> Sequence[RecordDispositionRevision]: + query = session.query(RecordDispositionRevision).filter( + RecordDispositionRevision.tenant_id == tenant_id, + RecordDispositionRevision.record_id == record_id, + ) + return ( + apply_temporal_revision_filter( + query, + RecordDispositionRevision, + valid_from=None, + valid_to=None, + ) + .order_by(RecordDispositionRevision.recorded_at.desc()) + .all() + ) + + +def _record_transfer_packages( + session: Session, *, tenant_id: str, record_id: str +) -> Sequence[RecordTransferPackageRevision]: + query = session.query(RecordTransferPackageRevision).filter( + RecordTransferPackageRevision.tenant_id == tenant_id, + RecordTransferPackageRevision.record_id == record_id, + ) + return ( + apply_temporal_revision_filter( + query, + RecordTransferPackageRevision, + valid_from=None, + valid_to=None, + ) + .order_by(RecordTransferPackageRevision.recorded_at.desc()) + .all() + ) + + +def _active_holds( + session: Session, + *, + tenant_id: str, + record_id: str, + at: datetime, +) -> Sequence[RecordHoldRevision]: + return ( + session.query(RecordHoldRevision) + .filter( + RecordHoldRevision.tenant_id == tenant_id, + RecordHoldRevision.record_id == record_id, + RecordHoldRevision.superseded_at.is_(None), + RecordHoldRevision.status == "active", + RecordHoldRevision.effective_from <= at, + or_( + RecordHoldRevision.effective_to.is_(None), + RecordHoldRevision.effective_to > at, + ), + ) + .order_by(RecordHoldRevision.recorded_at) + .all() + ) + + +def _require_no_active_hold( + session: Session, tenant_id: str, record_id: str, *, at: datetime +) -> None: + holds = _active_holds(session, tenant_id=tenant_id, record_id=record_id, at=at) + if holds: + raise RecordConflictError( + f"Disposition is blocked by {len(holds)} active record hold(s)." + ) + + def _append_event( session: Session, principal: object, @@ -948,6 +2380,27 @@ def _replay_event( return event +def _replay_by_key( + session: Session, + model: type, + *, + tenant_id: str, + idempotency_key: str, + request_hash: str, +): + row = ( + session.query(model) + .filter( + model.tenant_id == tenant_id, + model.idempotency_key == idempotency_key, + ) + .one_or_none() + ) + if row is not None: + _verify_replay(row.request_sha256, request_hash) + return row + + def _verify_replay(actual_hash: str, requested_hash: str) -> None: if actual_hash != requested_hash: raise RecordConflictError( @@ -980,6 +2433,226 @@ def _validate_interval(valid_from: datetime | None, valid_to: datetime | None) - raise RecordStoreError("valid_to must be after valid_from.") +def _revise_record_state( + session: Session, + principal: object, + *, + current: RecordRevision, + recorded_at: datetime, + values: Mapping[str, object], + snapshot_patch: Mapping[str, object], +) -> tuple[RecordRevision, RecordIdentity]: + identity = session.get(RecordIdentity, current.identity_id) + if identity is None: + raise RecordStoreError("Record identity is missing.") + revision_values = _revision_values(current, {}) + revision_values.update(values) + current.superseded_at = recorded_at + row = RecordRevision( + tenant_id=current.tenant_id, + record_id=current.record_id, + identity_id=current.identity_id, + revision=current.revision + 1, + previous_revision_id=current.id, + recorded_at=recorded_at, + changed_by=_actor(principal), + search_text=current.search_text, + snapshot={**dict(current.snapshot), **dict(snapshot_patch)}, + **revision_values, + ) + session.add(row) + return row, identity + + +def _calculate_retention( + record_class: RecordClassRevision, + *, + closed_at: datetime, + explicit_trigger_at: datetime | None, +) -> tuple[datetime | None, datetime | None, dict[str, Any]]: + period_days = record_class.retention_period_days + trigger_mode = (record_class.closure_trigger or "record_closed").strip().lower() + if period_days is None: + return ( + None, + None, + { + "class_id": record_class.class_id, + "class_revision": record_class.revision, + "period_days": None, + "trigger": trigger_mode, + }, + ) + if trigger_mode in {"record_closed", "closed", "closure"}: + started_at = closed_at + elif trigger_mode in {"calendar_year_end", "year_end"}: + started_at = closed_at.replace( + year=closed_at.year + 1, + month=1, + day=1, + hour=0, + minute=0, + second=0, + microsecond=0, + ) + elif trigger_mode in {"calendar_month_end", "month_end"}: + if closed_at.month == 12: + started_at = closed_at.replace( + year=closed_at.year + 1, + month=1, + day=1, + hour=0, + minute=0, + second=0, + microsecond=0, + ) + else: + started_at = closed_at.replace( + month=closed_at.month + 1, + day=1, + hour=0, + minute=0, + second=0, + microsecond=0, + ) + elif trigger_mode == "explicit": + if explicit_trigger_at is None: + raise RecordStoreError( + "This record class requires an explicit retention trigger date." + ) + started_at = explicit_trigger_at + else: + raise RecordStoreError( + f"Unsupported record-class closure trigger: {trigger_mode}." + ) + due_at = started_at + timedelta(days=period_days) + return ( + started_at, + due_at, + { + "class_id": record_class.class_id, + "class_revision": record_class.revision, + "period_days": period_days, + "trigger": trigger_mode, + "started_at": _datetime_text(started_at), + "due_at": _datetime_text(due_at), + }, + ) + + +def _record_evidence_manifest( + session: Session, *, tenant_id: str, record: RecordRevision +) -> dict[str, Any]: + identity = session.get(RecordIdentity, record.identity_id) + if identity is None: + raise RecordStoreError("Record identity is missing.") + record_class = ( + session.query(RecordClassRevision) + .filter( + RecordClassRevision.tenant_id == tenant_id, + RecordClassRevision.class_id == record.class_id, + RecordClassRevision.revision + == int(record.retention_rule.get("class_revision") or 0), + ) + .one_or_none() + ) + if record_class is None: + record_class = _current_class(session, tenant_id, record.class_id) + items = ( + session.query(RecordItem) + .filter( + RecordItem.tenant_id == tenant_id, + RecordItem.record_id == record.record_id, + ) + .order_by(RecordItem.sequence) + .all() + ) + active_holds = _active_holds( + session, + tenant_id=tenant_id, + record_id=record.record_id, + at=datetime.now(UTC), + ) + return { + "record_id": record.record_id, + "record_number": identity.record_number, + "record_revision": record.revision, + "recorded_at": _datetime_text(record.recorded_at), + "state": record.state, + "class": ( + { + "class_id": record_class.class_id, + "revision": record_class.revision, + "key": record_class.key, + "retention_period_days": record_class.retention_period_days, + "closure_trigger": record_class.closure_trigger, + } + if record_class is not None + else {"class_id": record.class_id, "revision": None} + ), + "retention": { + "closed_at": _datetime_text(record.closed_at), + "started_at": _datetime_text(record.retention_started_at), + "due_at": _datetime_text(record.retention_due_at), + "rule": dict(record.retention_rule), + }, + "appraisal": dict(record.appraisal), + "items": [ + { + "item_id": item.id, + "sequence": item.sequence, + "source_module": item.source_module, + "resource_type": item.resource_type, + "resource_id": item.resource_id, + "source_revision": item.source_revision, + "authority_mode": item.authority_mode, + "content_sha256": item.content_sha256, + "filed_at": _datetime_text(item.filed_at), + } + for item in items + ], + "active_holds": [ + { + "hold_id": hold.hold_id, + "revision": hold.revision, + "authority": hold.authority, + "effective_from": _datetime_text(hold.effective_from), + "effective_to": _datetime_text(hold.effective_to), + } + for hold in active_holds + ], + "institutional_context": dict(record.institutional_context), + } + + +def _disposition_consequence_preview( + *, + action: str, + record: RecordRevision, + evidence: Mapping[str, object], + approvals_available: bool, +) -> dict[str, Any]: + descriptions = { + "retain": "Retain the record under its current classification without an external effect.", + "transfer": "Prepare an archive-neutral package; dispatch remains a separate governed action.", + "destroy": "Mark destruction as pending only; no content or source object is deleted by approval.", + "reclassify": "Return the record to an open state so a governed classification revision can follow.", + } + return { + "action": action, + "description": descriptions[action], + "record_revision": record.revision, + "item_count": len(list(evidence.get("items") or [])), + "destructive_effect": False, + "external_effect": False, + "requires_independent_approval": True, + "approval_capability_available": approvals_available, + "limitations": [ + "Approval changes lifecycle state only; destructive and external effects require separate evidence-bound operations." + ], + } + + def _revision_values( current: RecordRevision, payload: Mapping[str, object] ) -> dict[str, Any]: @@ -1010,6 +2683,12 @@ def _revision_values( if "valid_to" in payload else current.valid_to ), + "closed_at": current.closed_at, + "retention_started_at": current.retention_started_at, + "retention_due_at": current.retention_due_at, + "retention_rule": dict(current.retention_rule), + "appraisal_state": current.appraisal_state, + "appraisal": dict(current.appraisal), } @@ -1095,6 +2774,12 @@ def _record_dict(row: RecordRevision, identity: RecordIdentity) -> dict[str, Any "valid_from": _datetime_text(row.valid_from), "valid_to": _datetime_text(row.valid_to), "recorded_at": _datetime_text(row.recorded_at), + "closed_at": _datetime_text(row.closed_at), + "retention_started_at": _datetime_text(row.retention_started_at), + "retention_due_at": _datetime_text(row.retention_due_at), + "retention_rule": dict(row.retention_rule), + "appraisal_state": row.appraisal_state, + "appraisal": dict(row.appraisal), } @@ -1158,6 +2843,70 @@ def _chronology_dict(row: RecordChronologyEntry) -> dict[str, Any]: } +def _hold_dict(row: RecordHoldRevision) -> dict[str, Any]: + return { + "hold_id": row.hold_id, + "record_id": row.record_id, + "revision": row.revision, + "status": row.status, + "reason": row.reason, + "authority": row.authority, + "scope": dict(row.scope), + "effective_from": _datetime_text(row.effective_from), + "effective_to": _datetime_text(row.effective_to), + "released_at": _datetime_text(row.released_at), + "policy_refs": list(row.policy_refs), + "institutional_context": dict(row.institutional_context), + "recorded_at": _datetime_text(row.recorded_at), + "changed_by": row.changed_by, + } + + +def _disposition_dict(row: RecordDispositionRevision) -> dict[str, Any]: + return { + "disposition_id": row.disposition_id, + "record_id": row.record_id, + "revision": row.revision, + "action": row.action, + "status": row.status, + "reason": row.reason, + "subject_revision": row.subject_revision, + "subject_sha256": row.subject_sha256, + "consequence_preview": dict(row.consequence_preview), + "policy_refs": list(row.policy_refs), + "approval_request_id": row.approval_request_id, + "proposed_by": row.proposed_by, + "reviewed_by": row.reviewed_by, + "reviewed_at": _datetime_text(row.reviewed_at), + "institutional_context": dict(row.institutional_context), + "recorded_at": _datetime_text(row.recorded_at), + } + + +def _transfer_package_dict(row: RecordTransferPackageRevision) -> dict[str, Any]: + return { + "package_id": row.package_id, + "record_id": row.record_id, + "disposition_id": row.disposition_id, + "revision": row.revision, + "record_revision": row.record_revision, + "provider_id": row.provider_id, + "profile": row.profile, + "status": row.status, + "authority_mode": row.authority_mode, + "manifest": dict(row.manifest), + "manifest_sha256": row.manifest_sha256, + "receipt": dict(row.receipt), + "receipt_sha256": row.receipt_sha256, + "external_reference": row.external_reference, + "recovery_operation_id": row.recovery_operation_id, + "simulated": row.simulated, + "institutional_context": dict(row.institutional_context), + "recorded_at": _datetime_text(row.recorded_at), + "changed_by": row.changed_by, + } + + def _filing_request_mapping(request: RecordFilingRequest) -> dict[str, object]: return { "tenant_id": request.tenant_id, diff --git a/tests/fixtures/service_to_decision_journey.json b/tests/fixtures/service_to_decision_journey.json new file mode 100644 index 0000000..07d6df6 --- /dev/null +++ b/tests/fixtures/service_to_decision_journey.json @@ -0,0 +1,18 @@ +{ + "id": "assisted-service-to-decision-record", + "description": "A completed assisted service is filed, closed, held, appraised, independently approved, and exercised through the archive simulation boundary.", + "record": { + "record_id": "record-1", + "record_number": "2026/0001", + "class_id": "class-permit" + }, + "expected": { + "closed_state": "retention_running", + "appraised_state": "appraised", + "approved_state": "transfer_pending", + "hold_blocks_disposition": true, + "disposition_status": "approved", + "transfer_status": "simulated_accepted", + "custody_transferred": false + } +} diff --git a/tests/test_migrations.py b/tests/test_migrations.py index 8ab4f91..d5b04ab 100644 --- a/tests/test_migrations.py +++ b/tests/test_migrations.py @@ -30,9 +30,12 @@ class RecordsMigrationTests(unittest.TestCase): "record_chronology_entries", "record_class_revisions", "record_file_plan_revisions", + "record_hold_revisions", "record_identities", "record_items", + "record_disposition_revisions", "record_revisions", + "record_transfer_package_revisions", "record_volume_revisions", }.issubset(inspector.get_table_names()) ) @@ -42,7 +45,7 @@ class RecordsMigrationTests(unittest.TestCase): ) with engine.connect() as connection: self.assertIn( - "6e4a2c8f1d9b", + "7f5b3d9a2c1e", set(MigrationContext.configure(connection).get_current_heads()), ) finally: diff --git a/tests/test_records.py b/tests/test_records.py index 1aff738..4c39f32 100644 --- a/tests/test_records.py +++ b/tests/test_records.py @@ -2,16 +2,22 @@ from __future__ import annotations from dataclasses import dataclass from datetime import UTC, datetime, timedelta +import json +from pathlib import Path import unittest from sqlalchemy import create_engine from sqlalchemy.orm import Session from govoplan_core.core.records import ( + RecordArchiveProviderState, + RecordArchiveReceipt, RecordFilingRequest, RecordSourceLocator, RecordSourceReference, ) +from govoplan_core.core.recovery import RecoveryCheckpoint, RecoveryOperation +from govoplan_core.core.approvals import ApprovalCheck, ApprovalRequestRef from govoplan_core.core.temporal import ( TemporalDataContext, bind_temporal_data_context, @@ -20,12 +26,16 @@ from govoplan_core.core.temporal import ( from govoplan_records.backend.db.models import ( RecordChronologyEntry, RecordClassRevision, + RecordDispositionRevision, RecordFilePlanRevision, + RecordHoldRevision, RecordIdentity, RecordItem, RecordRevision, + RecordTransferPackageRevision, RecordVolumeRevision, ) +from govoplan_records.backend.archive import SimulatedRecordArchiveProvider from govoplan_records.backend.service import RecordConflictError, SqlRecordRegistry @@ -63,21 +73,114 @@ class SourceProvider: ) +class UnknownOutcomeArchiveProvider: + provider_id = "unknown_simulation" + + def __init__(self) -> None: + self.dispatch_count = 0 + + def state(self): + return RecordArchiveProviderState( + provider_id=self.provider_id, + label="Unknown-outcome simulation", + profiles=("govoplan-unknown-outcome-v1",), + authority_modes=("linked_reference",), + healthy=True, + checked_at=NOW, + simulated=True, + ) + + def dispatch(self, session, principal, *, request): + del session, principal + self.dispatch_count += 1 + return RecordArchiveReceipt( + provider_id=self.provider_id, + package_id=request.package.package_id, + outcome="outcome_unknown", + observed_at=NOW + timedelta(minutes=12), + receipt_sha256="b" * 64, + retry_safe=False, + simulated=True, + metadata={"custody_transferred": False}, + ) + + class Registry: def __init__(self) -> None: self.provider = SourceProvider() + self.archive_provider = SimulatedRecordArchiveProvider() + self.unknown_archive_provider = UnknownOutcomeArchiveProvider() + self.approvals = ApprovalProvider() def capability_names(self): - return ("records.source.files",) + return ( + "records.source.files", + "records.archive.simulation", + "records.archive.unknown_simulation", + ) def tenant_capability(self, name, session, *, tenant_id): del session - return ( - self.provider - if name == "records.source.files" and tenant_id == "tenant-1" - else None + if tenant_id != "tenant-1": + return None + return { + "records.source.files": self.provider, + "records.archive.simulation": self.archive_provider, + "records.archive.unknown_simulation": self.unknown_archive_provider, + "approvals.requests": self.approvals, + }.get(name) + + +class ApprovalProvider: + def __init__(self) -> None: + self.approved = False + self.request = None + + def create_request(self, session, principal, *, command, idempotency_key): + del session, principal, idempotency_key + self.request = command + return ApprovalRequestRef(id="approval-1", revision=1, state="pending") + + def check_approved( + self, + session, + principal, + *, + request_id, + subject_module, + subject_type, + subject_id, + subject_version, + subject_digest, + ): + del session, principal + return ApprovalCheck( + request_id=request_id, + revision=2 if self.approved else 1, + state="approved" if self.approved else "pending", + approved=self.approved, + subject_module=subject_module, + subject_type=subject_type, + subject_id=subject_id, + subject_version=subject_version, + subject_digest=subject_digest, ) + def create_template(self, *args, **kwargs): + raise NotImplementedError + + def revise_template(self, *args, **kwargs): + raise NotImplementedError + + def publish_template(self, *args, **kwargs): + raise NotImplementedError + + def get_request(self, *args, **kwargs): + return None + + def decide(self, *args, **kwargs): + raise NotImplementedError + class RecordsTests(unittest.TestCase): def setUp(self) -> None: @@ -90,6 +193,11 @@ class RecordsTests(unittest.TestCase): RecordVolumeRevision.__table__, RecordItem.__table__, RecordChronologyEntry.__table__, + RecordHoldRevision.__table__, + RecordDispositionRevision.__table__, + RecordTransferPackageRevision.__table__, + RecoveryOperation.__table__, + RecoveryCheckpoint.__table__, ): table.create(self.engine) self.session = Session(self.engine) @@ -400,6 +508,309 @@ class RecordsTests(unittest.TestCase): self.records.source_providers(self.session, self.principal), ) + def test_governed_lifecycle_hold_disposition_and_transfer_simulation(self) -> None: + fixture = json.loads( + ( + Path(__file__).parent / "fixtures/service_to_decision_journey.json" + ).read_text(encoding="utf-8") + ) + expected = fixture["expected"] + self._create_record() + closed = self.records.close_record( + self.session, + self.principal, + record_id="record-1", + payload={ + "expected_revision": 1, + "purpose": "close completed case record", + "reason": "The administrative decision is final.", + "recorded_at": NOW + timedelta(minutes=2), + "idempotency_key": "record-close-1", + }, + ) + self.assertEqual(expected["closed_state"], closed["state"]) + self.assertEqual( + NOW + timedelta(minutes=2, days=3650), + datetime.fromisoformat(str(closed["retention_due_at"])), + ) + appraised = self.records.appraise_record( + self.session, + self.principal, + record_id="record-1", + payload={ + "expected_revision": 2, + "outcome": "transfer", + "purpose": "appraise completed record", + "reason": "Transfer to the institutional archive after review.", + "policy_refs": ["records-policy:v1"], + "override_retention_not_due": True, + "recorded_at": NOW + timedelta(minutes=3), + "idempotency_key": "record-appraise-1", + }, + ) + self.assertEqual(expected["appraised_state"], appraised["state"]) + hold = self.records.apply_hold( + self.session, + self.principal, + record_id="record-1", + payload={ + "expected_record_revision": 3, + "reason": "Pending judicial review.", + "authority": "Court order 2026-17", + "purpose": "preserve evidence", + "recorded_at": NOW + timedelta(minutes=4), + "idempotency_key": "record-hold-1", + }, + ) + with self.assertRaisesRegex(RecordConflictError, "active record hold"): + self.records.propose_disposition( + self.session, + self.principal, + record_id="record-1", + payload={ + "expected_record_revision": 3, + "action": "transfer", + "reason": "Transfer after retention.", + "purpose": "dispose record", + "recorded_at": NOW + timedelta(minutes=5), + "idempotency_key": "record-disposition-blocked", + }, + ) + self.records.release_hold( + self.session, + self.principal, + record_id="record-1", + hold_id=str(hold["hold_id"]), + payload={ + "expected_hold_revision": 1, + "reason": "Judicial review is complete.", + "purpose": "resume disposition", + "recorded_at": NOW + timedelta(minutes=6), + "idempotency_key": "record-hold-release-1", + }, + ) + disposition = self.records.propose_disposition( + self.session, + self.principal, + record_id="record-1", + payload={ + "disposition_id": "disposition-1", + "expected_record_revision": 3, + "action": "transfer", + "reason": "Transfer after independent review.", + "purpose": "dispose record", + "policy_refs": ["records-policy:v1"], + "recorded_at": NOW + timedelta(minutes=7), + "idempotency_key": "record-disposition-1", + }, + ) + self.assertEqual("review_pending", disposition["status"]) + self.assertFalse(disposition["consequence_preview"]["external_effect"]) + self.records.registry.approvals.approved = True + finalized = self.records.finalize_disposition( + self.session, + self.principal, + record_id="record-1", + disposition_id="disposition-1", + payload={ + "expected_disposition_revision": 1, + "purpose": "approve disposition", + "recorded_at": NOW + timedelta(minutes=8), + "idempotency_key": "record-disposition-finalize-1", + }, + ) + self.assertEqual(expected["approved_state"], finalized["record"]["state"]) + package = self.records.prepare_transfer_package( + self.session, + self.principal, + record_id="record-1", + payload={ + "package_id": "package-1", + "disposition_id": "disposition-1", + "expected_record_revision": 4, + "provider_id": "simulation", + "profile": "govoplan-simulation-v1", + "purpose": "validate archive transfer", + "recorded_at": NOW + timedelta(minutes=9), + "idempotency_key": "record-package-1", + }, + ) + self.assertEqual("prepared", package["status"]) + receipt = self.records.dispatch_transfer_package( + self.session, + self.principal, + record_id="record-1", + package_id="package-1", + payload={ + "expected_package_revision": 1, + "purpose": "validate archive transfer", + "recorded_at": NOW + timedelta(minutes=10), + "idempotency_key": "record-package-dispatch-1", + }, + ) + self.session.commit() + self.assertEqual(expected["transfer_status"], receipt["status"]) + self.assertTrue(receipt["simulated"]) + self.assertEqual( + expected["custody_transferred"], + receipt["receipt"]["metadata"]["custody_transferred"], + ) + + detail = self.records.get_record( + self.session, self.principal, record_id="record-1" + ) + self.assertEqual("released", detail["holds"][0]["status"]) + self.assertEqual( + expected["disposition_status"], detail["dispositions"][0]["status"] + ) + self.assertEqual("simulated_accepted", detail["transfer_packages"][0]["status"]) + + unknown_package = self.records.prepare_transfer_package( + self.session, + self.principal, + record_id="record-1", + payload={ + "package_id": "package-unknown", + "disposition_id": "disposition-1", + "expected_record_revision": 4, + "provider_id": "unknown_simulation", + "profile": "govoplan-unknown-outcome-v1", + "purpose": "exercise unknown transfer recovery", + "recorded_at": NOW + timedelta(minutes=11), + "idempotency_key": "record-package-unknown", + }, + ) + unknown_receipt = self.records.dispatch_transfer_package( + self.session, + self.principal, + record_id="record-1", + package_id="package-unknown", + payload={ + "expected_package_revision": unknown_package["revision"], + "purpose": "exercise unknown transfer recovery", + "recorded_at": NOW + timedelta(minutes=12), + "idempotency_key": "record-package-unknown-dispatch", + }, + ) + self.assertEqual("outcome_unknown", unknown_receipt["status"]) + self.assertFalse(unknown_receipt["receipt"]["retry_safe"]) + with self.assertRaisesRegex( + ValueError, "unknown outcomes require reconciliation" + ): + self.records.dispatch_transfer_package( + self.session, + self.principal, + record_id="record-1", + package_id="package-unknown", + payload={ + "expected_package_revision": 2, + "purpose": "must not retry unknown transfer", + "recorded_at": NOW + timedelta(minutes=13), + "idempotency_key": "record-package-unknown-retry", + }, + ) + self.assertEqual( + 1, self.records.registry.unknown_archive_provider.dispatch_count + ) + self.assertEqual( + "transfer_pending", + self.records.get_record(self.session, self.principal, record_id="record-1")[ + "record" + ]["state"], + ) + recovery = self.records.recovery_status( + self.session, self.principal, record_id="record-1" + ) + self.assertTrue(recovery["healthy"]) + self.assertTrue(recovery["package_checks"][0]["manifest_verified"]) + + def test_disposition_waits_for_optional_approvals_and_can_be_withdrawn( + self, + ) -> None: + self.records.registry.approvals = None + self._create_record() + self.records.close_record( + self.session, + self.principal, + record_id="record-1", + payload={ + "expected_revision": 1, + "purpose": "close record", + "reason": "Work completed.", + "recorded_at": NOW + timedelta(minutes=2), + "idempotency_key": "close-without-approvals", + }, + ) + self.records.appraise_record( + self.session, + self.principal, + record_id="record-1", + payload={ + "expected_revision": 2, + "outcome": "retain", + "purpose": "appraise record", + "reason": "Retain permanently.", + "override_retention_not_due": True, + "recorded_at": NOW + timedelta(minutes=3), + "idempotency_key": "appraise-without-approvals", + }, + ) + proposal = self.records.propose_disposition( + self.session, + self.principal, + record_id="record-1", + payload={ + "expected_record_revision": 3, + "action": "retain", + "purpose": "dispose record", + "reason": "Retain permanently.", + "recorded_at": NOW + timedelta(minutes=4), + "idempotency_key": "proposal-without-approvals", + }, + ) + self.assertEqual("review_unavailable", proposal["status"]) + with self.assertRaisesRegex(ValueError, "not awaiting an available approval"): + self.records.finalize_disposition( + self.session, + self.principal, + record_id="record-1", + disposition_id=str(proposal["disposition_id"]), + payload={ + "expected_disposition_revision": 1, + "purpose": "finalize disposition", + "recorded_at": NOW + timedelta(minutes=5), + "idempotency_key": "finalize-without-approvals", + }, + ) + withdrawn = self.records.withdraw_disposition( + self.session, + self.principal, + record_id="record-1", + disposition_id=str(proposal["disposition_id"]), + payload={ + "expected_disposition_revision": 1, + "purpose": "correct disposition", + "reason": "The appraisal evidence needs correction.", + "recorded_at": NOW + timedelta(minutes=6), + "idempotency_key": "withdraw-without-approvals", + }, + ) + self.assertEqual("withdrawn", withdrawn["status"]) + replacement = self.records.propose_disposition( + self.session, + self.principal, + record_id="record-1", + payload={ + "expected_record_revision": 3, + "action": "retain", + "purpose": "dispose record", + "reason": "Corrected permanent-retention proposal.", + "recorded_at": NOW + timedelta(minutes=7), + "idempotency_key": "replacement-without-approvals", + }, + ) + self.assertEqual("review_unavailable", replacement["status"]) + if __name__ == "__main__": unittest.main() diff --git a/tests/test_recovery.py b/tests/test_recovery.py new file mode 100644 index 0000000..853abcf --- /dev/null +++ b/tests/test_recovery.py @@ -0,0 +1,181 @@ +from __future__ import annotations + +from pathlib import Path +import sqlite3 +import tempfile +import unittest + +from sqlalchemy import create_engine +from sqlalchemy.orm import Session + +from govoplan_core.core.recovery import ( + RecoveryCheckpoint, + RecoveryOperation, + RecoveryStatus, + verify_recovery_evidence_chain, +) +from govoplan_core.core.runtime_coordination import ( + DistributedLease, + RuntimeIdentity, + bind_process_runtime_identity, +) +from govoplan_records.backend.recovery import ( + RecordRecoveryError, + begin_record_atomic_recovery, +) + + +class RecordsRecoveryTests(unittest.TestCase): + def setUp(self) -> None: + self.directory = tempfile.TemporaryDirectory(prefix="records-recovery-") + self.database_path = Path(self.directory.name) / "recovery.db" + self.engine = create_engine(f"sqlite:///{self.database_path}") + for table in ( + DistributedLease.__table__, + RecoveryOperation.__table__, + RecoveryCheckpoint.__table__, + ): + table.create(self.engine) + self.session = Session(self.engine) + bind_process_runtime_identity( + RuntimeIdentity( + installation_id="test-installation", + node_id="records-test-node", + incarnation="11111111-1111-4111-8111-111111111111", + role="api", + software_version="test", + composition_hash="a" * 64, + ) + ) + + def tearDown(self) -> None: + bind_process_runtime_identity(None) + self.session.close() + self.engine.dispose() + self.directory.cleanup() + + def test_atomic_evidence_chain_replays_only_the_same_request(self) -> None: + started = begin_record_atomic_recovery( + self.session, + tenant_id="tenant-1", + operation_type="record.close", + idempotency_key="close-1", + request={"record_id": "record-1", "expected_revision": 1}, + resource_type="record", + resource_id="record-1", + ) + started.commit_success( + self.session, + result={"record_id": "record-1", "revision": 2}, + resource_id="record-1", + ) + + operation = self.session.get(RecoveryOperation, started.operation_id) + self.assertIsNotNone(operation) + self.assertEqual(RecoveryStatus.SUCCEEDED.value, operation.status) + self.assertTrue( + verify_recovery_evidence_chain(self.session, started.operation_id) + ) + + replay = begin_record_atomic_recovery( + self.session, + tenant_id="tenant-1", + operation_type="record.close", + idempotency_key="close-1", + request={"record_id": "record-1", "expected_revision": 1}, + resource_type="record", + resource_id="record-1", + ) + self.assertTrue(replay.replayed) + self.assertIsNone(replay.operation) + + with self.assertRaisesRegex(RecordRecoveryError, "not started"): + begin_record_atomic_recovery( + self.session, + tenant_id="tenant-1", + operation_type="record.close", + idempotency_key="close-1", + request={"record_id": "record-1", "expected_revision": 9}, + resource_type="record", + resource_id="record-1", + ) + + def test_unresolved_resource_blocks_another_runtime_effect(self) -> None: + started = begin_record_atomic_recovery( + self.session, + tenant_id="tenant-1", + operation_type="transfer.simulate", + idempotency_key="dispatch-1", + request={"record_id": "record-1", "package_id": "package-1"}, + resource_type="record", + resource_id="record-1", + ) + bind_process_runtime_identity( + RuntimeIdentity( + installation_id="test-installation", + node_id="other-records-test-node", + incarnation="22222222-2222-4222-8222-222222222222", + role="api", + software_version="test", + composition_hash="a" * 64, + ) + ) + with self.assertRaisesRegex(RecordRecoveryError, "Another runtime"): + begin_record_atomic_recovery( + self.session, + tenant_id="tenant-1", + operation_type="record.reopen", + idempotency_key="reopen-1", + request={"record_id": "record-1"}, + resource_type="record", + resource_id="record-1", + ) + bind_process_runtime_identity( + RuntimeIdentity( + installation_id="test-installation", + node_id="records-test-node", + incarnation="11111111-1111-4111-8111-111111111111", + role="api", + software_version="test", + composition_hash="a" * 64, + ) + ) + started.fail(summary="Test cleanup", error_type="TestInterruption") + + def test_recovery_evidence_survives_database_restore(self) -> None: + started = begin_record_atomic_recovery( + self.session, + tenant_id="tenant-1", + operation_type="record.close", + idempotency_key="restore-close-1", + request={"record_id": "record-restore", "expected_revision": 1}, + resource_type="record", + resource_id="record-restore", + ) + started.commit_success( + self.session, + result={"record_id": "record-restore", "revision": 2}, + resource_id="record-restore", + ) + + restored_path = Path(self.directory.name) / "restored.db" + with ( + sqlite3.connect(self.database_path) as source, + sqlite3.connect(restored_path) as target, + ): + source.backup(target) + restored_engine = create_engine(f"sqlite:///{restored_path}") + try: + with Session(restored_engine) as restored: + operation = restored.get(RecoveryOperation, started.operation_id) + self.assertIsNotNone(operation) + self.assertEqual(RecoveryStatus.SUCCEEDED.value, operation.status) + self.assertTrue( + verify_recovery_evidence_chain(restored, started.operation_id) + ) + finally: + restored_engine.dispose() + + +if __name__ == "__main__": + unittest.main() diff --git a/webui/src/api/records.ts b/webui/src/api/records.ts index 3e61ab3..ce90c0a 100644 --- a/webui/src/api/records.ts +++ b/webui/src/api/records.ts @@ -12,6 +12,12 @@ export type FilePlanNode = { valid_from?: string | null; valid_to?: string | null; recorded_at: string; + closed_at?: string | null; + retention_started_at?: string | null; + retention_due_at?: string | null; + retention_rule: Record; + appraisal_state?: string | null; + appraisal: Record; institutional_context: Record; }; @@ -99,11 +105,63 @@ export type RecordChronology = { payload: Record; }; +export type RecordHold = { + hold_id: string; + record_id: string; + revision: number; + status: string; + reason: string; + authority: string; + scope: Record; + effective_from: string; + effective_to?: string | null; + released_at?: string | null; + policy_refs: string[]; + recorded_at: string; +}; + +export type RecordDisposition = { + disposition_id: string; + record_id: string; + revision: number; + action: "retain" | "transfer" | "destroy" | "reclassify"; + status: string; + reason: string; + subject_revision: number; + subject_sha256: string; + consequence_preview: Record; + policy_refs: string[]; + approval_request_id?: string | null; + proposed_by?: string | null; + reviewed_by?: string | null; + reviewed_at?: string | null; + recorded_at: string; +}; + +export type RecordTransferPackage = { + package_id: string; + record_id: string; + disposition_id: string; + revision: number; + record_revision: number; + provider_id: string; + profile: string; + status: string; + manifest_sha256: string; + receipt_sha256?: string | null; + external_reference?: string | null; + simulated: boolean; + recorded_at: string; +}; + export type RecordDetail = { record: RecordEntry; volumes: RecordVolume[]; items: RecordItem[]; chronology: RecordChronology[]; + holds: RecordHold[]; + dispositions: RecordDisposition[]; + transfer_packages: RecordTransferPackage[]; access_explanation: { decision: string; reason: string; @@ -120,6 +178,19 @@ export type RecordSourceProvider = { resource_types: string[]; }; +export type RecordArchiveProvider = { + id: string; + label: string; + profiles: string[]; + authority_modes: string[]; + healthy: boolean; + checked_at: string; + last_success_at?: string | null; + freshness_seconds?: number | null; + limitations: string[]; + simulated: boolean; +}; + export function listRecords( settings: ApiSettings, options: { @@ -154,6 +225,24 @@ export function getRecordSources(settings: ApiSettings, signal?: AbortSignal): P return apiFetch(settings, "/api/v1/records/sources", { signal }); } +export function getRecordArchiveProviders(settings: ApiSettings, signal?: AbortSignal): Promise<{ providers: RecordArchiveProvider[] }> { + return apiFetch(settings, "/api/v1/records/archive-providers", { signal }); +} + +export function writeFilePlanNode(settings: ApiSettings, payload: Record): Promise { + return apiFetch(settings, "/api/v1/records/catalog/file-plan", { + method: "POST", + body: JSON.stringify(payload) + }); +} + +export function writeRecordClass(settings: ApiSettings, payload: Record): Promise { + return apiFetch(settings, "/api/v1/records/catalog/classes", { + method: "POST", + body: JSON.stringify(payload) + }); +} + export function createRecord(settings: ApiSettings, payload: Record): Promise { return apiFetch(settings, "/api/v1/records", { method: "POST", @@ -182,3 +271,99 @@ export function fileRecordItem( body: JSON.stringify(payload) }); } + +export function createRecordVolume( + settings: ApiSettings, + recordId: string, + payload: Record +): Promise { + return recordMutation(settings, recordId, "volumes", payload); +} + +export function closeRecord(settings: ApiSettings, recordId: string, payload: Record): Promise { + return recordMutation(settings, recordId, "close", payload); +} + +export function reopenRecord(settings: ApiSettings, recordId: string, payload: Record): Promise { + return recordMutation(settings, recordId, "reopen", payload); +} + +export function appraiseRecord(settings: ApiSettings, recordId: string, payload: Record): Promise { + return recordMutation(settings, recordId, "appraise", payload); +} + +export function applyRecordHold(settings: ApiSettings, recordId: string, payload: Record): Promise { + return recordMutation(settings, recordId, "holds", payload); +} + +export function releaseRecordHold( + settings: ApiSettings, + recordId: string, + holdId: string, + payload: Record +): Promise { + return recordMutation(settings, recordId, `holds/${encodeURIComponent(holdId)}/release`, payload); +} + +export function proposeRecordDisposition( + settings: ApiSettings, + recordId: string, + payload: Record +): Promise { + return recordMutation(settings, recordId, "dispositions", payload); +} + +export function finalizeRecordDisposition( + settings: ApiSettings, + recordId: string, + dispositionId: string, + payload: Record +): Promise<{ disposition: RecordDisposition; record: RecordEntry }> { + return recordMutation(settings, recordId, `dispositions/${encodeURIComponent(dispositionId)}/finalize`, payload); +} + +export function withdrawRecordDisposition( + settings: ApiSettings, + recordId: string, + dispositionId: string, + payload: Record +): Promise { + return recordMutation(settings, recordId, `dispositions/${encodeURIComponent(dispositionId)}/withdraw`, payload); +} + +export function prepareRecordTransfer( + settings: ApiSettings, + recordId: string, + payload: Record +): Promise { + return recordMutation(settings, recordId, "transfer-packages", payload); +} + +export function dispatchRecordTransfer( + settings: ApiSettings, + recordId: string, + packageId: string, + payload: Record +): Promise { + return recordMutation(settings, recordId, `transfer-packages/${encodeURIComponent(packageId)}/dispatch`, payload); +} + +export function getRecordRecoveryStatus( + settings: ApiSettings, + recordId: string, + signal?: AbortSignal +): Promise> { + return apiFetch(settings, `/api/v1/records/${encodeURIComponent(recordId)}/recovery`, { signal }); +} + +function recordMutation( + settings: ApiSettings, + recordId: string, + path: string, + payload: Record +): Promise { + return apiFetch(settings, `/api/v1/records/${encodeURIComponent(recordId)}/${path}`, { + method: "POST", + body: JSON.stringify(payload) + }); +} diff --git a/webui/src/features/records/RecordCatalogDialog.tsx b/webui/src/features/records/RecordCatalogDialog.tsx new file mode 100644 index 0000000..eca4cec --- /dev/null +++ b/webui/src/features/records/RecordCatalogDialog.tsx @@ -0,0 +1,174 @@ +import { useEffect, useState, type FormEvent } from "react"; +import { + Button, + Dialog, + DismissibleAlert, + FormField, + type PlatformRouteContext +} from "@govoplan/core-webui"; +import { + writeFilePlanNode, + writeRecordClass, + type RecordCatalog +} from "../../api/records"; + + +export function RecordCatalogDialog({ + open, + settings, + catalog, + onClose, + onSaved +}: { + open: boolean; + settings: PlatformRouteContext["settings"]; + catalog: RecordCatalog; + onClose: () => void; + onSaved: () => void; +}) { + const [kind, setKind] = useState<"file-plan" | "class">("file-plan"); + const [selectedId, setSelectedId] = useState(""); + const [identifier, setIdentifier] = useState(""); + const [code, setCode] = useState(""); + const [label, setLabel] = useState(""); + const [description, setDescription] = useState(""); + const [parentId, setParentId] = useState(""); + const [filePlanNodeId, setFilePlanNodeId] = useState(""); + const [retentionDays, setRetentionDays] = useState(""); + const [closureTrigger, setClosureTrigger] = useState("record_closed"); + const [saving, setSaving] = useState(false); + const [error, setError] = useState(""); + + useEffect(() => { + if (!open) return; + setKind("file-plan"); + setSelectedId(""); + resetFields(); + }, [open]); + + useEffect(() => { + if (!selectedId) { + resetFields(); + return; + } + if (kind === "file-plan") { + const item = catalog.file_plan.find((candidate) => candidate.node_id === selectedId); + if (!item) return; + setIdentifier(item.node_id); + setCode(item.code); + setLabel(item.label); + setDescription(item.description ?? ""); + setParentId(item.parent_node_id ?? ""); + } else { + const item = catalog.classes.find((candidate) => candidate.class_id === selectedId); + if (!item) return; + setIdentifier(item.class_id); + setCode(item.key); + setLabel(item.label); + setDescription(item.description ?? ""); + setFilePlanNodeId(item.file_plan_node_id); + setRetentionDays(item.retention_period_days == null ? "" : String(item.retention_period_days)); + setClosureTrigger(item.closure_trigger ?? "record_closed"); + } + }, [catalog.classes, catalog.file_plan, kind, selectedId]); + + function resetFields() { + setIdentifier(""); + setCode(""); + setLabel(""); + setDescription(""); + setParentId(""); + setFilePlanNodeId(catalog.file_plan[0]?.node_id ?? ""); + setRetentionDays(""); + setClosureTrigger("record_closed"); + setError(""); + } + + async function submit(event: FormEvent) { + event.preventDefault(); + const current = kind === "file-plan" + ? catalog.file_plan.find((item) => item.node_id === selectedId) + : catalog.classes.find((item) => item.class_id === selectedId); + setSaving(true); + setError(""); + try { + if (kind === "file-plan") { + await writeFilePlanNode(settings, { + node_id: identifier.trim(), + code: code.trim(), + label: label.trim(), + description: description.trim() || null, + parent_node_id: parentId || null, + active: true, + recorded_at: new Date().toISOString(), + expected_revision: current?.revision, + idempotency_key: randomId(), + institutional_context: {} + }); + } else { + await writeRecordClass(settings, { + class_id: identifier.trim(), + file_plan_node_id: filePlanNodeId, + key: code.trim(), + label: label.trim(), + description: description.trim() || null, + metadata_requirements: [], + allowed_source_types: [], + retention_period_days: retentionDays === "" ? null : Number(retentionDays), + closure_trigger: closureTrigger, + access_mode: "tenant", + active: true, + recorded_at: new Date().toISOString(), + expected_revision: current?.revision, + idempotency_key: randomId(), + institutional_context: {} + }); + } + onSaved(); + } catch (reason) { + setError(reason instanceof Error && reason.message ? reason.message : "The Records catalog could not be saved."); + } finally { + setSaving(false); + } + } + + const isValid = identifier.trim() && code.trim() && label.trim() && (kind === "file-plan" || filePlanNodeId); + const options = kind === "file-plan" ? catalog.file_plan : catalog.classes; + return ( + } + > + {error && {error}} +
+ + + + + + + setIdentifier(event.target.value)} disabled={Boolean(selectedId)} required /> + setCode(event.target.value)} required /> + setLabel(event.target.value)} required /> + {kind === "file-plan" ? : <> setRetentionDays(event.target.value)} />} +