Files
govoplan/tests/test_institutional_service_journey.py
T
zemion f8b06887d2
Dependency Audit / dependency-audit (push) Successful in 1m46s
Deployment Installer / deployment-installer (push) Successful in 9s
Security Audit / security-audit (push) Failing after 1s
feat: extend resident permit service journey
2026-08-19 12:33:34 +02:00

868 lines
34 KiB
Python

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