from __future__ import annotations import importlib from types import SimpleNamespace from unittest.mock import patch import pytest from alembic.migration import MigrationContext from alembic.operations import Operations from fastapi import HTTPException from sqlalchemy import create_engine, inspect, text from sqlalchemy.orm import Session, sessionmaker from govoplan_access.backend.db.models import Account, Group, User from govoplan_campaign.backend.db.models import ( Campaign, CampaignShare, CampaignVersion, CampaignWorkAssignment, CampaignWorkAssignmentEvent, ) from govoplan_campaign.backend.routes.assignments import ( create_campaign_work_assignment, list_campaign_work_assignment_history, reassign_campaign_work_assignment, reconcile_campaign_work_assignments, transition_campaign_work_assignment, ) from govoplan_campaign.backend.schemas import ( CampaignWorkAssigneeInput, CampaignWorkAssignmentCreateRequest, CampaignWorkAssignmentReassignRequest, CampaignWorkAssignmentTransitionRequest, ) 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 TENANT_ID = "tenant-1" class _Principal: tenant_id = TENANT_ID api_key = None def __init__(self, user_id: str, account_id: str, *scopes: str) -> None: self.user = SimpleNamespace( id=user_id, display_name=f"User {user_id}", email=f"{user_id}@example.test", ) self.account_id = account_id self.scopes = frozenset(scopes) def has(self, scope: str) -> bool: return scope in self.scopes or "tenant:*" in self.scopes class _Directory: def __init__(self) -> None: self.user_2_active = True def users_for_tenant(self, tenant_id: str): if tenant_id != TENANT_ID: return () return ( UserRef(id="user-1", account_id="account-1", tenant_id=TENANT_ID, display_name="Owner"), UserRef( id="user-2", account_id="account-2", tenant_id=TENANT_ID, display_name="Collaborator", status="active" if self.user_2_active else "inactive", ), UserRef(id="user-3", account_id="account-3", tenant_id=TENANT_ID, display_name="Unrelated"), ) def groups_for_tenant(self, tenant_id: str): return (GroupRef(id="group-1", tenant_id=TENANT_ID, name="Campaign group"),) if tenant_id == TENANT_ID else () def groups_for_user(self, user_id: str, *, tenant_id: str): if tenant_id == TENANT_ID and user_id == "user-2": return (GroupRef(id="group-1", tenant_id=TENANT_ID, name="Campaign group"),) return () class _Tasks: def __init__(self) -> None: self.commands = [] def create_task(self, _session, _principal, *, command): self.commands.append(command) return WorkItem( id=f"task-{len(self.commands)}", provider_id="tasks", owner_module="tasks", tenant_id=command.tenant_id, title=command.title, assignments=command.assignments, sources=command.sources, ) class _Notifications: def __init__(self) -> None: self.requests = [] def tenant_id_for_notification(self, _session, *, notification_id: str): del notification_id return TENANT_ID def enqueue_notification(self, _session, request, *, enqueue_delivery: bool = True): self.requests.append((request, enqueue_delivery)) return {"id": f"notification-{len(self.requests)}"} def deliver_notification(self, _session, *, notification_id: str): return {"id": notification_id} def deliver_pending(self, _session, *, tenant_id=None, limit: int = 50): return {"tenant_id": tenant_id, "limit": limit} class _FailingTasks(_Tasks): def create_task(self, _session, _principal, *, command): del command raise RuntimeError("Tasks is temporarily unavailable") class _Registry: def __init__(self, tasks: _Tasks | None = None, notifications: _Notifications | None = None) -> None: self.tasks = tasks self.notifications = notifications def has_capability(self, name: str) -> bool: return ( (name == "tasks.commands" and self.tasks is not None) or (name == "notifications.dispatch" and self.notifications is not None) ) def capability(self, name: str): if name == "tasks.commands": return self.tasks if name == "notifications.dispatch": return self.notifications return None class _Organizations: def get_function(self, function_id: str): if function_id != "function-1": return None return OrganizationFunctionRef( id=function_id, tenant_id=TENANT_ID, organization_unit_id="unit-1", slug="campaign-review", name="Campaign review function", ) class _Idm: def organization_function_assignments_for_function(self, function_id: str, *, tenant_id=None, effective_at=None): del effective_at if function_id == "function-1" and tenant_id == TENANT_ID: return (SimpleNamespace(account_id="account-2", status="active", valid_from=None, valid_until=None),) return () def organization_function_incumbencies(self, function_ids, *, tenant_id, effective_at=None): del effective_at return {item: SimpleNamespace(assignments=self.organization_function_assignments_for_function(item, tenant_id=tenant_id)) for item in function_ids} @pytest.fixture() def session() -> Session: engine = create_engine("sqlite+pysqlite:///:memory:") Base.metadata.create_all( engine, tables=[ Account.__table__, User.__table__, Group.__table__, ChangeSequenceEntry.__table__, Campaign.__table__, CampaignVersion.__table__, CampaignShare.__table__, CampaignWorkAssignment.__table__, CampaignWorkAssignmentEvent.__table__, ], ) factory = sessionmaker(bind=engine, class_=Session, expire_on_commit=False) database = factory() database.add_all( [ Account(id=f"account-{number}", email=f"user-{number}@example.test", normalized_email=f"user-{number}@example.test") for number in range(1, 4) ] ) database.flush() database.add_all( [ User(id=f"user-{number}", tenant_id=TENANT_ID, account_id=f"account-{number}", email=f"user-{number}@example.test") for number in range(1, 4) ] ) database.add(Group(id="group-1", tenant_id=TENANT_ID, slug="campaign-group", name="Campaign group")) campaign = Campaign( id="campaign-1", tenant_id=TENANT_ID, owner_user_id="user-1", external_id="campaign-1", name="Campaign One", current_version_id="version-1", ) database.add_all( [ campaign, CampaignVersion(id="version-1", campaign_id=campaign.id, version_number=1, raw_json={"version": "1.0"}), CampaignShare( id="share-1", tenant_id=TENANT_ID, campaign_id=campaign.id, target_type="group", target_id="group-1", permission="read", ), ] ) database.commit() try: yield database finally: database.close() engine.dispose() def _manager() -> _Principal: return _Principal( "user-1", "account-1", "campaigns:campaign:read", "campaigns:assignment:read", "campaigns:assignment:manage", "campaigns:assignment:complete", ) def _assignee() -> _Principal: return _Principal( "user-2", "account-2", "campaigns:campaign:read", "campaigns:assignment:read", "campaigns:assignment:complete", ) 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() def test_assignment_is_authorization_neutral_and_mirrors_through_optional_tasks(session: Session) -> None: directory = _Directory() tasks = _Tasks() notifications = _Notifications() with ( patch("govoplan_campaign.backend.routes.assignments._access_directory", return_value=directory), patch("govoplan_campaign.backend.routes.assignments.get_registry", return_value=_Registry(tasks, notifications)), patch("govoplan_campaign.backend.routes.assignments.audit_from_principal", side_effect=_commit_audit), ): response = create_campaign_work_assignment( "campaign-1", CampaignWorkAssignmentCreateRequest( purpose="Review the frozen recipient import", assignee=CampaignWorkAssigneeInput(type="account", id="account-2"), reference={"kind": "campaign_version", "id": "version-1"}, ), session, _manager(), ) assert response.assignee_label_snapshot == "Collaborator" assert response.assignee_resolution_state == "resolved" assert response.resolution_provenance["policy_code"] == "assignment_does_not_grant_access" assert response.task_mirror_id == "task-1" assert tasks.commands[0].assignments[0].kind == "account" assert tasks.commands[0].provenance["authorization_neutral"] is True assert notifications.requests[0][0].recipient_id == "account-2" assert notifications.requests[0][0].payload["purpose_disclosed"] is False assert notifications.requests[0][1] is False assert session.query(CampaignShare).count() == 1 assert [item.event_kind for item in session.query(CampaignWorkAssignmentEvent).all()] == ["assigned"] def test_assignment_rejects_target_without_existing_campaign_access(session: Session) -> None: with patch("govoplan_campaign.backend.routes.assignments._access_directory", return_value=_Directory()): with pytest.raises(HTTPException) as exc_info: create_campaign_work_assignment( "campaign-1", CampaignWorkAssignmentCreateRequest( purpose="Should not grant access", assignee=CampaignWorkAssigneeInput(type="account", id="account-3"), ), session, _manager(), ) assert exc_info.value.status_code == 422 assert exc_info.value.detail["code"] == "campaign_assignment_assignee_inaccessible" assert session.query(CampaignWorkAssignment).count() == 0 assert session.query(CampaignShare).count() == 1 def test_optional_tasks_failure_is_recorded_without_blocking_campaign_work(session: Session) -> None: with ( patch("govoplan_campaign.backend.routes.assignments._access_directory", return_value=_Directory()), patch("govoplan_campaign.backend.routes.assignments.get_registry", return_value=_Registry(_FailingTasks())), patch("govoplan_campaign.backend.routes.assignments.audit_from_principal", side_effect=_commit_audit), ): created = create_campaign_work_assignment( "campaign-1", CampaignWorkAssignmentCreateRequest( purpose="Continue even without Tasks", assignee=CampaignWorkAssigneeInput(type="account", id="account-2"), ), session, _manager(), ) assert created.task_mirror_status == "failed" assert "RuntimeError" in (created.task_mirror_error or "") assert session.get(CampaignWorkAssignment, created.id) is not None def test_assignee_can_start_and_complete_but_not_cancel(session: Session) -> None: directory = _Directory() 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.audit_from_principal", side_effect=_commit_audit), ): created = create_campaign_work_assignment( "campaign-1", CampaignWorkAssignmentCreateRequest( purpose="Review campaign", assignee=CampaignWorkAssigneeInput(type="account", id="account-2"), mirror_to_tasks=False, ), session, _manager(), ) started = transition_campaign_work_assignment( "campaign-1", created.id, CampaignWorkAssignmentTransitionRequest(expected_revision=1, action="start"), session, _assignee(), ) completed = transition_campaign_work_assignment( "campaign-1", created.id, CampaignWorkAssignmentTransitionRequest(expected_revision=2, action="complete"), session, _assignee(), ) assert started.status == "in_progress" assert completed.status == "completed" assert completed.resource_revision == 3 assert [item.event_kind for item in session.query(CampaignWorkAssignmentEvent).order_by(CampaignWorkAssignmentEvent.created_at, CampaignWorkAssignmentEvent.id)] == ["assigned", "started", "completed"] def test_reassignment_and_deactivation_reconciliation_preserve_history(session: Session) -> None: directory = _Directory() with ( patch("govoplan_campaign.backend.routes.assignments._access_directory", return_value=directory), patch("govoplan_campaign.backend.routes.assignments.audit_from_principal", side_effect=_commit_audit), ): created = create_campaign_work_assignment( "campaign-1", CampaignWorkAssignmentCreateRequest( purpose="Coordinate delivery", assignee=CampaignWorkAssigneeInput(type="group", id="group-1"), mirror_to_tasks=False, ), session, _manager(), ) reassigned = reassign_campaign_work_assignment( "campaign-1", created.id, CampaignWorkAssignmentReassignRequest( expected_revision=1, assignee=CampaignWorkAssigneeInput(type="account", id="account-2"), reason="Named accountability is now required.", mirror_to_tasks=False, ), session, _manager(), ) directory.user_2_active = False result = reconcile_campaign_work_assignments("campaign-1", 100, session, _manager()) history = list_campaign_work_assignment_history("campaign-1", created.id, 50, None, session, _manager()) 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 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()), patch("govoplan_campaign.backend.routes.assignments.organization_directory", return_value=_Organizations()), patch("govoplan_campaign.backend.routes.assignments._idm_function_directory", return_value=_Idm()), patch("govoplan_campaign.backend.routes.assignments.audit_from_principal", side_effect=_commit_audit), ): created = create_campaign_work_assignment( "campaign-1", CampaignWorkAssignmentCreateRequest( purpose="Approve the recipient segment", assignee=CampaignWorkAssigneeInput(type="organization_function", id="function-1"), mirror_to_tasks=False, ), session, _manager(), ) assert created.assignee_label_snapshot == "Campaign review function" assert created.resolution_provenance["resolved_members"] == 1 assert created.resolution_provenance["all_current_incumbents_authorized"] is True def test_assignment_migration_is_repeatable_and_creates_history_indexes() -> None: migration = importlib.import_module( "govoplan_campaign.backend.migrations.versions." "d8e9f0a1b2c3_v0121_campaign_work_assignments" ) 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(migration, "op", Operations(context)): migration.upgrade() migration.upgrade() inspector = inspect(connection) assignment_columns = {item["name"] for item in inspector.get_columns("campaign_work_assignments")} assignment_indexes = {item["name"] for item in inspector.get_indexes("campaign_work_assignments")} event_indexes = {item["name"] for item in inspector.get_indexes("campaign_work_assignment_events")} assert { "purpose", "assignee_type", "assignee_id", "assignee_label_snapshot", "assignee_resolution_state", "resolution_provenance", "task_mirror_status", "resource_revision", }.issubset(assignment_columns) assert "ix_campaign_work_assignments_campaign_status" in assignment_indexes assert "ix_campaign_work_assignment_events_history" in event_indexes with patch.object(migration, "op", Operations(context)): 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" ) }