Reconcile product inputs and prove durable work handoffs
This commit is contained in:
@@ -2,11 +2,14 @@ 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
|
||||
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,
|
||||
@@ -18,6 +21,17 @@ from govoplan_core.core.institutional import (
|
||||
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,
|
||||
@@ -37,6 +51,34 @@ from govoplan_forms_runtime.backend.service import (
|
||||
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)
|
||||
@@ -74,10 +116,14 @@ class _Provider:
|
||||
def __init__(self, definition: ServiceDefinition) -> None:
|
||||
self.definition = definition
|
||||
|
||||
def get_service_definition(self, session, principal, *, reference, effective_at=None):
|
||||
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):
|
||||
def list_service_definitions(
|
||||
self, session, principal, *, tenant_id, query="", limit=100
|
||||
):
|
||||
return (self.definition,)
|
||||
|
||||
|
||||
@@ -111,17 +157,58 @@ class _Principal:
|
||||
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.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()
|
||||
@@ -266,6 +353,149 @@ class InstitutionalServiceJourneyTests(unittest.TestCase):
|
||||
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()
|
||||
|
||||
Reference in New Issue
Block a user