feat(records): complete eAkte reference journey
Module Package Release / publish-packages (push) Successful in 12s

This commit is contained in:
2026-08-22 18:42:09 +02:00
parent 38f203a906
commit a5eee2c23f
15 changed files with 1020 additions and 36 deletions
+503 -1
View File
@@ -2,8 +2,11 @@ from __future__ import annotations
from dataclasses import dataclass
from datetime import UTC, datetime, timedelta
import hashlib
import json
from pathlib import Path
import sqlite3
import tempfile
import unittest
from sqlalchemy import create_engine
@@ -73,6 +76,48 @@ class SourceProvider:
)
class ReferenceJourneySourceProvider:
def __init__(self, provider_id: str, resource_types: tuple[str, ...]) -> None:
self.provider_id = provider_id
self._resource_types = resource_types
def resource_types(self):
return self._resource_types
def resolve(self, session, principal, *, locator, purpose):
del session, purpose
if principal.tenant_id != locator.tenant_id:
raise ValueError("Source access denied.")
if locator.source_module != self.provider_id:
raise ValueError("The source module does not match this provider.")
if locator.resource_type not in self._resource_types:
raise ValueError("The source type is not supported by this provider.")
metadata = dict(locator.metadata)
reference = dict(metadata.get("reference") or {})
digest_input = ":".join(
(
locator.source_module,
locator.resource_type,
locator.resource_id,
locator.source_revision,
)
).encode("utf-8")
return RecordSourceReference(
locator=locator,
label=str(reference.get("label") or locator.resource_id),
authority_mode=str(
reference.get("authority_mode") or "linked_reference"
),
content_sha256=hashlib.sha256(digest_input).hexdigest(),
content_type=str(reference.get("content_type") or "application/json"),
size_bytes=int(reference.get("size_bytes") or 0),
valid_from=NOW,
recorded_at=NOW,
launch_url=str(reference.get("launch_url") or "") or None,
metadata=metadata,
)
class UnknownOutcomeArchiveProvider:
provider_id = "unknown_simulation"
@@ -131,6 +176,33 @@ class Registry:
}.get(name)
class ReferenceJourneyRegistry(Registry):
def __init__(self) -> None:
super().__init__()
self.source_providers_by_module = {
"files": ReferenceJourneySourceProvider("files", ("file_version",)),
"forms_runtime": ReferenceJourneySourceProvider(
"forms_runtime", ("form_submission_revision",)
),
"cases": ReferenceJourneySourceProvider("cases", ("case_revision",)),
"decisions": ReferenceJourneySourceProvider(
"decisions", ("decision_revision",)
),
}
def capability_names(self):
return (
*(f"records.source.{name}" for name in self.source_providers_by_module),
"records.archive.simulation",
"records.archive.unknown_simulation",
)
def tenant_capability(self, name, session, *, tenant_id):
if name.startswith("records.source.") and tenant_id == "tenant-1":
return self.source_providers_by_module.get(name.removeprefix("records.source."))
return super().tenant_capability(name, session, tenant_id=tenant_id)
class ApprovalProvider:
def __init__(self) -> None:
self.approved = False
@@ -231,7 +303,12 @@ class RecordsTests(unittest.TestCase):
"file_plan_node_id": "plan-permits",
"key": "permit.application",
"label": "Permit application",
"allowed_source_types": ["files:file_version"],
"allowed_source_types": [
"files:file_version",
"forms_runtime:form_submission_revision",
"cases:case_revision",
"decisions:decision_revision",
],
"retention_period_days": 3650,
"access_mode": "tenant",
"recorded_at": NOW,
@@ -727,6 +804,431 @@ class RecordsTests(unittest.TestCase):
self.assertTrue(recovery["healthy"])
self.assertTrue(recovery["package_checks"][0]["manifest_verified"])
def test_assisted_service_reference_journey_survives_backup_and_restore(
self,
) -> None:
fixture = json.loads(
(
Path(__file__).parent / "fixtures/service_to_decision_journey.json"
).read_text(encoding="utf-8")
)
expected = fixture["expected"]
context = fixture["institutional_context"]
self.records = SqlRecordRegistry(ReferenceJourneyRegistry())
def create_journey_record(
record_data: dict[str, object],
record_context: dict[str, object],
*,
idempotency_suffix: str,
channel: str,
) -> dict[str, object]:
return self.records.create_record(
self.session,
self.principal,
payload={
**record_data,
"file_plan_node_id": "plan-permits",
"description": fixture["description"],
"state": "open",
"source_authority_mode": "native_authoritative",
"access_mode": "tenant",
"purpose": record_context["purpose"],
"responsible_unit_id": "unit-mobility-services",
"responsible_function_id": record_context[
"responsible_function_id"
],
"institutional_context": {**record_context, "channel": channel},
"recorded_at": NOW + timedelta(minutes=1),
"valid_from": NOW,
"change_reason": f"The {channel} application was accepted for processing.",
"idempotency_key": f"reference-record-create-{idempotency_suffix}",
},
)
def filing_request(
record_id: str,
source: dict[str, object],
record_context: dict[str, object],
*,
idempotency_suffix: str = "",
) -> RecordFilingRequest:
metadata = dict(source["metadata"])
return RecordFilingRequest(
tenant_id="tenant-1",
record_id=record_id,
source=RecordSourceLocator(
tenant_id="tenant-1",
source_module=str(source["source_module"]),
resource_type=str(source["resource_type"]),
resource_id=str(source["resource_id"]),
source_revision=str(source["source_revision"]),
metadata={
"reference": {
"label": source["label"],
"authority_mode": source["authority_mode"],
"content_type": source["content_type"],
"size_bytes": source["size_bytes"],
"launch_url": source["launch_url"],
},
**metadata,
},
),
purpose=str(source["purpose"]),
filing_reason=str(source["filing_reason"]),
relationship=str(source["relationship"]),
institutional_context={
**record_context,
"evidence_role": metadata["evidence_role"],
},
metadata={"fixture_source_id": source["id"]},
idempotency_key=(
f"reference-file-{source['id']}{idempotency_suffix}"
),
)
record = create_journey_record(
fixture["record"], context, idempotency_suffix="assisted", channel="assisted"
)
self.assertEqual(fixture["record"]["title"], record["title"])
filed_results = []
for source in fixture["sources"]:
filed_results.append(
self.records.file(
self.session,
self.principal,
request=filing_request(str(record["record_id"]), source, context),
)
)
replay = self.records.file(
self.session,
self.principal,
request=filing_request(
str(record["record_id"]), fixture["sources"][0], context
),
)
self.assertTrue(replay.replayed)
self.assertEqual(filed_results[0].item_id, replay.item_id)
detail = self.records.get_record(
self.session, self.principal, record_id=str(record["record_id"])
)
self.assertEqual(expected["filed_item_count"], len(detail["items"]))
self.assertEqual(
[source["source_revision"] for source in fixture["sources"]],
[item["source"]["source_revision"] for item in detail["items"]],
)
self.assertEqual(
[source["relationship"] for source in fixture["sources"]],
[item["relationship"] for item in detail["items"]],
)
self.assertEqual(
{
"application",
"application_attachment",
"case_context",
"formal_decision",
"delivery_receipt",
"correction",
},
{item["source_metadata"]["evidence_role"] for item in detail["items"]},
)
digital_fixture = fixture["digital_equivalent"]
digital_context = {
**context,
**digital_fixture["institutional_context_overrides"],
}
digital_record = create_journey_record(
digital_fixture["record"],
digital_context,
idempotency_suffix="digital",
channel="digital",
)
digital_sources = [digital_fixture["intake_source"]]
for source in fixture["sources"][1:]:
digital_sources.append(
{
**source,
"id": f"digital-{source['id']}",
"resource_id": f"{source['resource_id']}-digital",
"label": f"Digital equivalent · {source['label']}",
}
)
for source in digital_sources:
self.records.file(
self.session,
self.principal,
request=filing_request(
str(digital_record["record_id"]),
source,
digital_context,
idempotency_suffix="-digital",
),
)
digital_detail = self.records.get_record(
self.session,
self.principal,
record_id=str(digital_record["record_id"]),
)
def equivalence_fields(item):
return (
item["source"]["source_module"],
item["source"]["resource_type"],
item["relationship"],
item["source_metadata"]["evidence_role"],
)
self.assertEqual(
[equivalence_fields(item) for item in detail["items"]],
[equivalence_fields(item) for item in digital_detail["items"]],
)
all_records, total_records = self.records.list_records(
self.session, self.principal
)
self.assertEqual(expected["equivalent_record_count"], total_records)
self.assertEqual(
{fixture["record"]["record_id"], digital_fixture["record"]["record_id"]},
{item["record_id"] for item in all_records},
)
lifecycle_start = datetime.now(UTC).replace(microsecond=0) + timedelta(
minutes=10
)
closed = self.records.close_record(
self.session,
self.principal,
record_id=str(record["record_id"]),
payload={
"expected_revision": 1,
"purpose": "close completed resident parking permit record",
"reason": "Decision, delivery, and correction evidence are complete.",
"recorded_at": lifecycle_start,
"idempotency_key": "reference-close",
},
)
self.assertEqual(expected["closed_state"], closed["state"])
appraised = self.records.appraise_record(
self.session,
self.principal,
record_id=str(record["record_id"]),
payload={
"expected_revision": 2,
"outcome": "transfer",
"purpose": "appraise completed resident parking permit record",
"reason": "Offer the complete record after governed review.",
"policy_refs": [context["retention_policy_ref"]],
"override_retention_not_due": True,
"recorded_at": lifecycle_start + timedelta(minutes=1),
"idempotency_key": "reference-appraise",
},
)
self.assertEqual(expected["appraised_state"], appraised["state"])
hold = self.records.apply_hold(
self.session,
self.principal,
record_id=str(record["record_id"]),
payload={
"expected_record_revision": 3,
"reason": "Preserve the record while a correction is reviewed.",
"authority": "Mobility authority review 2026-17",
"purpose": "preserve correction evidence",
"policy_refs": [context["retention_policy_ref"]],
"institutional_context": context,
"recorded_at": lifecycle_start + timedelta(minutes=2),
"idempotency_key": "reference-hold",
},
)
with self.assertRaisesRegex(RecordConflictError, "active record hold"):
self.records.propose_disposition(
self.session,
self.principal,
record_id=str(record["record_id"]),
payload={
"expected_record_revision": 3,
"action": "transfer",
"reason": "Attempt transfer while evidence is held.",
"purpose": "prove hold enforcement",
"recorded_at": lifecycle_start + timedelta(minutes=3),
"idempotency_key": "reference-disposition-blocked",
},
)
self.records.release_hold(
self.session,
self.principal,
record_id=str(record["record_id"]),
hold_id=str(hold["hold_id"]),
payload={
"expected_hold_revision": 1,
"reason": "Correction review completed with prior evidence preserved.",
"purpose": "resume governed disposition",
"recorded_at": lifecycle_start + timedelta(minutes=4),
"idempotency_key": "reference-hold-release",
},
)
disposition = self.records.propose_disposition(
self.session,
self.principal,
record_id=str(record["record_id"]),
payload={
"disposition_id": "reference-disposition",
"expected_record_revision": 3,
"action": "transfer",
"reason": "The complete record is ready for an independently reviewed offer.",
"purpose": "prepare governed archive offer",
"policy_refs": [context["retention_policy_ref"]],
"institutional_context": context,
"recorded_at": lifecycle_start + timedelta(minutes=5),
"idempotency_key": "reference-disposition",
},
)
self.records.registry.approvals.approved = True
finalized = self.records.finalize_disposition(
self.session,
self.principal,
record_id=str(record["record_id"]),
disposition_id=str(disposition["disposition_id"]),
payload={
"expected_disposition_revision": 1,
"purpose": "record independent disposition approval",
"recorded_at": lifecycle_start + timedelta(minutes=6),
"idempotency_key": "reference-disposition-finalize",
},
)
self.assertEqual(expected["approved_state"], finalized["record"]["state"])
package = self.records.prepare_transfer_package(
self.session,
self.principal,
record_id=str(record["record_id"]),
payload={
"package_id": "reference-transfer-package",
"disposition_id": str(disposition["disposition_id"]),
"expected_record_revision": 4,
"provider_id": "simulation",
"profile": "govoplan-simulation-v1",
"purpose": "validate archive transfer boundary",
"recorded_at": lifecycle_start + timedelta(minutes=7),
"idempotency_key": "reference-transfer-prepare",
},
)
receipt = self.records.dispatch_transfer_package(
self.session,
self.principal,
record_id=str(record["record_id"]),
package_id=str(package["package_id"]),
payload={
"expected_package_revision": 1,
"purpose": "validate archive receipt handling",
"recorded_at": lifecycle_start + timedelta(minutes=8),
"idempotency_key": "reference-transfer-dispatch",
},
)
self.assertEqual(expected["transfer_status"], receipt["status"])
self.assertEqual(
expected["custody_transferred"],
receipt["receipt"]["metadata"]["custody_transferred"],
)
self.session.commit()
with tempfile.TemporaryDirectory(prefix="govoplan-records-restore-") as temp_dir:
backup_path = Path(temp_dir) / "restored-records.sqlite3"
source_connection = self.engine.raw_connection()
try:
with sqlite3.connect(backup_path) as backup_connection:
source_connection.driver_connection.backup(backup_connection)
finally:
source_connection.close()
restored_engine = create_engine(f"sqlite+pysqlite:///{backup_path}")
restored_session = Session(restored_engine)
try:
restored_records = SqlRecordRegistry(ReferenceJourneyRegistry())
matches, total = restored_records.list_records(
restored_session,
self.principal,
query=expected["restored_search_query"],
)
self.assertEqual(1, total)
self.assertEqual(record["record_id"], matches[0]["record_id"])
restored = restored_records.get_record(
restored_session,
self.principal,
record_id=str(record["record_id"]),
)
self.assertEqual(expected["filed_item_count"], len(restored["items"]))
self.assertEqual(
{**context, "channel": "assisted"},
restored["record"]["institutional_context"],
)
self.assertEqual("released", restored["holds"][0]["status"])
self.assertEqual(
expected["disposition_status"],
restored["dispositions"][0]["status"],
)
self.assertEqual(
[context["retention_policy_ref"]],
restored["dispositions"][0]["policy_refs"],
)
self.assertEqual(
expected["transfer_status"],
restored["transfer_packages"][0]["status"],
)
self.assertTrue(restored["transfer_packages"][0]["simulated"])
self.assertEqual(
expected["filed_item_count"],
sum(
entry["event_type"] == "record.item_filed"
for entry in restored["chronology"]
),
)
recovery = restored_records.recovery_status(
restored_session,
self.principal,
record_id=str(record["record_id"]),
)
self.assertTrue(recovery["healthy"])
self.assertEqual(
expected["filed_item_count"], len(recovery["source_checks"])
)
self.assertTrue(
all(
check["status"] == "verified"
for check in recovery["source_checks"]
)
)
self.assertTrue(
all(
check["manifest_verified"]
for check in recovery["package_checks"]
)
)
token = bind_temporal_data_context(
TemporalDataContext(
validity_mode="at",
valid_at=lifecycle_start - timedelta(minutes=1),
recorded_at=lifecycle_start - timedelta(minutes=1),
)
)
try:
historical = restored_records.get_record(
restored_session,
self.principal,
record_id=str(record["record_id"]),
)
self.assertEqual("open", historical["record"]["state"])
self.assertEqual(
expected["filed_item_count"], len(historical["items"])
)
self.assertEqual([], historical["transfer_packages"])
finally:
reset_temporal_data_context(token)
finally:
restored_session.close()
restored_engine.dispose()
def test_disposition_waits_for_optional_approvals_and_can_be_withdrawn(
self,
) -> None: