feat: orchestrate accountable Campaign work
Module Package Release / publish-packages (push) Successful in 13s
Module Package Release / publish-packages (push) Successful in 13s
This commit is contained in:
@@ -32,8 +32,14 @@ from govoplan_campaign.backend.schemas import (
|
||||
CampaignWorkAssignmentReassignRequest,
|
||||
CampaignWorkAssignmentTransitionRequest,
|
||||
)
|
||||
from govoplan_core.core.access import GroupRef, UserRef
|
||||
from govoplan_campaign.backend.work_orchestration import (
|
||||
SqlCampaignWorkOrchestrationProvider,
|
||||
)
|
||||
from govoplan_core.auth import ApiPrincipal
|
||||
from govoplan_core.core.access import GroupRef, PrincipalRef, UserRef
|
||||
from govoplan_core.core.campaigns import CampaignWorkHandoffRequest
|
||||
from govoplan_core.core.change_sequence import ChangeSequenceEntry
|
||||
from govoplan_core.core.events import EventBus, event_bus_context
|
||||
from govoplan_core.core.organizations import OrganizationFunctionRef
|
||||
from govoplan_core.core.tasks import WorkItem
|
||||
from govoplan_core.db.base import Base
|
||||
@@ -257,6 +263,31 @@ def _assignee() -> _Principal:
|
||||
)
|
||||
|
||||
|
||||
def _api_principal(user_id: str, account_id: str) -> ApiPrincipal:
|
||||
return ApiPrincipal(
|
||||
principal=PrincipalRef(
|
||||
account_id=account_id,
|
||||
membership_id=user_id,
|
||||
tenant_id=TENANT_ID,
|
||||
scopes=frozenset(
|
||||
{
|
||||
"campaigns:campaign:read",
|
||||
"campaigns:campaign:create",
|
||||
"campaigns:assignment:read",
|
||||
"campaigns:assignment:manage",
|
||||
"campaigns:assignment:complete",
|
||||
}
|
||||
),
|
||||
),
|
||||
account=SimpleNamespace(id=account_id),
|
||||
user=SimpleNamespace(
|
||||
id=user_id,
|
||||
display_name=f"User {user_id}",
|
||||
email=f"{user_id}@example.test",
|
||||
),
|
||||
)
|
||||
|
||||
|
||||
def _commit_audit(session: Session, *_args, **_kwargs) -> None:
|
||||
session.commit()
|
||||
|
||||
@@ -407,10 +438,189 @@ def test_reassignment_and_deactivation_reconciliation_preserve_history(session:
|
||||
assert reassigned.assignee_type == "account"
|
||||
assert result.changed == 1
|
||||
assert result.assignments[0].assignee_resolution_state == "unavailable"
|
||||
assert [item.event_kind for item in reversed(history.items)] == ["assigned", "reassigned", "assignee_unavailable"]
|
||||
assert [item.event_kind for item in reversed(history.items)] == [
|
||||
"assigned",
|
||||
"reassigned",
|
||||
"assignee_unavailable",
|
||||
]
|
||||
assert history.items[1].details["assignee_id"] == "group-1"
|
||||
|
||||
|
||||
def test_workflow_provider_is_idempotent_emits_typed_events_and_rechecks_access(
|
||||
session: Session,
|
||||
) -> None:
|
||||
directory = _Directory()
|
||||
registry = _Registry(_Tasks(), _Notifications())
|
||||
provider = SqlCampaignWorkOrchestrationProvider(registry=registry)
|
||||
manager = _api_principal("user-1", "account-1")
|
||||
assignee = _api_principal("user-2", "account-2")
|
||||
request = CampaignWorkHandoffRequest(
|
||||
tenant_id=TENANT_ID,
|
||||
campaign_id="campaign-1",
|
||||
expected_campaign_revision=1,
|
||||
idempotency_key="workflow-handoff-1",
|
||||
purpose="Review the Campaign evidence",
|
||||
assignee_kind="account",
|
||||
assignee_id="account-2",
|
||||
correlation_id="workflow-correlation-1",
|
||||
workflow_instance_id="workflow-instance-1",
|
||||
workflow_step_id="workflow-step-1",
|
||||
)
|
||||
bus = EventBus()
|
||||
events = []
|
||||
bus.subscribe("campaign.work.changed", events.append)
|
||||
with (
|
||||
patch(
|
||||
"govoplan_campaign.backend.routes.assignments._access_directory",
|
||||
return_value=directory,
|
||||
),
|
||||
patch(
|
||||
"govoplan_campaign.backend.route_support._access_directory",
|
||||
return_value=directory,
|
||||
),
|
||||
patch(
|
||||
"govoplan_campaign.backend.routes.assignments.get_registry",
|
||||
return_value=registry,
|
||||
),
|
||||
patch(
|
||||
"govoplan_campaign.backend.work_orchestration.audit_from_principal",
|
||||
return_value=SimpleNamespace(id="audit-workflow-1"),
|
||||
),
|
||||
patch(
|
||||
"govoplan_campaign.backend.routes.assignments.audit_from_principal",
|
||||
side_effect=_commit_audit,
|
||||
),
|
||||
event_bus_context(bus),
|
||||
):
|
||||
created = provider.prepare_handoff(session, manager, request=request)
|
||||
session.commit()
|
||||
replayed = provider.prepare_handoff(session, manager, request=request)
|
||||
accepted = transition_campaign_work_assignment(
|
||||
"campaign-1",
|
||||
created.assignment_id,
|
||||
CampaignWorkAssignmentTransitionRequest(
|
||||
expected_revision=1,
|
||||
action="accept",
|
||||
),
|
||||
session,
|
||||
assignee,
|
||||
)
|
||||
completed = transition_campaign_work_assignment(
|
||||
"campaign-1",
|
||||
created.assignment_id,
|
||||
CampaignWorkAssignmentTransitionRequest(
|
||||
expected_revision=2,
|
||||
action="complete",
|
||||
),
|
||||
session,
|
||||
assignee,
|
||||
)
|
||||
|
||||
assert created.replayed is False
|
||||
assert replayed.replayed is True
|
||||
assert created.assignment_id == replayed.assignment_id
|
||||
assert created.campaign_ref == "campaign:campaign-1:version:version-1:r1"
|
||||
assert created.assignment_ref.endswith(":r1")
|
||||
assert created.optional_capabilities == {"tasks": True, "notifications": True}
|
||||
assert session.query(CampaignWorkAssignment).filter(
|
||||
CampaignWorkAssignment.orchestration_idempotency_key
|
||||
== "workflow-handoff-1"
|
||||
).count() == 1
|
||||
assert accepted.status == "in_progress"
|
||||
assert completed.status == "completed"
|
||||
assert [event.payload["outcome"] for event in events] == [
|
||||
"assigned",
|
||||
"accepted",
|
||||
"completed",
|
||||
]
|
||||
assert [event.payload["assignment_revision"] for event in events] == [1, 2, 3]
|
||||
assert all(event.correlation_id == "workflow-correlation-1" for event in events)
|
||||
|
||||
with patch(
|
||||
"govoplan_campaign.backend.route_support._access_directory",
|
||||
return_value=directory,
|
||||
):
|
||||
allowed = provider.inspect_handoff(
|
||||
session,
|
||||
assignee,
|
||||
tenant_id=TENANT_ID,
|
||||
assignment_id=created.assignment_id,
|
||||
expected_revision=3,
|
||||
)
|
||||
share = session.get(CampaignShare, "share-1")
|
||||
assert share is not None
|
||||
share.revoked_at = completed.updated_at
|
||||
session.flush()
|
||||
revoked = provider.inspect_handoff(
|
||||
session,
|
||||
assignee,
|
||||
tenant_id=TENANT_ID,
|
||||
assignment_id=created.assignment_id,
|
||||
expected_revision=3,
|
||||
)
|
||||
|
||||
assert allowed.allowed is True
|
||||
assert revoked.allowed is False
|
||||
assert revoked.provenance["code"] == "campaign_handoff_access_revoked"
|
||||
|
||||
|
||||
def test_workflow_provider_can_create_self_assigned_campaign_and_rejects_stale_revision(
|
||||
session: Session,
|
||||
) -> None:
|
||||
directory = _Directory()
|
||||
registry = _Registry()
|
||||
provider = SqlCampaignWorkOrchestrationProvider(registry=registry)
|
||||
manager = _api_principal("user-1", "account-1")
|
||||
with (
|
||||
patch(
|
||||
"govoplan_campaign.backend.routes.assignments._access_directory",
|
||||
return_value=directory,
|
||||
),
|
||||
patch(
|
||||
"govoplan_campaign.backend.routes.assignments.get_registry",
|
||||
return_value=registry,
|
||||
),
|
||||
patch(
|
||||
"govoplan_campaign.backend.work_orchestration.audit_from_principal",
|
||||
return_value=SimpleNamespace(id="audit-workflow-create"),
|
||||
),
|
||||
):
|
||||
created = provider.prepare_handoff(
|
||||
session,
|
||||
manager,
|
||||
request=CampaignWorkHandoffRequest(
|
||||
tenant_id=TENANT_ID,
|
||||
create_external_id="workflow-created",
|
||||
create_name="Workflow-created Campaign",
|
||||
idempotency_key="workflow-create-1",
|
||||
purpose="Prepare the Campaign",
|
||||
assignee_kind="account",
|
||||
assignee_id="account-1",
|
||||
workflow_instance_id="workflow-instance-create",
|
||||
workflow_step_id="workflow-step-create",
|
||||
),
|
||||
)
|
||||
|
||||
assert session.get(Campaign, created.campaign_id).external_id == "workflow-created"
|
||||
version = session.get(CampaignVersion, "version-1")
|
||||
assert version is not None
|
||||
version.edit_revision = 2
|
||||
with pytest.raises(ValueError, match="Campaign revision changed"):
|
||||
provider.prepare_handoff(
|
||||
session,
|
||||
manager,
|
||||
request=CampaignWorkHandoffRequest(
|
||||
tenant_id=TENANT_ID,
|
||||
campaign_id="campaign-1",
|
||||
expected_campaign_revision=1,
|
||||
idempotency_key="workflow-stale-1",
|
||||
purpose="Review stale Campaign",
|
||||
assignee_kind="account",
|
||||
assignee_id="account-1",
|
||||
),
|
||||
)
|
||||
|
||||
|
||||
def test_organization_function_requires_authorized_current_incumbencies(session: Session) -> None:
|
||||
with (
|
||||
patch("govoplan_campaign.backend.routes.assignments._access_directory", return_value=_Directory()),
|
||||
@@ -470,3 +680,52 @@ def test_assignment_migration_is_repeatable_and_creates_history_indexes() -> Non
|
||||
migration.downgrade()
|
||||
assert not inspect(connection).has_table("campaign_work_assignment_events")
|
||||
assert not inspect(connection).has_table("campaign_work_assignments")
|
||||
|
||||
|
||||
def test_workflow_orchestration_migration_is_repeatable() -> None:
|
||||
assignments = importlib.import_module(
|
||||
"govoplan_campaign.backend.migrations.versions."
|
||||
"d8e9f0a1b2c3_v0121_campaign_work_assignments"
|
||||
)
|
||||
orchestration = importlib.import_module(
|
||||
"govoplan_campaign.backend.migrations.versions."
|
||||
"f3c7a9d2e6b1_v0123_campaign_work_orchestration"
|
||||
)
|
||||
engine = create_engine("sqlite+pysqlite:///:memory:")
|
||||
with engine.begin() as connection:
|
||||
connection.execute(text("CREATE TABLE access_users (id VARCHAR(36) PRIMARY KEY)"))
|
||||
connection.execute(text("CREATE TABLE campaigns (id VARCHAR(36) PRIMARY KEY)"))
|
||||
connection.execute(text("CREATE TABLE campaign_versions (id VARCHAR(36) PRIMARY KEY, campaign_id VARCHAR(36) NOT NULL)"))
|
||||
context = MigrationContext.configure(connection)
|
||||
with patch.object(assignments, "op", Operations(context)):
|
||||
assignments.upgrade()
|
||||
with patch.object(orchestration, "op", Operations(context)):
|
||||
orchestration.upgrade()
|
||||
orchestration.upgrade()
|
||||
|
||||
inspector = inspect(connection)
|
||||
columns = {
|
||||
item["name"]
|
||||
for item in inspector.get_columns("campaign_work_assignments")
|
||||
}
|
||||
indexes = {
|
||||
item["name"]
|
||||
for item in inspector.get_indexes("campaign_work_assignments")
|
||||
}
|
||||
assert {
|
||||
"orchestration_idempotency_key",
|
||||
"orchestration_request_sha256",
|
||||
"orchestration_correlation_id",
|
||||
"workflow_instance_id",
|
||||
"workflow_step_id",
|
||||
}.issubset(columns)
|
||||
assert "uq_campaign_work_assignment_orchestration_key" in indexes
|
||||
|
||||
with patch.object(orchestration, "op", Operations(context)):
|
||||
orchestration.downgrade()
|
||||
assert "workflow_instance_id" not in {
|
||||
item["name"]
|
||||
for item in inspect(connection).get_columns(
|
||||
"campaign_work_assignments"
|
||||
)
|
||||
}
|
||||
|
||||
@@ -484,8 +484,9 @@ def test_static_campaign_handbook_has_unique_ids_help_contexts_and_no_planned_re
|
||||
"campaign.work",
|
||||
"campaign.work.create",
|
||||
"campaign.work.action.start",
|
||||
"campaign.work.action.complete",
|
||||
"campaign.work.action.reassign",
|
||||
"campaign.work.action.complete",
|
||||
"campaign.work.action.reject",
|
||||
"campaign.work.action.reassign",
|
||||
"campaign.work.action.cancel",
|
||||
"campaign.work.history",
|
||||
}
|
||||
|
||||
@@ -0,0 +1,42 @@
|
||||
from __future__ import annotations
|
||||
|
||||
from govoplan_campaign.backend.manifest import get_manifest
|
||||
from govoplan_campaign.backend.workflow_definitions import (
|
||||
campaign_workflow_definitions,
|
||||
)
|
||||
from govoplan_core.core.campaigns import CAPABILITY_CAMPAIGNS_WORK_ORCHESTRATION
|
||||
|
||||
|
||||
def test_campaign_work_handoff_is_an_opt_in_reusable_template() -> None:
|
||||
contribution = campaign_workflow_definitions(module_version="0.1.23")[0]
|
||||
|
||||
assert contribution.origin_module_id == "campaigns"
|
||||
assert contribution.definition_kind == "template"
|
||||
assert contribution.activate_on_install is False
|
||||
assert contribution.allow_reuse is True
|
||||
assert contribution.required_capabilities == (
|
||||
CAPABILITY_CAMPAIGNS_WORK_ORCHESTRATION,
|
||||
)
|
||||
assert contribution.policy_metadata["assignment_authorization_neutral"] is True
|
||||
|
||||
nodes = {
|
||||
str(node["id"]): node
|
||||
for node in contribution.graph["nodes"] # type: ignore[index]
|
||||
}
|
||||
prepare = nodes["prepare"]
|
||||
handoff = nodes["campaign_work"]
|
||||
assert prepare["config"]["operation"] == "campaigns.work.prepare" # type: ignore[index]
|
||||
assert handoff["type"] == "workflow.external_handoff"
|
||||
assert handoff["config"]["event_type"] == "campaign.work.changed" # type: ignore[index]
|
||||
assert handoff["config"]["terminal_outcomes"] == { # type: ignore[index]
|
||||
"completed": "completed",
|
||||
"rejected": "rejected",
|
||||
"cancelled": "cancelled",
|
||||
}
|
||||
|
||||
|
||||
def test_manifest_contributes_the_current_campaign_work_template() -> None:
|
||||
manifest = get_manifest()
|
||||
|
||||
assert len(manifest.workflow_definitions) == 1
|
||||
assert manifest.workflow_definitions[0].origin_module_version == manifest.version
|
||||
Reference in New Issue
Block a user