from __future__ import annotations from dataclasses import dataclass from datetime import UTC, datetime, timedelta import json from pathlib import Path from types import SimpleNamespace from urllib.parse import parse_qs, urlparse import unittest from sqlalchemy import create_engine from sqlalchemy.orm import Session, sessionmaker from govoplan_core.auth import ApiPrincipal from govoplan_core.core.access import PrincipalRef from govoplan_core.core.institutional import ( CAPABILITY_FORM_DEFINITIONS, CAPABILITY_SERVICE_DEFINITIONS, EvidenceReference, FormDefinition, FormFieldDefinition, InstitutionalReference, ServiceBinding, ServiceDefinition, TemporalRevision, service_launch_capability, ) from govoplan_core.core.notifications import CAPABILITY_NOTIFICATIONS_DISPATCH from govoplan_core.core.payments import ( ManualPaymentReconciliationCommand, PaymentRequestCommand, ) from govoplan_core.core.runtime_coordination import ( DistributedLease, RuntimeIdentity, bind_process_runtime_identity, ) from govoplan_core.core.tasks import ( RegisteredWorkItemProvider, WorkItemProviderRegistration, WorkItemQuery, ) from govoplan_core.core.recovery import RecoveryCheckpoint, RecoveryOperation from govoplan_cases.backend.service_intake import ( CAPABILITY_CASES_SERVICE_INTAKE, CaseServiceIntake, ) from govoplan_forms.backend.db.models import FormDefinitionRevision from govoplan_forms.backend.service import ( SqlFormDefinitionProvider, record_form_definition, ) from govoplan_forms_runtime.backend.db.models import ( FormAssistedConfirmation, FormInstanceEvent, FormInstanceIdentity, FormInstanceRevision, FormIntakeProfile, FormIntakeSession, FormStatusAccessGrant, FormStatusAccessPolicy, FormStatusAccessToken, ) from govoplan_forms_runtime.backend.intake import FormIntakeService from govoplan_forms_runtime.backend.service import ( FormRuntimeError, FormRuntimeService, FormsServiceLauncher, ) from govoplan_forms_runtime.backend.status_access import FormStatusAccessService from govoplan_portal.backend.service_directory import PortalServiceDirectory from govoplan_payments.backend.db.models import ( PaymentEvent, PaymentObligation, PaymentReconciliation, ) from govoplan_payments.backend.service import SqlPaymentRequestProvider from govoplan_tasks.backend.aggregation import aggregate_work_items from govoplan_workflow_engine.backend.db.models import ( WorkflowDefinition, WorkflowDefinitionRevision, WorkflowInstance, WorkflowInstanceEvent, WorkflowInstanceStep, WorkflowTrigger, WorkflowTriggerDelivery, WorkflowWaitState, ) from govoplan_workflow_engine.backend.instance_service import ( resolve_step, start_instance, ) from govoplan_workflow_engine.backend.schemas import ( WorkflowDefinitionCreateRequest, WorkflowEdge, WorkflowGraph, WorkflowInstanceStartRequest, WorkflowNode, WorkflowStepActionRequest, ) from govoplan_workflow_engine.backend.service import ( activate_definition, create_definition, ) from govoplan_workflow_engine.backend.work_items import WorkflowWorkItemProvider NOW = datetime(2026, 8, 1, 10, 0, tzinfo=UTC) JOURNEY = json.loads( (Path(__file__).parent / "fixtures/resident_parking_permit_journey.json").read_text( encoding="utf-8" ) ) def _service() -> ServiceDefinition: return ServiceDefinition( reference=InstitutionalReference( kind="service", owner_module="portal", object_id=JOURNEY["service"]["object_id"], tenant_id="tenant-1", version=JOURNEY["service"]["version"], ), key=JOURNEY["service"]["key"], temporal=TemporalRevision( revision=JOURNEY["service"]["version"], valid_from=NOW - timedelta(days=1), valid_to=NOW + timedelta(days=1), recorded_at=NOW - timedelta(days=2), ), title=JOURNEY["title"], audience=(JOURNEY["service"]["audience"],), required_evidence_types=tuple(JOURNEY["service"]["required_evidence_types"]), bindings=( ServiceBinding("capability", CAPABILITY_CASES_SERVICE_INTAKE), ServiceBinding("case", JOURNEY["case"]["type_key"]), ServiceBinding("workflow", "workflow:resident-parking-permit-review"), ), publication_state="published", ) class _Provider: def __init__(self, definition: ServiceDefinition) -> None: self.definition = definition def get_service_definition( self, session, principal, *, reference, effective_at=None ): return self.definition def list_service_definitions( self, session, principal, *, tenant_id, query="", limit=100 ): return (self.definition,) class _Registry: def __init__(self, definition: ServiceDefinition) -> None: self.capabilities = { CAPABILITY_SERVICE_DEFINITIONS: _Provider(definition), CAPABILITY_CASES_SERVICE_INTAKE: CaseServiceIntake(), service_launch_capability("case"): object(), } def has(self, module_id: str) -> bool: return module_id in {"portal", "cases"} def has_capability(self, name: str) -> bool: return name in self.capabilities def capability(self, name: str) -> object: return self.capabilities[name] def require_capability(self, name: str) -> object: return self.capabilities[name] @dataclass class _Principal: tenant_id: str = "tenant-1" account_id: str = "account-1" class _FormRegistry(_Registry): def __init__(self, definition: ServiceDefinition) -> None: super().__init__(definition) self.capabilities[CAPABILITY_FORM_DEFINITIONS] = SqlFormDefinitionProvider() self.capabilities[service_launch_capability("form")] = FormsServiceLauncher( self ) self.notifications = _NotificationProvider() self.capabilities[CAPABILITY_NOTIFICATIONS_DISPATCH] = self.notifications def has(self, module_id: str) -> bool: return module_id in {"portal", "forms", "forms_runtime", "notifications"} class _NotificationProvider: def __init__(self) -> None: self.requests: list[object] = [] def tenant_id_for_notification(self, session, *, notification_id): return "tenant-1" def enqueue_notification(self, session, request, *, enqueue_delivery=True): self.requests.append(request) return {"id": f"notification-{len(self.requests)}"} def deliver_notification(self, session, *, notification_id): return {"id": notification_id} def deliver_pending(self, session, *, tenant_id=None, limit=50): return {"delivered": 0} class _WorkflowTaskRegistry: def __init__(self) -> None: self.provider = WorkflowWorkItemProvider(registry=self) self.registered = RegisteredWorkItemProvider( module_id="workflow_engine", registration=WorkItemProviderRegistration( id="workflow_engine.handoffs", factory=lambda _context: self.provider, order=20, ), ) def has_capability(self, _name: str) -> bool: return False def capability(self, name: str) -> object: raise KeyError(name) def work_item_providers(self): return ((self.registered, self.provider),) def _workflow_principal() -> ApiPrincipal: return ApiPrincipal( principal=PrincipalRef( account_id="account-1", membership_id="membership-1", tenant_id="tenant-1", scopes=frozenset( { "tasks:item:read", "workflow:definition:read", "workflow:instance:read", "workflow:instance:start", "workflow:instance:transition", } ), ), account=SimpleNamespace(id="account-1"), user=SimpleNamespace(id="membership-1"), ) class InstitutionalServiceJourneyTests(unittest.TestCase): def test_one_service_version_drives_portal_and_case_intake(self) -> None: definition = _service() registry = _Registry(definition) entries = PortalServiceDirectory(registry).list_entries( None, None, tenant_id="tenant-1", effective_at=NOW, audiences=("resident",), ) plan = registry.capability(CAPABILITY_CASES_SERVICE_INTAKE).plan( entries[0].definition, case_id="case-1", effective_at=NOW, ) self.assertTrue(entries[0].available) self.assertIs(definition, entries[0].definition) self.assertEqual(definition.reference, plan.service_ref) self.assertEqual(JOURNEY["service"]["version"], plan.context.service_ref.version) self.assertEqual("workflow:resident-parking-permit-review", plan.workflow_refs[0]) def test_reference_fixture_names_remaining_manual_target_evidence(self) -> None: self.assertEqual("Anwohnerparkausweis", JOURNEY["title_de"]) self.assertEqual("de-DE", JOURNEY["locale"]) self.assertEqual("email_link", JOURNEY["status_access"]["mode"]) self.assertEqual("manual", JOURNEY["payment"]["mode"]) self.assertEqual(8, len(JOURNEY["acceptance"]["automated"])) self.assertEqual(6, len(JOURNEY["acceptance"]["manual_or_target"])) def test_portal_launches_exact_form_revision_and_persists_submission(self) -> None: engine = create_engine("sqlite+pysqlite:///:memory:") for table in ( FormDefinitionRevision.__table__, FormInstanceIdentity.__table__, FormInstanceRevision.__table__, FormInstanceEvent.__table__, ): table.create(engine) session = Session(engine) principal = _Principal() try: form = record_form_definition( session, principal, definition=FormDefinition( reference=InstitutionalReference( kind="form", owner_module="forms", object_id=JOURNEY["form"]["object_id"], tenant_id="tenant-1", version="3", ), key=JOURNEY["form"]["object_id"], temporal=TemporalRevision( revision="3", valid_from=NOW - timedelta(days=1), valid_to=NOW + timedelta(days=1), recorded_at=NOW - timedelta(days=2), change_reason="Publish the resident parking permit application.", ), title=JOURNEY["title"], fields=( FormFieldDefinition( key="applicant_name", label="Applicant name", required=True, constraints={"min_length": 2}, ), FormFieldDefinition( key="applicant_email", label="Applicant email", value_type="email", required=True, constraints={"min_length": 5}, ), FormFieldDefinition( key="residence_address", label="Primary residence address", required=True, constraints={"min_length": 5}, ), FormFieldDefinition( key="licence_plate", label="Vehicle licence plate", required=True, constraints={"min_length": 3}, ), ), publication_state="published", allow_drafts=True, handoff_kinds=("case",), ), ) binding = ServiceBinding( "form", f"{form.reference.object_id}/{form.reference.version}", ) service = ServiceDefinition( reference=InstitutionalReference( kind="service", owner_module="services", object_id=JOURNEY["service"]["object_id"], tenant_id="tenant-1", version=JOURNEY["service"]["version"], ), key=JOURNEY["service"]["key"], temporal=TemporalRevision( revision=JOURNEY["service"]["version"], valid_from=NOW - timedelta(days=1), valid_to=NOW + timedelta(days=1), recorded_at=NOW - timedelta(days=2), ), title=JOURNEY["title"], audience=("public",), bindings=(binding,), publication_state="published", ) registry = _FormRegistry(service) directory = PortalServiceDirectory(registry) launched = directory.launch_service( session, principal, reference=service.reference, requested_at=NOW, idempotency_key="portal-form-launch-1", parameters=JOURNEY["form"]["fields"], ) replay = directory.launch_service( session, principal, reference=service.reference, requested_at=NOW, idempotency_key="portal-form-launch-1", parameters=JOURNEY["form"]["fields"], ) instance = FormRuntimeService(registry).get_instance( session, principal, instance_id=str(launched.metadata["form_instance_id"]), ) self.assertEqual("form_submission", launched.target_ref.kind) self.assertEqual(service.reference, launched.service_ref) self.assertEqual("3", launched.metadata["form_definition_revision"]) self.assertTrue(replay.replayed) self.assertEqual(launched.target_ref, replay.target_ref) self.assertEqual( ( form.reference.owner_module, form.reference.object_id, form.reference.tenant_id, form.reference.version, ), ( instance.definition_ref.owner_module, instance.definition_ref.object_id, instance.definition_ref.tenant_id, instance.definition_ref.version, ), ) self.assertEqual(NOW, instance.definition_ref.valid_at) self.assertEqual(service.reference, instance.service_ref) self.assertEqual(JOURNEY["form"]["fields"], instance.values) finally: session.close() engine.dispose() def test_assisted_intake_reuses_exact_form_and_persists_readback_provenance( self, ) -> None: engine = create_engine("sqlite+pysqlite:///:memory:") for table in ( FormDefinitionRevision.__table__, FormInstanceIdentity.__table__, FormInstanceRevision.__table__, FormInstanceEvent.__table__, FormIntakeProfile.__table__, FormIntakeSession.__table__, FormAssistedConfirmation.__table__, FormStatusAccessPolicy.__table__, FormStatusAccessGrant.__table__, FormStatusAccessToken.__table__, ): table.create(engine) sessions = sessionmaker(bind=engine) principal = _Principal() assisted = JOURNEY["assisted_intake"] try: with sessions() as session: form = record_form_definition( session, principal, definition=FormDefinition( reference=InstitutionalReference( kind="form", owner_module="forms", object_id=JOURNEY["form"]["object_id"], tenant_id="tenant-1", version=JOURNEY["form"]["version"], ), key=JOURNEY["form"]["object_id"], temporal=TemporalRevision( revision=JOURNEY["form"]["version"], recorded_at=NOW - timedelta(days=2), change_reason="Publish the resident parking permit application.", ), title=JOURNEY["title"], fields=tuple( FormFieldDefinition( key=key, label=key.replace("_", " ").title(), value_type=( "email" if key == "applicant_email" else "text" ), required=True, constraints={"min_length": 2}, ) for key in JOURNEY["form"]["fields"] ), publication_state="published", allow_drafts=True, handoff_kinds=("case",), ), ) registry = _FormRegistry(_service()) status_access = JOURNEY["status_access"] FormStatusAccessService(registry).upsert_policy( session, principal, definition_ref=form.reference, mode=status_access["mode"], enabled=True, email_field_key=status_access["email_field_key"], token_ttl_seconds=status_access["token_ttl_seconds"], request_limit_per_hour=status_access[ "request_limit_per_hour" ], recorded_at=NOW, ) intake = FormIntakeService(registry) profile = intake.create_profile( session, principal, definition_ref=form.reference, mode="assisted", custodian_ref=assisted["responsible_function_ref"], recorded_at=NOW, ) started = intake.start_assisted( session, principal, profile_id=profile.profile_id, values=JOURNEY["form"]["fields"], channel=assisted["channel"], affected_party_ref=assisted["affected_party_ref"], represented_party_ref=assisted["represented_party_ref"], authority_basis=assisted["authority_basis"], purpose=assisted["purpose"], legal_basis_ref=assisted["legal_basis_ref"], consent_basis=assisted["consent_basis"], notice_given=assisted["notice_given"], responsible_function_ref=assisted["responsible_function_ref"], language=assisted["language"], accessibility_needs=assisted["accessibility_needs"], field_sources={ key: { "source": "person_statement", "confidence": "stated", "declared_by_ref": assisted["affected_party_ref"], } for key in JOURNEY["form"]["fields"] }, idempotency_key="resident-permit-assisted-start", recorded_at=NOW + timedelta(minutes=1), ) self.assertEqual(form.reference, started.instance.definition_ref) self.assertEqual(JOURNEY["form"]["fields"], started.instance.values) self.assertEqual( assisted["purpose"], started.instance.metadata["intake"]["purpose"] ) instance_id = started.instance.instance_id session.commit() with sessions() as resumed: runtime = FormRuntimeService(registry) current = runtime.get_instance( resumed, principal, instance_id=instance_id, ) self.assertIsNotNone(current) with self.assertRaisesRegex(FormRuntimeError, "read-back confirmation"): runtime.submit_instance( resumed, principal, instance_id=instance_id, expected_revision=current.revision, values=current.values, attachment_refs=(), signature_refs=(), idempotency_key="resident-permit-assisted-unconfirmed", recorded_at=NOW + timedelta(minutes=2), ) confirmation = FormIntakeService(registry).record_assisted_confirmation( resumed, principal, instance_id=instance_id, expected_revision=current.revision, values=current.values, attachment_refs=(), signature_refs=(), outcome=assisted["confirmation_outcome"], method=assisted["confirmation_method"], confirmed_by_ref=assisted["affected_party_ref"], confirmed_at=NOW + timedelta(minutes=3), idempotency_key="resident-permit-assisted-readback", field_sources={ key: { "source": "person_statement", "confidence": "stated", "declared_by_ref": assisted["affected_party_ref"], } for key in JOURNEY["form"]["fields"] }, ) submitted = runtime.submit_instance( resumed, principal, instance_id=instance_id, expected_revision=current.revision, values=current.values, attachment_refs=(), signature_refs=(), idempotency_key="resident-permit-assisted-submit", recorded_at=NOW + timedelta(minutes=4), ) resumed.commit() self.assertEqual("submitted", submitted.status) self.assertEqual(current.revision, confirmation.instance_revision) self.assertEqual(assisted["affected_party_ref"], confirmation.confirmed_by_ref) status_service = FormStatusAccessService(registry) access = status_service.access_summary_for_instance( resumed, tenant_id="tenant-1", instance_id=instance_id, ) self.assertIsNotNone(access) tracking_id = str(access["tracking_id"]) challenge = status_service.public_access_challenge( resumed, tracking_id=tracking_id, ) self.assertEqual("email_link", challenge["mode"]) self.assertFalse( status_service.request_email_link( resumed, tracking_id=tracking_id, email="wrong@example.test", requested_at=NOW + timedelta(minutes=5), ) ) self.assertTrue( status_service.request_email_link( resumed, tracking_id=tracking_id, email=JOURNEY["form"]["fields"]["applicant_email"], requested_at=NOW + timedelta(minutes=6), ) ) notification = registry.notifications.requests[-1] query = parse_qs(urlparse(notification.action_url).query) projection = status_service.get_public_projection( resumed, tracking_id=tracking_id, token=query["token"][0], observed_at=NOW + timedelta(minutes=7), ) self.assertEqual("submitted", projection["status"]) self.assertEqual(JOURNEY["title"], projection["title"]) self.assertEqual( ["submitted"], [item["status"] for item in projection["timeline"]], ) self.assertNotIn("values", projection) finally: engine.dispose() def test_workflow_handoff_survives_session_reopen_and_projects_into_tasks( self, ) -> None: engine = create_engine("sqlite+pysqlite:///:memory:") tables = ( DistributedLease.__table__, RecoveryOperation.__table__, RecoveryCheckpoint.__table__, WorkflowDefinition.__table__, WorkflowDefinitionRevision.__table__, WorkflowInstance.__table__, WorkflowInstanceStep.__table__, WorkflowInstanceEvent.__table__, WorkflowTrigger.__table__, WorkflowTriggerDelivery.__table__, WorkflowWaitState.__table__, ) for table in tables: table.create(engine) sessions = sessionmaker(bind=engine) registry = _WorkflowTaskRegistry() principal = _workflow_principal() bind_process_runtime_identity( RuntimeIdentity( installation_id="service-journey", node_id="journey-node", incarnation="journey-run", role="web", software_version="test", composition_hash="c" * 64, ) ) try: with sessions() as session: definition = create_definition( session, tenant_id="tenant-1", actor_id="account-1", payload=WorkflowDefinitionCreateRequest( name=JOURNEY["workflow"]["definition_name"], graph=WorkflowGraph( nodes=[ WorkflowNode( id="start", type="workflow.start.manual", ), WorkflowNode( id="review", type="workflow.activity", config={ "title": JOURNEY["workflow"]["work_item_title"], "instructions": JOURNEY["workflow"]["instructions"], "assignee": "account:account-1", "due_after": "2d", }, ), WorkflowNode( id="done", type="workflow.end.completed", ), ], edges=[ WorkflowEdge( id="start-review", source="start", target="review", ), WorkflowEdge( id="review-done", source="review", target="done", ), ], ), execution_mode="guided", ), ) activate_definition( session, tenant_id="tenant-1", definition_id=definition.id, actor_id="account-1", ) instance, replayed = start_instance( session, tenant_id="tenant-1", definition_id=definition.id, actor_id="account-1", principal=principal, registry=registry, payload=WorkflowInstanceStartRequest( idempotency_key="permit-decision-1", input={"case_id": "case-1"}, correlation_id="case-1", ), ) self.assertFalse(replayed) session.commit() instance_id = instance.id step_id = instance.current_step_id with sessions() as reopened: work = aggregate_work_items( registry, reopened, principal, query=WorkItemQuery(tenant_id="tenant-1"), ) self.assertEqual(1, work.total) self.assertEqual(step_id, work.items[0].id) self.assertEqual( JOURNEY["workflow"]["work_item_title"], work.items[0].title, ) self.assertTrue(work.items[0].action_url.startswith("/workflow?")) self.assertIn(f"run={instance_id}", work.items[0].action_url) resolve_step( reopened, tenant_id="tenant-1", instance_id=instance_id, step_id=step_id, actor_id="account-1", principal=principal, registry=registry, payload=WorkflowStepActionRequest(action="complete"), ) reopened.commit() with sessions() as verified: self.assertEqual( 0, aggregate_work_items( registry, verified, principal, query=WorkItemQuery(tenant_id="tenant-1"), ).total, ) finally: bind_process_runtime_identity(None) engine.dispose() def test_case_bound_payment_handoff_is_replay_safe_and_evidence_bound( self, ) -> None: engine = create_engine("sqlite+pysqlite:///:memory:") for table in ( PaymentObligation.__table__, PaymentReconciliation.__table__, PaymentEvent.__table__, ): table.create(engine) session = Session(engine) provider = SqlPaymentRequestProvider() payment = JOURNEY["payment"] try: command = PaymentRequestCommand( tenant_id="tenant-1", source_module="cases", source_resource_type="case", source_resource_id="case-1", amount_minor=payment["amount_minor"], currency=payment["currency"], subject=payment["subject"], idempotency_key="resident-permit-case-1-fee", requested_at=NOW + timedelta(days=1), requested_by_ref="workflow:resident-parking-permit-review", due_at=NOW + timedelta(days=1 + payment["due_days"]), context_refs={ "case": "case-1", "workflow": "workflow:resident-parking-permit-review", }, ) requested = provider.request_payment(session, command) replay = provider.request_payment(session, command) self.assertEqual(requested["payment_id"], replay["payment_id"]) self.assertTrue(replay["replayed"]) self.assertEqual("case-1", requested["source"]["resource_id"]) paid = provider.reconcile_manual_payment( session, ManualPaymentReconciliationCommand( tenant_id="tenant-1", payment_id=str(requested["payment_id"]), amount_minor=payment["amount_minor"], currency=payment["currency"], transaction_reference="BANK-RPP-2026-0001", evidence_ref=EvidenceReference( kind="document", owner_module=payment["evidence_owner"], evidence_id="file-payment-rpp-1", tenant_id="tenant-1", version="1", checksum="b" * 64, ), idempotency_key="resident-permit-bank-receipt-1", received_at=NOW + timedelta(days=2), recorded_at=NOW + timedelta(days=2, minutes=5), recorded_by_ref="account:payment-officer-1", ), ) session.commit() self.assertEqual("paid", paid["status"]) self.assertEqual( "BANK-RPP-2026-0001", paid["reconciliation"]["transaction_reference"], ) self.assertEqual(2, len(paid["events"])) finally: session.close() engine.dispose() if __name__ == "__main__": unittest.main()