Compare commits

...
2 Commits
Author SHA1 Message Date
zemion 1bd24f9b5b feat: orchestrate accountable Campaign work
Module Package Release / publish-packages (push) Successful in 13s
2026-08-22 02:14:35 +02:00
zemion 4f52f010ee fix(campaigns): declare work view surface
Module Package Release / publish-packages (push) Successful in 12s
2026-08-22 00:56:37 +02:00
17 changed files with 1613 additions and 26 deletions
+29
View File
@@ -138,6 +138,35 @@ Notifications. The thread displays only human discussion; approvals, workflow
state, delivery events, and durable system evidence remain on their owning
surfaces and in Tenant audit.
### Assign accountable campaign work
Open **Work** to assign one bounded purpose to an account, group, or
organization function that already has Campaign access. Assignment records
responsibility only: it never creates a share, transfers ownership, or grants a
permission. Assignees may accept, complete, or reject their work; rejection is
distinct from administrative cancellation. Managers may reassign or cancel
open work, and every transition retains the expected revision, actor snapshot,
typed target, and append-only event history.
Workflow may create or reference a Campaign and open the same assignment through
the optional `campaigns.workOrchestration` capability. Those assignments pin the
Campaign version and store the Workflow instance, step, correlation, and
idempotency provenance. Campaign emits `campaign.work.changed` for assignment,
acceptance, start, reassignment, completion, rejection, and cancellation.
Workflow uses the assignment ID and event revision, rechecks current Campaign
access, and then resumes the matching durable external hand-off without browser
polling. A missing Tasks or Notifications capability only removes the optional
projection or notification. A missing Campaign provider, revoked Campaign
access, or stale event revision keeps the Workflow blocked and inspectable.
Campaign also contributes the opt-in **Accountable Campaign work hand-off**
Workflow template. It is deliberately not activated on installation. A
configurator must copy or activate it and supply either `campaign_id` or
`create_campaign`; unused optional input keys must be present with `null`
values. The template prepares the assignment idempotently, opens the exact
Campaign work URL, and waits for completion, rejection, cancellation, or the
configured timeout. Opening the link never completes the Workflow.
### Prepare a campaign
1. Create a campaign and confirm its owner or owning group.
+2 -2
View File
@@ -4,14 +4,14 @@ build-backend = "setuptools.build_meta"
[project]
name = "govoplan-campaign"
version = "0.1.21"
version = "0.1.23"
description = "GovOPlaN campaigns module with backend and WebUI integration."
readme = "README.md"
requires-python = ">=3.12"
license = { file = "LICENSE" }
authors = [{ name = "GovOPlaN" }]
dependencies = [
"govoplan-core>=0.1.18",
"govoplan-core>=0.1.28",
"jsonschema>=4,<5",
"pydantic>=2,<3",
"SQLAlchemy>=2,<3",
@@ -226,6 +226,11 @@ class CampaignCollaborationEntry(Base, TimestampMixin):
class CampaignWorkAssignment(Base, TimestampMixin):
__tablename__ = "campaign_work_assignments"
__table_args__ = (
UniqueConstraint(
"tenant_id",
"orchestration_idempotency_key",
name="uq_campaign_work_assignment_orchestration_key",
),
Index(
"ix_campaign_work_assignments_campaign_status",
"tenant_id",
@@ -275,6 +280,21 @@ class CampaignWorkAssignment(Base, TimestampMixin):
)
task_mirror_error: Mapped[str | None] = mapped_column(String(500), nullable=True)
task_mirrored_at: Mapped[datetime | None] = mapped_column(DateTime(timezone=True), nullable=True)
orchestration_idempotency_key: Mapped[str | None] = mapped_column(
String(255), nullable=True, index=True
)
orchestration_request_sha256: Mapped[str | None] = mapped_column(
String(64), nullable=True
)
orchestration_correlation_id: Mapped[str | None] = mapped_column(
String(128), nullable=True, index=True
)
workflow_instance_id: Mapped[str | None] = mapped_column(
String(36), nullable=True, index=True
)
workflow_step_id: Mapped[str | None] = mapped_column(
String(36), nullable=True, index=True
)
resource_revision: Mapped[int] = mapped_column(Integer, default=1, nullable=False)
@@ -252,7 +252,7 @@ CAMPAIGN_USER_DOCUMENTATION = (
topic_id="campaigns.workflow.assign-accountable-work",
title="Assign accountable Campaign work without granting access",
summary="Record bounded work for an account, group, or organization function while keeping authorization and Campaign ownership separate.",
body="Campaign work assignments record responsibility, not authority. Every reader and actor must still pass the parent Campaign access check, and a new account, group, or organization-function target is accepted only when it already resolves to active principals with Campaign access. Each assignment retains its purpose, optional due date, assigner, typed assignee reference, human-readable snapshot, current resolution state, stable Campaign or child reference, optimistic revision, and append-only transition history. Assignees with the separate completion permission can start and complete their own work; managers can reassign or cancel it. Reconciliation records vacancy, deactivation, or restored resolution without deleting history or transferring ownership. Notifications and Tasks mirroring are optional and cannot make the Campaign transaction fail.",
body="Campaign work assignments record responsibility, not authority. Every reader and actor must still pass the parent Campaign access check, and a new account, group, or organization-function target is accepted only when it already resolves to active principals with Campaign access. Each assignment retains its purpose, optional due date, assigner, typed assignee reference, human-readable snapshot, current resolution state, stable Campaign or child reference, optimistic revision, and append-only transition history. Assignees with the separate completion permission can accept, complete, or reject their own work; rejection is distinct from manager cancellation. Managers can also reassign or cancel it. Workflow-opened work additionally retains the correlation, idempotency, Workflow instance and step, exact Campaign version, and emits a common revision-bearing lifecycle event for assignment, acceptance, start, reassignment, completion, rejection, or cancellation. Workflow rechecks Campaign access before it resumes; the assignment itself never grants access. Reconciliation records vacancy, deactivation, or restored resolution without deleting history or transferring ownership. Notifications and Tasks mirroring are optional and cannot make the Campaign transaction fail.",
order=34,
audience=("campaign_manager", "campaign_reviewer", "campaign_sender"),
required_scopes=("campaigns:campaign:read", "campaigns:assignment:read"),
@@ -263,6 +263,7 @@ CAMPAIGN_USER_DOCUMENTATION = (
"campaign.work.create",
"campaign.work.action.start",
"campaign.work.action.complete",
"campaign.work.action.reject",
"campaign.work.action.reassign",
"campaign.work.action.cancel",
"campaign.work.history",
@@ -276,11 +277,11 @@ CAMPAIGN_USER_DOCUMENTATION = (
"Open Work in the selected Campaign workspace and choose Add assignment.",
"Enter a bounded purpose, optional due date, typed target, and optional stable Campaign evidence reference.",
"Resolve any authorization-neutral rejection by granting access through the separate Campaign sharing workflow or choosing another assignee; creating the assignment itself never grants access.",
"Start and complete your own assignment, or use manager actions to reassign or cancel open work.",
"Accept and complete your own assignment, reject it explicitly when it cannot be taken on, or use manager actions to reassign or cancel open work.",
"Reload and reconcile assignments after account, group, organization-function, or incumbency changes; inspect the retained history before acting on unavailable work.",
),
outcome="A durable accountability record whose lifecycle is independent from Campaign ownership, authorization, and delivery state.",
verification="Reload Work, inspect the typed target, resolution provenance, revision and history, and confirm Campaign shares and ownership did not change. When Tasks is installed, confirm the optional mirror links back to this Campaign assignment.",
verification="Reload Work, inspect the typed target, resolution provenance, revision and history, and confirm Campaign shares and ownership did not change. For Workflow-opened work, follow the focused assignment link and verify the exact terminal event and revision resume only the pinned Workflow step. When Tasks is installed, confirm the optional mirror links back to this Campaign assignment.",
related_modules=("access", "organizations", "idm", "tasks", "notifications", "audit", "policy"),
limitations=(
"Organizations and IDM are optional; organization-function assignment is unavailable until both directory and incumbency capabilities are active.",
@@ -291,7 +292,7 @@ CAMPAIGN_USER_DOCUMENTATION = (
"de": {
"title": "Verantwortliche Kampagnenarbeit zuweisen, ohne Zugriff zu vergeben",
"summary": "Begrenzte Arbeit für Konto, Gruppe oder Organisationsfunktion erfassen und Berechtigung sowie Kampagneneigentum getrennt halten.",
"body": "Kampagnenzuweisungen dokumentieren Verantwortung, nicht Berechtigung. Lesende und Handelnde müssen weiterhin den Zugriff auf die übergeordnete Kampagne nachweisen. Neue Ziele werden nur angenommen, wenn Konto, Gruppe oder alle aktuellen Funktionsinhabenden bereits Kampagnenzugriff besitzen. Zweck, optionale Fälligkeit, zuweisende Person, typisierte Referenz, lesbarer Schnappschuss, aktueller Auflösungszustand, Revision und unveränderliche Übergangshistorie bleiben erhalten. Deaktivierung oder Vakanz wird beim Abgleich als nicht verfügbar dokumentiert. Optionale Benachrichtigungen und Tasks-Spiegelungen dürfen die Kampagnentransaktion nicht blockieren.",
"body": "Kampagnenzuweisungen dokumentieren Verantwortung, nicht Berechtigung. Lesende und Handelnde müssen weiterhin den Zugriff auf die übergeordnete Kampagne nachweisen. Neue Ziele werden nur angenommen, wenn Konto, Gruppe oder alle aktuellen Funktionsinhabenden bereits Kampagnenzugriff besitzen. Zweck, optionale Fälligkeit, zuweisende Person, typisierte Referenz, lesbarer Schnappschuss, aktueller Auflösungszustand, Revision und unveränderliche Übergangshistorie bleiben erhalten. Zugewiesene Personen können Arbeit annehmen, abschließen oder ausdrücklich ablehnen; Ablehnung bleibt von einer administrativen Stornierung getrennt. Durch Workflow eröffnete Arbeit bewahrt Korrelation, Idempotenz, Workflow-Instanz und -Schritt sowie die genaue Kampagnenversion und erzeugt revisionsgebundene Lebenszyklusereignisse. Workflow prüft den Kampagnenzugriff vor der Fortsetzung erneut. Deaktivierung oder Vakanz wird beim Abgleich als nicht verfügbar dokumentiert. Optionale Benachrichtigungen und Tasks-Spiegelungen dürfen die Kampagnentransaktion nicht blockieren.",
}
},
),
+34 -3
View File
@@ -14,6 +14,7 @@ from govoplan_core.core.campaigns import (
CAPABILITY_CAMPAIGNS_POLICY_CONTEXT,
CAPABILITY_CAMPAIGNS_RETENTION,
CAPABILITY_CAMPAIGNS_SCHEDULES,
CAPABILITY_CAMPAIGNS_WORK_ORCHESTRATION,
)
from govoplan_core.core.calendar import CAPABILITY_CALENDAR_INVITATIONS
from govoplan_core.core.module_guards import (
@@ -72,6 +73,9 @@ from govoplan_campaign.backend.documentation import (
)
from govoplan_campaign.backend.dsar_provider import CAMPAIGN_DSAR_CAPABILITY
from govoplan_campaign.backend.search_source import create_campaign_search_source
from govoplan_campaign.backend.workflow_definitions import (
campaign_workflow_definitions,
)
register_campaign_change_tracking()
@@ -133,13 +137,13 @@ PERMISSIONS = (
_permission(
"campaigns:assignment:manage",
"Manage campaign work assignments",
"Create, reassign, cancel, and reconcile authorization-neutral campaign work assignments.",
"Create, reassign, cancel, and reconcile authorization-neutral campaign work assignments, including Workflow-opened hand-offs.",
"Campaign work",
),
_permission(
"campaigns:assignment:complete",
"Complete assigned campaign work",
"Start or complete campaign work assigned to the current account, group, or organization function.",
"Accept, complete, or reject campaign work assigned to the current account, group, or organization function.",
"Campaign work",
),
_permission(
@@ -439,7 +443,8 @@ def _campaigns_router(context: ModuleContext):
manifest = ModuleManifest(
id="campaigns",
name="Campaigns",
version="0.1.21",
version="0.1.23",
workflow_definitions=campaign_workflow_definitions(module_version="0.1.23"),
required_capabilities=(
CAPABILITY_AUTH_PRINCIPAL_RESOLVER,
CAPABILITY_AUTH_PERMISSION_EVALUATOR,
@@ -475,6 +480,10 @@ manifest = ModuleManifest(
ModuleInterfaceProvider(name="campaigns.mail_policy_context", version="0.1.6"),
ModuleInterfaceProvider(name="campaigns.policy_context", version="0.1.6"),
ModuleInterfaceProvider(name="campaigns.retention", version="0.1.6"),
ModuleInterfaceProvider(
name="campaigns.work_orchestration",
version="1.0.0",
),
ModuleInterfaceProvider(
name=REPORT_PROVIDER_CAPABILITY_PREFIX + "campaigns",
version="1.0.0",
@@ -691,12 +700,20 @@ manifest = ModuleManifest(
"campaigns.route.operator-redirect",
OPERATOR_QUEUE_SURFACE_ID,
REPORTS_SURFACE_ID,
"campaigns.page.work",
"campaigns.page.activity",
),
order=40,
),
),
view_surfaces=(
ViewSurface(
id="campaigns.page.work",
module_id="campaigns",
kind="page",
label="Campaign work",
order=44,
),
ViewSurface(
id="campaigns.page.activity",
module_id="campaigns",
@@ -1598,6 +1615,10 @@ manifest = ModuleManifest(
"govoplan_campaign.backend.capabilities",
fromlist=["retention_capability"],
).retention_capability(context),
CAPABILITY_CAMPAIGNS_WORK_ORCHESTRATION: lambda context: __import__(
"govoplan_campaign.backend.work_orchestration",
fromlist=["SqlCampaignWorkOrchestrationProvider"],
).SqlCampaignWorkOrchestrationProvider(registry=context.registry),
REPORT_PROVIDER_CAPABILITY_PREFIX + "campaigns": lambda context: __import__(
"govoplan_campaign.backend.reports.provider",
fromlist=["CampaignAggregateReportProvider"],
@@ -1605,6 +1626,16 @@ manifest = ModuleManifest(
CAMPAIGN_DSAR_CAPABILITY: _dsar_provider,
},
capability_documentation={
CAPABILITY_CAMPAIGNS_WORK_ORCHESTRATION: CapabilityDocumentation(
label="Campaign work orchestration",
summary=(
"Creates or references Campaign work idempotently and exposes "
"revision-bearing lifecycle events without granting access."
),
contract_version="1.0",
documentation_types=("admin", "user"),
audience=("campaign_manager", "workflow_designer", "module_admin"),
),
REPORT_PROVIDER_CAPABILITY_PREFIX + "campaigns": CapabilityDocumentation(
label="Campaign aggregate report provider",
summary=(
@@ -0,0 +1,103 @@
"""add durable Campaign work orchestration provenance
revision = "f3c7a9d2e6b1"
down_revision = "d8e9f0a1b2c3"
"""
from __future__ import annotations
import sqlalchemy as sa
from alembic import op
revision = "f3c7a9d2e6b1"
down_revision = "d8e9f0a1b2c3"
branch_labels = None
depends_on = None
_COLUMN_SPECS = (
("orchestration_idempotency_key", sa.String(length=255)),
("orchestration_request_sha256", sa.String(length=64)),
("orchestration_correlation_id", sa.String(length=128)),
("workflow_instance_id", sa.String(length=36)),
("workflow_step_id", sa.String(length=36)),
)
def upgrade() -> None:
inspector = sa.inspect(op.get_bind())
if not inspector.has_table("campaign_work_assignments"):
return
existing = {
item["name"]
for item in inspector.get_columns("campaign_work_assignments")
}
with op.batch_alter_table("campaign_work_assignments") as batch:
for name, column_type in _COLUMN_SPECS:
if name not in existing:
batch.add_column(sa.Column(name, column_type, nullable=True))
inspector = sa.inspect(op.get_bind())
indexes = {
item["name"]
for item in inspector.get_indexes("campaign_work_assignments")
}
for name, columns in (
(
"ix_campaign_work_assignments_orchestration_idempotency_key",
["orchestration_idempotency_key"],
),
(
"ix_campaign_work_assignments_orchestration_correlation_id",
["orchestration_correlation_id"],
),
(
"ix_campaign_work_assignments_workflow_instance_id",
["workflow_instance_id"],
),
(
"ix_campaign_work_assignments_workflow_step_id",
["workflow_step_id"],
),
(
"uq_campaign_work_assignment_orchestration_key",
["tenant_id", "orchestration_idempotency_key"],
),
):
if name not in indexes:
op.create_index(
name,
"campaign_work_assignments",
columns,
unique=name.startswith("uq_"),
)
def downgrade() -> None:
inspector = sa.inspect(op.get_bind())
if not inspector.has_table("campaign_work_assignments"):
return
indexes = {
item["name"]
for item in inspector.get_indexes("campaign_work_assignments")
}
for name in (
"uq_campaign_work_assignment_orchestration_key",
"ix_campaign_work_assignments_workflow_step_id",
"ix_campaign_work_assignments_workflow_instance_id",
"ix_campaign_work_assignments_orchestration_correlation_id",
"ix_campaign_work_assignments_orchestration_idempotency_key",
):
if name in indexes:
op.drop_index(name, table_name="campaign_work_assignments")
existing = {
item["name"]
for item in sa.inspect(op.get_bind()).get_columns(
"campaign_work_assignments"
)
}
with op.batch_alter_table("campaign_work_assignments") as batch:
for name, _column_type in reversed(_COLUMN_SPECS):
if name in existing:
batch.drop_column(name)
@@ -43,6 +43,13 @@ from govoplan_core.core.idm import (
CAPABILITY_IDM_FUNCTION_ASSIGNMENTS,
IdmFunctionAssignmentDirectory,
)
from govoplan_core.core.events import (
EventActorRef,
EventObjectRef,
EventTenantRef,
PlatformEvent,
emit_platform_event,
)
from govoplan_core.core.notifications import (
NotificationDispatchRequest,
notification_dispatch_provider,
@@ -127,7 +134,10 @@ def list_campaign_work_assignments(
campaign = _get_campaign_for_principal(session, campaign_id, principal)
_require_permission(principal, "campaigns:campaign:read")
statuses = tuple(dict.fromkeys(item.strip() for item in assignment_status if item.strip()))
if any(item not in {"open", "in_progress", "completed", "cancelled"} for item in statuses):
if any(
item not in {"open", "in_progress", "completed", "rejected", "cancelled"}
for item in statuses
):
raise HTTPException(status_code=422, detail="Unsupported assignment status filter.")
query = session.query(CampaignWorkAssignment).filter(
CampaignWorkAssignment.tenant_id == principal.tenant_id,
@@ -288,12 +298,20 @@ def transition_campaign_work_assignment(
if payload.action == "cancel":
raise HTTPException(status_code=403, detail="Only an assignment manager may cancel work.")
_require_revision(assignment, payload.expected_revision)
target_status = {"start": "in_progress", "complete": "completed", "cancel": "cancelled"}[payload.action]
target_status = {
"accept": "in_progress",
"start": "in_progress",
"complete": "completed",
"reject": "rejected",
"cancel": "cancelled",
}[payload.action]
if assignment.status == target_status:
return _assignment_response(assignment)
allowed = {
"accept": {"open"},
"start": {"open"},
"complete": {"open", "in_progress"},
"reject": {"open", "in_progress"},
"cancel": {"open", "in_progress"},
}
if assignment.status not in allowed[payload.action]:
@@ -306,7 +324,13 @@ def transition_campaign_work_assignment(
session,
assignment=assignment,
principal=principal,
event_kind={"start": "started", "complete": "completed", "cancel": "cancelled"}[payload.action],
event_kind={
"accept": "accepted",
"start": "started",
"complete": "completed",
"reject": "rejected",
"cancel": "cancelled",
}[payload.action],
details={"reason": payload.reason},
)
_notify_assignment(session, campaign=campaign, assignment=assignment, event_kind=payload.action)
@@ -697,6 +721,51 @@ def _record_event(
)
session.add(event)
session.flush()
emit_platform_event(
session,
PlatformEvent(
type="campaign.work.changed",
module_id="campaigns",
event_id=event.id,
occurred_at=event.created_at,
correlation_id=assignment.orchestration_correlation_id,
causation_id=assignment.workflow_step_id,
actor=EventActorRef(
type="account",
id=principal.account_id,
label=_actor_label(principal),
),
tenant=EventTenantRef(id=assignment.tenant_id),
subject=EventObjectRef(
type="campaign_work_assignment",
id=assignment.id,
),
resource=EventObjectRef(
type="campaign",
id=assignment.campaign_id,
),
classification="internal",
payload={
"campaign_id": assignment.campaign_id,
"campaign_version_id": assignment.campaign_version_id,
"assignment_id": assignment.id,
"assignment_revision": assignment.resource_revision,
"assignment_ref": (
"campaign-work-assignment:"
f"{assignment.id}:r{assignment.resource_revision}"
),
"outcome": event_kind,
"status": assignment.status,
"assignee_type": assignment.assignee_type,
"assignee_id": assignment.assignee_id,
"action_url": (
f"/campaigns/{assignment.campaign_id}/work"
f"?assignment={assignment.id}"
),
},
),
registry=get_registry(),
)
return event
+8 -2
View File
@@ -134,7 +134,13 @@ class CampaignCollaborationListResponse(BaseModel):
CampaignWorkAssigneeType = Literal["account", "group", "organization_function"]
CampaignWorkAssignmentStatus = Literal["open", "in_progress", "completed", "cancelled"]
CampaignWorkAssignmentStatus = Literal[
"open",
"in_progress",
"completed",
"rejected",
"cancelled",
]
CampaignWorkAssigneeResolutionState = Literal[
"resolved",
"unavailable",
@@ -195,7 +201,7 @@ class CampaignWorkAssignmentTransitionRequest(BaseModel):
model_config = ConfigDict(extra="forbid")
expected_revision: int = Field(ge=1)
action: Literal["start", "complete", "cancel"]
action: Literal["accept", "start", "complete", "reject", "cancel"]
reason: str | None = Field(default=None, max_length=500)
@field_validator("reason")
@@ -0,0 +1,778 @@
from __future__ import annotations
from collections.abc import Mapping
from dataclasses import replace
from datetime import datetime
import hashlib
import json
from fastapi import HTTPException
from sqlalchemy.orm import Session
from govoplan_campaign.backend.db.models import (
Campaign,
CampaignVersion,
CampaignWorkAssignment,
)
from govoplan_campaign.backend.persistence.versions import create_minimal_campaign
from govoplan_campaign.backend.route_support import _get_campaign_for_principal
from govoplan_campaign.backend.routes.assignments import (
_actor_label,
_mirror_assignment_to_tasks,
_notify_assignment,
_record_event,
_require_resolved_assignee,
_resolve_assignee,
)
from govoplan_campaign.backend.schemas import CampaignWorkAssigneeInput
from govoplan_core.audit.logging import audit_from_principal
from govoplan_core.auth import ApiPrincipal, has_scope
from govoplan_core.core.automation import (
ActionDefinition,
ActionExecutionRequest,
ActionExecutionResult,
ActionPreview,
EffectDefinition,
EffectPreview,
ObservedEffect,
)
from govoplan_core.core.campaigns import (
CampaignWorkHandoffInspection,
CampaignWorkHandoffRef,
CampaignWorkHandoffRequest,
)
from govoplan_core.core.notifications import CAPABILITY_NOTIFICATIONS_DISPATCH
from govoplan_core.core.tasks import CAPABILITY_TASK_COMMANDS
from govoplan_core.security.time import utc_now
ACTION_KEY = "campaigns.work.prepare"
ASSIGNMENT_EFFECT = "campaigns.work.assignment_created"
CAMPAIGN_EFFECT = "campaigns.work.campaign_created"
class SqlCampaignWorkOrchestrationProvider:
"""Campaign-owned adapter used through optional Core capabilities only."""
def __init__(self, *, registry: object | None = None) -> None:
self._registry = registry
def action_definitions(self) -> tuple[ActionDefinition, ...]:
return (
ActionDefinition(
action_key=ACTION_KEY,
owner_module="campaigns",
description=(
"Reference or create a Campaign and open one authorization-neutral "
"accountable work hand-off."
),
input_schema_ref="govoplan/campaigns/work-handoff.v1",
required_scopes=(
"campaigns:campaign:read",
"campaigns:campaign:create",
"campaigns:assignment:manage",
),
policy_checks=(
"campaign access is checked independently of assignment",
"the assignee must already have Campaign access",
"the expected Campaign revision must still be current",
),
risk_level="moderate",
reversibility="compensatable",
expected_effect_keys=(ASSIGNMENT_EFFECT, CAMPAIGN_EFFECT),
idempotency_strategy="caller_supplied",
audit_event_types=(
"campaign.assignment.created",
"campaign.created_minimal",
),
preview_required=True,
recovery_mode="atomic",
recovery_verification=(
"resolve the assignment by tenant and orchestration idempotency key",
"verify the exact Campaign version and assignment revisions",
"confirm the assigned principal still has independent Campaign access",
),
),
)
def effect_definitions(self) -> tuple[EffectDefinition, ...]:
return (
EffectDefinition(
effect_key=ASSIGNMENT_EFFECT,
owner_module="campaigns",
operation="created",
description="Create an accountable Campaign work assignment.",
resource_types=("campaign_work_assignment",),
audit_event_types=("campaign.assignment.created",),
compensation_hint="Cancel the open assignment through Campaign work.",
),
EffectDefinition(
effect_key=CAMPAIGN_EFFECT,
owner_module="campaigns",
operation="created",
description="Create a minimal Campaign draft when no campaign is referenced.",
resource_types=("campaign", "campaign_version"),
audit_event_types=("campaign.created_minimal",),
compensation_hint=(
"Delete the untouched draft under the normal Campaign lifecycle policy."
),
),
)
def preview_action(
self,
session: object,
principal: object,
*,
request: ActionExecutionRequest,
) -> ActionPreview:
if request.action_key != ACTION_KEY:
return _blocked_preview("The Campaign work action is not supported.")
try:
sql_session, api_principal = _context(session, principal)
handoff = _request(request)
_preview_handoff(sql_session, api_principal, handoff)
except (HTTPException, TypeError, ValueError) as exc:
return _blocked_preview(_message(exc))
creating = handoff.campaign_id is None
effects = [
EffectPreview(
effect_key=ASSIGNMENT_EFFECT,
summary="Open one revision-bearing Campaign work assignment.",
)
]
if creating:
effects.insert(
0,
EffectPreview(
effect_key=CAMPAIGN_EFFECT,
summary="Create one minimal Campaign draft and initial version.",
),
)
return ActionPreview(
action_key=ACTION_KEY,
allowed=True,
summary=(
"Create a Campaign draft and open accountable work."
if creating
else "Reference the current Campaign revision and open accountable work."
),
risk_level="moderate",
reversibility="compensatable",
effects=tuple(effects),
policy_provenance=(
{
"code": "campaign_assignment_does_not_grant_access",
"assignment_authorization_neutral": True,
"campaign_access_rechecked_on_resume": True,
},
),
preview_ref=f"campaign-work-preview:{_request_hash(handoff)}",
)
def execute_action(
self,
session: object,
principal: object,
*,
request: ActionExecutionRequest,
) -> ActionExecutionResult:
if request.action_key != ACTION_KEY:
raise ValueError("The Campaign work action is not supported.")
sql_session, api_principal = _context(session, principal)
handoff = _request(request)
ref = self.prepare_handoff(
sql_session,
api_principal,
request=handoff,
)
effects = [
ObservedEffect(
effect_key=ASSIGNMENT_EFFECT,
operation="created",
resource_ref=ref.assignment_ref,
summary=(
"Reused the existing idempotent Campaign work assignment."
if ref.replayed
else "Created the Campaign work assignment."
),
metadata={"replayed": ref.replayed},
)
]
if not ref.replayed and handoff.campaign_id is None:
effects.insert(
0,
ObservedEffect(
effect_key=CAMPAIGN_EFFECT,
operation="created",
resource_ref=ref.campaign_ref,
summary="Created the minimal Campaign draft.",
),
)
return ActionExecutionResult(
state="completed",
output=_ref_payload(ref),
observed_effects=tuple(effects),
audit_event_refs=(
str(ref.provenance["audit_event_ref"]),
)
if ref.provenance.get("audit_event_ref")
else (),
)
def prepare_handoff(
self,
session: object,
principal: object,
*,
request: CampaignWorkHandoffRequest,
) -> CampaignWorkHandoffRef:
sql_session, api_principal = _context(session, principal)
if api_principal.tenant_id != request.tenant_id:
raise ValueError("Campaign hand-off tenant does not match the principal")
request_hash = _request_hash(request)
existing = (
sql_session.query(CampaignWorkAssignment)
.filter(
CampaignWorkAssignment.tenant_id == request.tenant_id,
CampaignWorkAssignment.orchestration_idempotency_key
== request.idempotency_key,
)
.one_or_none()
)
if existing is not None:
if existing.orchestration_request_sha256 != request_hash:
raise ValueError(
"Campaign hand-off idempotency key was already used for "
"different input."
)
campaign = _get_campaign_for_principal(
sql_session,
existing.campaign_id,
api_principal,
)
return _handoff_ref(
sql_session,
campaign=campaign,
assignment=existing,
registry=self._registry,
replayed=True,
)
campaign, version, created = _campaign_and_version(
sql_session,
api_principal,
request,
create=True,
)
resolution = _resolve_assignee(
sql_session,
campaign=campaign,
assignee=CampaignWorkAssigneeInput(
type=request.assignee_kind,
id=request.assignee_id,
),
)
_require_resolved_assignee(resolution)
now = utc_now()
assignment = CampaignWorkAssignment(
tenant_id=campaign.tenant_id,
campaign_id=campaign.id,
campaign_version_id=version.id,
reference_kind="campaign_version",
reference_id=version.id,
reference_label=f"Campaign version {version.version_number}",
purpose=request.purpose.strip(),
status="open",
due_at=request.due_at,
assignee_type=request.assignee_kind,
assignee_id=request.assignee_id.strip(),
assignee_label_snapshot=(resolution.label or request.assignee_id)[:500],
assignee_current_label=resolution.label,
assignee_resolution_state=resolution.state,
resolution_provenance={
**resolution.provenance,
"source": "workflow",
"workflow_instance_id": request.workflow_instance_id,
"workflow_step_id": request.workflow_step_id,
"expected_campaign_revision": request.expected_campaign_revision,
},
resolution_checked_at=now,
assigned_by_user_id=api_principal.user.id,
assigned_by_label_snapshot=_actor_label(api_principal),
orchestration_idempotency_key=request.idempotency_key,
orchestration_request_sha256=request_hash,
orchestration_correlation_id=request.correlation_id,
workflow_instance_id=request.workflow_instance_id,
workflow_step_id=request.workflow_step_id,
)
sql_session.add(assignment)
sql_session.flush()
_record_event(
sql_session,
assignment=assignment,
principal=api_principal,
event_kind="assigned",
details={
"source": "workflow",
"workflow_instance_id": request.workflow_instance_id,
"workflow_step_id": request.workflow_step_id,
},
)
if request.mirror_to_tasks:
_mirror_assignment_to_tasks(
sql_session,
campaign=campaign,
assignment=assignment,
principal=api_principal,
)
else:
assignment.task_mirror_status = "skipped"
_notify_assignment(
sql_session,
campaign=campaign,
assignment=assignment,
event_kind="assigned",
)
audit_ref = audit_from_principal(
sql_session,
api_principal,
action="campaign.assignment.created",
object_type="campaign_work_assignment",
object_id=assignment.id,
details={
"campaign_id": campaign.id,
"campaign_version_id": version.id,
"campaign_revision": version.edit_revision,
"resource_revision": assignment.resource_revision,
"source": "workflow",
"workflow_instance_id": request.workflow_instance_id,
"workflow_step_id": request.workflow_step_id,
"assignment_authorization_neutral": True,
"purpose_disclosed": False,
},
correlation_id=request.correlation_id,
causation_id=request.workflow_step_id,
commit=False,
)
if created:
audit_from_principal(
sql_session,
api_principal,
action="campaign.created_minimal",
object_type="campaign",
object_id=campaign.id,
details={
"version_id": version.id,
"external_id": campaign.external_id,
"source": "workflow",
"workflow_instance_id": request.workflow_instance_id,
},
correlation_id=request.correlation_id,
causation_id=request.workflow_step_id,
commit=False,
)
sql_session.flush()
ref = _handoff_ref(
sql_session,
campaign=campaign,
assignment=assignment,
registry=self._registry,
)
return replace(
ref,
provenance={**dict(ref.provenance), "audit_event_ref": audit_ref.id},
)
def inspect_handoff(
self,
session: object,
principal: object,
*,
tenant_id: str,
assignment_id: str,
expected_revision: int | None = None,
) -> CampaignWorkHandoffInspection:
try:
sql_session, api_principal = _context(session, principal)
except TypeError as exc:
return CampaignWorkHandoffInspection(allowed=False, reason=str(exc))
if api_principal.tenant_id != tenant_id:
return CampaignWorkHandoffInspection(
allowed=False,
reason="Campaign hand-off tenant does not match the principal.",
provenance={"code": "campaign_handoff_tenant_mismatch"},
)
assignment = sql_session.get(CampaignWorkAssignment, assignment_id)
if assignment is None or assignment.tenant_id != tenant_id:
return CampaignWorkHandoffInspection(
allowed=False,
reason="Campaign work assignment is unavailable.",
provenance={"code": "campaign_handoff_missing"},
)
try:
_get_campaign_for_principal(
sql_session,
assignment.campaign_id,
api_principal,
)
except HTTPException as exc:
return CampaignWorkHandoffInspection(
allowed=False,
status=assignment.status, # type: ignore[arg-type]
assignment_revision=assignment.resource_revision,
reason=_message(exc),
provenance={
"code": "campaign_handoff_access_revoked",
"campaign_id": assignment.campaign_id,
"assignment_does_not_grant_access": True,
},
)
if (
expected_revision is not None
and assignment.resource_revision != expected_revision
):
return CampaignWorkHandoffInspection(
allowed=False,
status=assignment.status, # type: ignore[arg-type]
assignment_revision=assignment.resource_revision,
action_url=_action_url(assignment),
assignment_ref=_assignment_ref(assignment),
reason="Campaign work assignment revision changed; reload its event.",
provenance={
"code": "campaign_handoff_revision_conflict",
"expected_revision": expected_revision,
"current_revision": assignment.resource_revision,
},
)
return CampaignWorkHandoffInspection(
allowed=True,
status=assignment.status, # type: ignore[arg-type]
assignment_revision=assignment.resource_revision,
action_url=_action_url(assignment),
assignment_ref=_assignment_ref(assignment),
provenance={
"code": "campaign_handoff_access_rechecked",
"campaign_id": assignment.campaign_id,
"assignment_does_not_grant_access": True,
},
)
def _context(
session: object,
principal: object,
) -> tuple[Session, ApiPrincipal]:
if not isinstance(session, Session):
raise TypeError("Campaign work orchestration requires a SQLAlchemy Session.")
if not isinstance(principal, ApiPrincipal):
raise TypeError("Campaign work orchestration requires an API principal.")
return session, principal
def _request(request: ActionExecutionRequest) -> CampaignWorkHandoffRequest:
value = request.input
assignee = value.get("assignee")
if not isinstance(assignee, Mapping):
raise ValueError("Campaign work hand-offs require an assignee object.")
create = value.get("create_campaign")
if create is not None and not isinstance(create, Mapping):
raise ValueError("Campaign creation input must be an object.")
due_at = _date(value.get("due_at"))
return CampaignWorkHandoffRequest(
tenant_id=request.tenant_id,
idempotency_key=request.idempotency_key,
purpose=str(value.get("purpose") or ""),
assignee_kind=str(assignee.get("kind") or ""), # type: ignore[arg-type]
assignee_id=str(assignee.get("id") or ""),
campaign_id=_optional(value.get("campaign_id")),
create_external_id=_optional(create.get("external_id")) if create else None,
create_name=_optional(create.get("name")) if create else None,
create_description=(
_optional(create.get("description")) if create else None
),
expected_campaign_revision=_integer(
value.get("expected_campaign_revision")
),
due_at=due_at,
mirror_to_tasks=bool(value.get("mirror_to_tasks", True)),
correlation_id=request.invocation.correlation_id,
workflow_instance_id=_reference_id(
request.metadata.get("workflow_instance_ref"),
"workflow-instance:",
),
workflow_step_id=_reference_id(
request.metadata.get("workflow_step_ref"),
"workflow-step:",
),
)
def _preview_handoff(
session: Session,
principal: ApiPrincipal,
request: CampaignWorkHandoffRequest,
) -> None:
if principal.tenant_id != request.tenant_id:
raise ValueError("Campaign hand-off tenant does not match the principal")
for scope in (
"campaigns:campaign:read",
"campaigns:campaign:create",
"campaigns:assignment:manage",
):
if not has_scope(principal, scope):
raise ValueError(f"Campaign work hand-off requires {scope}.")
existing = (
session.query(CampaignWorkAssignment)
.filter(
CampaignWorkAssignment.tenant_id == request.tenant_id,
CampaignWorkAssignment.orchestration_idempotency_key
== request.idempotency_key,
)
.one_or_none()
)
if existing is not None:
if existing.orchestration_request_sha256 != _request_hash(request):
raise ValueError(
"Campaign hand-off idempotency key was already used for different input."
)
_get_campaign_for_principal(session, existing.campaign_id, principal)
return
if request.campaign_id is None:
if request.assignee_kind != "account" or (
request.assignee_id != principal.account_id
):
raise ValueError(
"A newly created Campaign can initially be assigned only to its "
"creating account; share it explicitly before assigning other principals."
)
duplicate = (
session.query(Campaign.id)
.filter(
Campaign.tenant_id == request.tenant_id,
Campaign.external_id == request.create_external_id,
)
.first()
)
if duplicate is not None:
raise ValueError("Campaign external ID already exists for this tenant.")
if request.expected_campaign_revision not in {None, 1}:
raise ValueError("A new Campaign starts at revision one.")
return
campaign, _version, _created = _campaign_and_version(
session,
principal,
request,
create=False,
)
resolution = _resolve_assignee(
session,
campaign=campaign,
assignee=CampaignWorkAssigneeInput(
type=request.assignee_kind,
id=request.assignee_id,
),
)
_require_resolved_assignee(resolution)
def _campaign_and_version(
session: Session,
principal: ApiPrincipal,
request: CampaignWorkHandoffRequest,
*,
create: bool,
) -> tuple[Campaign, CampaignVersion, bool]:
if request.campaign_id is None:
if not create:
raise ValueError("Campaign creation is not available during preview.")
campaign, version = create_minimal_campaign(
session,
tenant_id=request.tenant_id,
user_id=principal.user.id,
external_id=str(request.create_external_id),
name=str(request.create_name),
description=request.create_description,
current_flow="create",
current_step="basics",
commit=False,
)
return campaign, version, True
campaign = _get_campaign_for_principal(
session,
request.campaign_id,
principal,
)
version = session.get(CampaignVersion, campaign.current_version_id)
if version is None or version.campaign_id != campaign.id:
raise ValueError("The Campaign current version is unavailable.")
if (
request.expected_campaign_revision is not None
and version.edit_revision != request.expected_campaign_revision
):
raise ValueError(
"Campaign revision changed; reload the Campaign before opening work."
)
return campaign, version, False
def _handoff_ref(
session: Session,
*,
campaign: Campaign,
assignment: CampaignWorkAssignment,
registry: object | None,
replayed: bool = False,
) -> CampaignWorkHandoffRef:
version = session.get(CampaignVersion, assignment.campaign_version_id)
if version is None or version.campaign_id != campaign.id:
raise ValueError("The pinned Campaign hand-off version is unavailable.")
return CampaignWorkHandoffRef(
tenant_id=assignment.tenant_id,
campaign_id=campaign.id,
campaign_version_id=version.id,
campaign_revision=version.edit_revision,
assignment_id=assignment.id,
assignment_revision=assignment.resource_revision,
status=assignment.status, # type: ignore[arg-type]
action_url=_action_url(assignment),
campaign_ref=(
f"campaign:{campaign.id}:version:{version.id}:r{version.edit_revision}"
),
assignment_ref=_assignment_ref(assignment),
replayed=replayed,
optional_capabilities={
"tasks": _has_capability(registry, CAPABILITY_TASK_COMMANDS),
"notifications": _has_capability(
registry,
CAPABILITY_NOTIFICATIONS_DISPATCH,
),
},
provenance={
"assignment_authorization_neutral": True,
"campaign_access_checked": True,
"workflow_instance_id": assignment.workflow_instance_id,
"workflow_step_id": assignment.workflow_step_id,
"correlation_id": assignment.orchestration_correlation_id,
},
)
def _ref_payload(ref: CampaignWorkHandoffRef) -> dict[str, object]:
return {
"campaign_id": ref.campaign_id,
"campaign_version_id": ref.campaign_version_id,
"campaign_revision": ref.campaign_revision,
"assignment_id": ref.assignment_id,
"assignment_revision": ref.assignment_revision,
"status": ref.status,
"action_url": ref.action_url,
"campaign_ref": ref.campaign_ref,
"assignment_ref": ref.assignment_ref,
"event_type": ref.event_type,
"replayed": ref.replayed,
"optional_capabilities": dict(ref.optional_capabilities),
"provenance": dict(ref.provenance),
"outcome": "success",
}
def _request_hash(request: CampaignWorkHandoffRequest) -> str:
payload = {
"tenant_id": request.tenant_id,
"purpose": request.purpose.strip(),
"assignee_kind": request.assignee_kind,
"assignee_id": request.assignee_id.strip(),
"campaign_id": request.campaign_id,
"create_external_id": request.create_external_id,
"create_name": request.create_name,
"create_description": request.create_description,
"expected_campaign_revision": request.expected_campaign_revision,
"due_at": request.due_at.isoformat() if request.due_at else None,
"mirror_to_tasks": request.mirror_to_tasks,
"correlation_id": request.correlation_id,
"workflow_instance_id": request.workflow_instance_id,
"workflow_step_id": request.workflow_step_id,
}
encoded = json.dumps(payload, sort_keys=True, separators=(",", ":"))
return hashlib.sha256(encoded.encode("utf-8")).hexdigest()
def _blocked_preview(reason: str) -> ActionPreview:
return ActionPreview(
action_key=ACTION_KEY,
allowed=False,
summary=reason,
risk_level="moderate",
reversibility="compensatable",
blockers=(reason,),
policy_provenance=(
{
"code": "campaign_work_handoff_blocked",
"reason": reason,
},
),
)
def _message(exc: Exception) -> str:
if isinstance(exc, HTTPException):
detail = exc.detail
if isinstance(detail, Mapping):
return str(detail.get("explanation") or detail.get("code") or detail)
return str(detail)
return str(exc)
def _date(value: object) -> datetime | None:
if value is None or value == "":
return None
if isinstance(value, datetime):
return value
try:
return datetime.fromisoformat(str(value).replace("Z", "+00:00"))
except ValueError as exc:
raise ValueError("Campaign hand-off due date must use ISO 8601.") from exc
def _integer(value: object) -> int | None:
if value is None or value == "":
return None
if isinstance(value, bool):
raise ValueError("Campaign revisions must be integers.")
try:
return int(value)
except (TypeError, ValueError) as exc:
raise ValueError("Campaign revisions must be integers.") from exc
def _optional(value: object) -> str | None:
candidate = str(value or "").strip()
return candidate or None
def _reference_id(value: object, prefix: str) -> str | None:
candidate = str(value or "").strip()
return candidate.removeprefix(prefix) or None if candidate.startswith(prefix) else None
def _assignment_ref(assignment: CampaignWorkAssignment) -> str:
return f"campaign-work-assignment:{assignment.id}:r{assignment.resource_revision}"
def _action_url(assignment: CampaignWorkAssignment) -> str:
return (
f"/campaigns/{assignment.campaign_id}/work"
f"?assignment={assignment.id}"
)
def _has_capability(registry: object | None, name: str) -> bool:
return bool(
registry is not None
and hasattr(registry, "has_capability")
and registry.has_capability(name)
)
__all__ = ["ACTION_KEY", "SqlCampaignWorkOrchestrationProvider"]
@@ -0,0 +1,205 @@
from __future__ import annotations
from govoplan_core.core.campaigns import CAPABILITY_CAMPAIGNS_WORK_ORCHESTRATION
from govoplan_core.core.workflows import WorkflowDefinitionContribution
def campaign_workflow_definitions(
*,
module_version: str,
) -> tuple[WorkflowDefinitionContribution, ...]:
"""Return opt-in Campaign workflow templates owned by this module."""
return (
WorkflowDefinitionContribution(
origin_module_id="campaigns",
origin_module_version=module_version,
definition_key="accountable-campaign-work-handoff",
name="Accountable Campaign work hand-off",
description=(
"Create or reference a Campaign, assign bounded work, and wait "
"for its revision-bearing completion, rejection, cancellation, "
"or timeout event."
),
graph=_campaign_work_handoff_graph(),
definition_kind="template",
scope_type="system",
inherit_to_lower_scopes=True,
allow_start=True,
allow_reuse=True,
allow_automation=False,
execution_mode="guided",
activate_on_install=False,
required_capabilities=(CAPABILITY_CAMPAIGNS_WORK_ORCHESTRATION,),
required_interfaces=("campaigns.work_orchestration",),
metadata={
"domain": "campaigns.accountable_work",
"state_owner": "campaigns",
"template_requires_configuration": True,
},
policy_metadata={
"assignment_authorization_neutral": True,
"campaign_access_rechecked_on_resume": True,
"navigation_does_not_complete_work": True,
},
),
)
def _campaign_work_handoff_graph() -> dict[str, object]:
return {
"schema_version": 1,
"nodes": [
{
"id": "start",
"type": "workflow.start.manual",
"label": "Campaign work requested",
"position": {"x": 20, "y": 140},
"config": {
"input_schema_ref": "govoplan/campaigns/work-handoff.v1",
},
},
{
"id": "prepare",
"type": "workflow.capability",
"label": "Prepare Campaign work",
"position": {"x": 250, "y": 140},
"config": {
"capability": CAPABILITY_CAMPAIGNS_WORK_ORCHESTRATION,
"operation": "campaigns.work.prepare",
"input_mapping": {
"campaign_id": "$input.campaign_id",
"create_campaign": "$input.create_campaign",
"expected_campaign_revision": (
"$input.expected_campaign_revision"
),
"purpose": "$input.purpose",
"assignee": "$input.assignee",
"due_at": "$input.due_at",
"mirror_to_tasks": "$input.mirror_to_tasks",
},
"idempotency_key": "workflow-step",
"failure_policy": "manual",
"view_surface_ids": ["campaigns.page.work"],
},
},
{
"id": "campaign_work",
"type": "workflow.external_handoff",
"label": "Complete Campaign work",
"position": {"x": 510, "y": 140},
"config": {
"provider_capability": (
CAPABILITY_CAMPAIGNS_WORK_ORCHESTRATION
),
"event_type": "campaign.work.changed",
"event_filter": {
"payload": {
"assignment_id": (
"$steps.prepare.execution.output.assignment_id"
)
}
},
"outcome_path": "payload.outcome",
"terminal_outcomes": {
"completed": "completed",
"rejected": "rejected",
"cancelled": "cancelled",
},
"observed_outcomes": [
"assigned",
"accepted",
"started",
"reassigned",
],
"external_id": (
"$steps.prepare.execution.output.assignment_id"
),
"expected_revision": (
"$steps.prepare.execution.output.assignment_revision"
),
"action_url": "$steps.prepare.execution.output.action_url",
"immutable_ref": (
"$steps.prepare.execution.output.assignment_ref"
),
"optional_capabilities": (
"$steps.prepare.execution.output.optional_capabilities"
),
"timeout_after": "$input.timeout_after",
"view_surface_ids": ["campaigns.page.work"],
},
},
{
"id": "completed",
"type": "workflow.end.completed",
"label": "Campaign work completed",
"position": {"x": 790, "y": 20},
"config": {"output_mapping": {}},
},
{
"id": "rejected",
"type": "workflow.end.cancelled",
"label": "Campaign work rejected",
"position": {"x": 790, "y": 120},
"config": {"reason": "Campaign work was rejected"},
},
{
"id": "cancelled",
"type": "workflow.end.cancelled",
"label": "Campaign work cancelled",
"position": {"x": 790, "y": 220},
"config": {"reason": "Campaign work was cancelled"},
},
{
"id": "timed_out",
"type": "workflow.end.cancelled",
"label": "Campaign work timed out",
"position": {"x": 790, "y": 320},
"config": {"reason": "Campaign work timed out"},
},
],
"edges": [
{"id": "start-prepare", "source": "start", "target": "prepare"},
{
"id": "prepare-work",
"source": "prepare",
"source_port": "success",
"target": "campaign_work",
},
{
"id": "work-completed",
"source": "campaign_work",
"source_port": "completed",
"target": "completed",
},
{
"id": "work-rejected",
"source": "campaign_work",
"source_port": "rejected",
"target": "rejected",
},
{
"id": "work-cancelled",
"source": "campaign_work",
"source_port": "cancelled",
"target": "cancelled",
},
{
"id": "work-timeout",
"source": "campaign_work",
"source_port": "timed_out",
"target": "timed_out",
},
],
"metadata": {
"notation": "govoplan.workflow.native",
"domain": "campaigns.accountable_work",
"configuration_notes": (
"Provide either campaign_id or create_campaign and explicit null "
"values for unused optional inputs."
),
},
}
__all__ = ["campaign_workflow_definitions"]
+261 -2
View File
@@ -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"
)
}
+3 -2
View File
@@ -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",
}
+42
View File
@@ -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
+1 -1
View File
@@ -1,6 +1,6 @@
{
"name": "@govoplan/campaign-webui",
"version": "0.1.21",
"version": "0.1.23",
"private": true,
"type": "module",
"main": "src/index.ts",
+2 -2
View File
@@ -82,7 +82,7 @@ export type CampaignCollaborationCreate = {
};
export type CampaignWorkAssigneeType = "account" | "group" | "organization_function";
export type CampaignWorkAssignmentStatus = "open" | "in_progress" | "completed" | "cancelled";
export type CampaignWorkAssignmentStatus = "open" | "in_progress" | "completed" | "rejected" | "cancelled";
export type CampaignWorkAssignmentResolutionState = "resolved" | "unavailable" | "provider_unavailable";
export type CampaignWorkAssignment = {
@@ -1976,7 +1976,7 @@ export async function transitionCampaignWorkAssignment(
settings: ApiSettings,
campaignId: string,
assignment: Pick<CampaignWorkAssignment, "id" | "resource_revision">,
action: "start" | "complete" | "cancel"
action: "accept" | "start" | "complete" | "reject" | "cancel"
): Promise<CampaignWorkAssignment> {
return apiFetch<CampaignWorkAssignment>(
settings,
@@ -1,5 +1,6 @@
import { useCallback, useEffect, useMemo, useState } from "react";
import { CheckCircle2, History, Play, Plus, RefreshCw, UserRoundCog } from "lucide-react";
import { useSearchParams } from "react-router";
import {
Button,
Card,
@@ -59,6 +60,8 @@ export default function CampaignWorkPage({
campaignId: string;
}) {
const workspace = useCampaignWorkspaceData(settings, campaignId);
const [searchParams] = useSearchParams();
const requestedAssignmentId = searchParams.get("assignment");
const [assignments, setAssignments] = useState<CampaignWorkAssignment[]>([]);
const [nextCursor, setNextCursor] = useState<string | null>(null);
const [hasMore, setHasMore] = useState(false);
@@ -71,6 +74,7 @@ export default function CampaignWorkPage({
const [createOpen, setCreateOpen] = useState(false);
const [reassigning, setReassigning] = useState<CampaignWorkAssignment | null>(null);
const [cancelling, setCancelling] = useState<CampaignWorkAssignment | null>(null);
const [rejecting, setRejecting] = useState<CampaignWorkAssignment | null>(null);
const [historyFor, setHistoryFor] = useState<CampaignWorkAssignment | null>(null);
const [history, setHistory] = useState<CampaignWorkAssignmentEvent[]>([]);
const [historyCursor, setHistoryCursor] = useState<string | null>(null);
@@ -102,6 +106,13 @@ export default function CampaignWorkPage({
void loadAssignments();
}, [loadAssignments]);
useEffect(() => {
if (!requestedAssignmentId || loading) return;
document
.getElementById(`campaign-assignment-${requestedAssignmentId}`)
?.focus({ preventScroll: false });
}, [assignments, loading, requestedAssignmentId]);
function replaceAssignment(updated: CampaignWorkAssignment) {
setAssignments((current) => current.map((item) => item.id === updated.id ? updated : item));
}
@@ -175,7 +186,7 @@ export default function CampaignWorkPage({
}
}
async function transition(assignment: CampaignWorkAssignment, action: "start" | "complete" | "cancel") {
async function transition(assignment: CampaignWorkAssignment, action: "accept" | "start" | "complete" | "reject" | "cancel") {
if (busyId) return;
setBusyId(assignment.id);
setError("");
@@ -183,7 +194,15 @@ export default function CampaignWorkPage({
const updated = await transitionCampaignWorkAssignment(settings, campaignId, assignment, action);
replaceAssignment(updated);
setCancelling(null);
setMessage(`Work ${action === "start" ? "started" : action === "complete" ? "completed" : "cancelled"}.`);
setRejecting(null);
const resultLabel = {
accept: "accepted",
start: "started",
complete: "completed",
reject: "rejected",
cancel: "cancelled"
}[action];
setMessage(`Work ${resultLabel}.`);
} catch (err) {
setError(errorText(err));
} finally {
@@ -285,7 +304,12 @@ export default function CampaignWorkPage({
) : (
<ol className="campaign-work-list" aria-label="Campaign work assignments">
{assignments.map((assignment) => (
<li key={assignment.id} className="campaign-work-item">
<li
key={assignment.id}
id={`campaign-assignment-${assignment.id}`}
className={`campaign-work-item${assignment.id === requestedAssignmentId ? " is-focused" : ""}`}
tabIndex={-1}
>
<article>
<header className="campaign-work-item-header">
<div>
@@ -319,8 +343,8 @@ export default function CampaignWorkPage({
<History size={16} aria-hidden="true" /> History
</Button>
{canComplete && assignment.status === "open" ? (
<Button onClick={() => void transition(assignment, "start")} disabled={Boolean(busyId)} helpContextId="campaign.work.action.start" helpModuleId="campaign">
<Play size={16} aria-hidden="true" /> Start
<Button onClick={() => void transition(assignment, "accept")} disabled={Boolean(busyId)} helpContextId="campaign.work.action.start" helpModuleId="campaign">
<Play size={16} aria-hidden="true" /> Accept
</Button>
) : null}
{canComplete && (assignment.status === "open" || assignment.status === "in_progress") ? (
@@ -333,6 +357,13 @@ export default function CampaignWorkPage({
<UserRoundCog size={16} aria-hidden="true" /> Reassign
</Button>
) : null}
{canComplete && (assignment.status === "open" || assignment.status === "in_progress") ? (
<span className="campaign-work-destructive-action">
<Button variant="danger" onClick={() => setRejecting(assignment)} disabled={Boolean(busyId)}>
Reject work
</Button>
</span>
) : null}
{canManage && (assignment.status === "open" || assignment.status === "in_progress") ? (
<span className="campaign-work-destructive-action">
<Button variant="danger" onClick={() => setCancelling(assignment)} disabled={Boolean(busyId)} helpContextId="campaign.work.action.cancel" helpModuleId="campaign">
@@ -391,6 +422,17 @@ export default function CampaignWorkPage({
{historyHasMore ? <Button onClick={() => void loadOlderHistory()} disabled={historyLoading}>Load older history</Button> : null}
</Dialog>
<ConfirmDialog
open={Boolean(rejecting)}
title="Reject assigned work?"
message="The work will close as rejected, separately from cancellation. Its purpose, assignee and transition history remain durable evidence."
confirmLabel="Reject work"
tone="danger"
busy={Boolean(busyId)}
onCancel={() => setRejecting(null)}
onConfirm={() => rejecting ? void transition(rejecting, "reject") : undefined}
/>
<ConfirmDialog
open={Boolean(cancelling)}
title="Cancel assigned work?"
+1
View File
@@ -2802,6 +2802,7 @@
.campaign-work-list,
.campaign-work-history { display: grid; gap: 12px; margin: 0; padding: 0; list-style: none; }
.campaign-work-item { border: var(--border-line); border-radius: var(--radius-sm); background: var(--panel-bg); }
.campaign-work-item.is-focused { box-shadow: var(--focus-ring-strong); }
.campaign-work-item article { display: grid; gap: 14px; padding: 16px; }
.campaign-work-item-header { display: flex; flex-wrap: wrap; align-items: flex-start; justify-content: space-between; gap: 12px; }
.campaign-work-item-header h3 { margin: 0 0 4px; font-size: var(--font-size-md); }