Files
govoplan/tests/test_institutional_service_journey.py
T
zemion b269791c48
Dependency Audit / dependency-audit (push) Successful in 1m48s
Deployment Installer / deployment-installer (push) Successful in 6s
Security Audit / security-audit (push) Successful in 12m2s
Reconcile product inputs and prove durable work handoffs
2026-08-06 16:06:18 +02:00

502 lines
18 KiB
Python

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()