from __future__ import annotations from dataclasses import dataclass from datetime import UTC, datetime, timedelta from types import SimpleNamespace 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, FormDefinition, FormFieldDefinition, InstitutionalReference, ServiceBinding, ServiceDefinition, TemporalRevision, service_launch_capability, ) 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 ( FormInstanceEvent, FormInstanceIdentity, FormInstanceRevision, ) from govoplan_forms_runtime.backend.service import ( FormRuntimeService, FormsServiceLauncher, ) from govoplan_portal.backend.service_directory import PortalServiceDirectory 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) def _service() -> ServiceDefinition: return ServiceDefinition( reference=InstitutionalReference( kind="service", owner_module="portal", object_id="permit", tenant_id="tenant-1", version="5", ), key="permit.apply", temporal=TemporalRevision( revision="5", valid_from=NOW - timedelta(days=1), valid_to=NOW + timedelta(days=1), recorded_at=NOW - timedelta(days=2), ), title="Apply for a permit", audience=("resident",), required_evidence_types=("application",), bindings=( ServiceBinding("capability", CAPABILITY_CASES_SERVICE_INTAKE), ServiceBinding("case", "permit-application"), ServiceBinding("workflow", "workflow: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 ) def has(self, module_id: str) -> bool: return module_id in {"portal", "forms", "forms_runtime"} 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("5", plan.context.service_ref.version) self.assertEqual("workflow:permit-review", plan.workflow_refs[0]) 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="permit-application", tenant_id="tenant-1", version="3", ), key="permit-application", 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 permit application.", ), title="Permit application", fields=( FormFieldDefinition( key="applicant_name", label="Applicant name", required=True, constraints={"min_length": 2}, ), ), 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="permit", tenant_id="tenant-1", version="6", ), key="permit.apply", temporal=TemporalRevision( revision="6", valid_from=NOW - timedelta(days=1), valid_to=NOW + timedelta(days=1), recorded_at=NOW - timedelta(days=2), ), title="Apply for a permit", 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={"applicant_name": "Ada Lovelace"}, ) replay = directory.launch_service( session, principal, reference=service.reference, requested_at=NOW, idempotency_key="portal-form-launch-1", parameters={"applicant_name": "Ada Lovelace"}, ) 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("Ada Lovelace", instance.values["applicant_name"]) finally: session.close() 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="Permit decision", graph=WorkflowGraph( nodes=[ WorkflowNode( id="start", type="workflow.start.manual", ), WorkflowNode( id="review", type="workflow.activity", config={ "title": "Decide the permit application", "instructions": "Review the filed evidence and record the decision.", "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( "Decide the permit application", 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() if __name__ == "__main__": unittest.main()