"""Tests for Workflow Engine — covers ACs 8-22.""" from __future__ import annotations import uuid from datetime import UTC, datetime, timedelta import pytest from httpx import AsyncClient from sqlalchemy import select from app.models.notification import Notification from app.models.workflow import Workflow, WorkflowInstance, WorkflowStepHistory from app.services.workflow_service import ( advance_instance, auto_reject_timeout, check_timeout, create_instance, create_workflow, delete_workflow, find_workflows_for_event, get_instance, list_workflows, start_instance_for_event, update_workflow, ) from tests.conftest import ORIGIN_HEADER, login_client, seed_tenant_and_users # ─── Workflow CRUD (ACs 8-12) ─── VALID_STEPS = [ {"name": "Step 1", "type": "action", "config": {"action_type": "noop"}}, {"name": "Step 2", "type": "approval", "config": {"required_role": "admin"}}, {"name": "Step 3", "type": "notification", "config": {"title": "Done", "body": "Completed"}}, ] @pytest.mark.asyncio async def test_ac8_create_workflow(client: AsyncClient, db_session): """AC8: POST /api/v1/workflows with valid steps JSONB returns 201 + workflow definition.""" await seed_tenant_and_users(db_session) await login_client(client, "admin@tenanta.com") resp = await client.post( "/api/v1/workflows", json={ "name": "Test Workflow", "description": "A test workflow", "steps": VALID_STEPS, }, headers=ORIGIN_HEADER, ) assert resp.status_code == 201 data = resp.json() assert data["name"] == "Test Workflow" assert len(data["steps"]) == 3 assert data["is_active"] is True assert "id" in data @pytest.mark.asyncio async def test_ac9_list_workflows(client: AsyncClient, db_session): """AC9: GET /api/v1/workflows returns 200 + paginated list.""" await seed_tenant_and_users(db_session) await login_client(client, "admin@tenanta.com") # Create a workflow first await client.post( "/api/v1/workflows", json={"name": "WF1", "steps": VALID_STEPS}, headers=ORIGIN_HEADER, ) resp = await client.get("/api/v1/workflows") assert resp.status_code == 200 data = resp.json() assert "items" in data assert "total" in data assert data["total"] >= 1 assert len(data["items"]) >= 1 @pytest.mark.asyncio async def test_ac10_get_workflow_by_id(client: AsyncClient, db_session): """AC10: GET /api/v1/workflows/{id} returns 200 + workflow detail with steps.""" await seed_tenant_and_users(db_session) await login_client(client, "admin@tenanta.com") create_resp = await client.post( "/api/v1/workflows", json={"name": "Detail WF", "steps": VALID_STEPS}, headers=ORIGIN_HEADER, ) wf_id = create_resp.json()["id"] resp = await client.get(f"/api/v1/workflows/{wf_id}") assert resp.status_code == 200 data = resp.json() assert data["id"] == wf_id assert data["name"] == "Detail WF" assert len(data["steps"]) == 3 @pytest.mark.asyncio async def test_ac11_update_workflow(client: AsyncClient, db_session): """AC11: PATCH /api/v1/workflows/{id} returns 200, updated.""" await seed_tenant_and_users(db_session) await login_client(client, "admin@tenanta.com") create_resp = await client.post( "/api/v1/workflows", json={"name": "Original", "steps": VALID_STEPS}, headers=ORIGIN_HEADER, ) wf_id = create_resp.json()["id"] resp = await client.patch( f"/api/v1/workflows/{wf_id}", json={"name": "Updated Name", "description": "Updated desc"}, headers=ORIGIN_HEADER, ) assert resp.status_code == 200 data = resp.json() assert data["name"] == "Updated Name" assert data["description"] == "Updated desc" @pytest.mark.asyncio async def test_ac12_delete_workflow(client: AsyncClient, db_session): """AC12: DELETE /api/v1/workflows/{id} returns 204.""" await seed_tenant_and_users(db_session) await login_client(client, "admin@tenanta.com") create_resp = await client.post( "/api/v1/workflows", json={"name": "To Delete", "steps": VALID_STEPS}, headers=ORIGIN_HEADER, ) wf_id = create_resp.json()["id"] resp = await client.delete(f"/api/v1/workflows/{wf_id}", headers=ORIGIN_HEADER) assert resp.status_code == 204 # Verify deleted get_resp = await client.get(f"/api/v1/workflows/{wf_id}") assert get_resp.status_code == 404 # ─── Instance Lifecycle (ACs 13-18) ─── @pytest.mark.asyncio async def test_ac13_create_instance(client: AsyncClient, db_session): """AC13: POST /api/v1/workflows/{id}/instances returns 201, instance created with status=pending.""" await seed_tenant_and_users(db_session) await login_client(client, "admin@tenanta.com") create_resp = await client.post( "/api/v1/workflows", json={"name": "Instance WF", "steps": VALID_STEPS}, headers=ORIGIN_HEADER, ) wf_id = create_resp.json()["id"] resp = await client.post( f"/api/v1/workflows/{wf_id}/instances", json={"context": {"key": "value"}}, headers=ORIGIN_HEADER, ) assert resp.status_code == 201 data = resp.json() assert data["status"] == "pending" assert data["current_step_index"] == 0 assert data["workflow_id"] == wf_id @pytest.mark.asyncio async def test_ac14_list_instances_filtered(client: AsyncClient, db_session): """AC14: GET /api/v1/workflows/instances?status=in_progress returns 200 + filtered list.""" await seed_tenant_and_users(db_session) await login_client(client, "admin@tenanta.com") # Create workflow and instance create_resp = await client.post( "/api/v1/workflows", json={"name": "Filter WF", "steps": VALID_STEPS}, headers=ORIGIN_HEADER, ) wf_id = create_resp.json()["id"] await client.post( f"/api/v1/workflows/{wf_id}/instances", json={}, headers=ORIGIN_HEADER, ) # Filter by status=pending (just created) resp = await client.get("/api/v1/workflows/instances?status=pending") assert resp.status_code == 200 data = resp.json() assert all(item["status"] == "pending" for item in data["items"]) assert data["total"] >= 1 @pytest.mark.asyncio async def test_ac15_get_instance_detail(client: AsyncClient, db_session): """AC15: GET /api/v1/workflows/instances/{id} returns 200 + current_step_index + history.""" await seed_tenant_and_users(db_session) await login_client(client, "admin@tenanta.com") create_resp = await client.post( "/api/v1/workflows", json={"name": "Detail Instance WF", "steps": VALID_STEPS}, headers=ORIGIN_HEADER, ) wf_id = create_resp.json()["id"] inst_resp = await client.post( f"/api/v1/workflows/{wf_id}/instances", json={}, headers=ORIGIN_HEADER, ) inst_id = inst_resp.json()["id"] resp = await client.get(f"/api/v1/workflows/instances/{inst_id}") assert resp.status_code == 200 data = resp.json() assert data["id"] == inst_id assert "current_step_index" in data assert "history" in data assert len(data["history"]) >= 1 # At least the initial "entered" entry @pytest.mark.asyncio async def test_ac16_advance_approve(client: AsyncClient, db_session): """AC16: POST /api/v1/workflows/instances/{id}/advance (approve) returns 200, step advanced.""" await seed_tenant_and_users(db_session) await login_client(client, "admin@tenanta.com") create_resp = await client.post( "/api/v1/workflows", json={"name": "Advance WF", "steps": VALID_STEPS}, headers=ORIGIN_HEADER, ) wf_id = create_resp.json()["id"] inst_resp = await client.post( f"/api/v1/workflows/{wf_id}/instances", json={}, headers=ORIGIN_HEADER, ) inst_id = inst_resp.json()["id"] # Advance (approve) — should move from step 0 to step 1 resp = await client.post( f"/api/v1/workflows/instances/{inst_id}/advance", json={"decision": "approve"}, headers=ORIGIN_HEADER, ) assert resp.status_code == 200 data = resp.json() assert data["current_step_index"] == 1 assert data["status"] in ("in_progress",) @pytest.mark.asyncio async def test_ac17_advance_reject(client: AsyncClient, db_session): """AC17: POST /api/v1/workflows/instances/{id}/advance (reject) returns 200, status=rejected, initiator notified.""" await seed_tenant_and_users(db_session) await login_client(client, "admin@tenanta.com") create_resp = await client.post( "/api/v1/workflows", json={"name": "Reject WF", "steps": VALID_STEPS}, headers=ORIGIN_HEADER, ) wf_id = create_resp.json()["id"] inst_resp = await client.post( f"/api/v1/workflows/{wf_id}/instances", json={}, headers=ORIGIN_HEADER, ) inst_id = inst_resp.json()["id"] # Reject resp = await client.post( f"/api/v1/workflows/instances/{inst_id}/advance", json={"decision": "reject", "comment": "Not approved"}, headers=ORIGIN_HEADER, ) assert resp.status_code == 200 data = resp.json() assert data["status"] == "rejected" assert data["completed_at"] is not None # Check notification was created for initiator result = await db_session.execute(select(Notification)) notifications = result.scalars().all() assert len(notifications) >= 1 assert any(n.type == "workflow_rejected" for n in notifications) @pytest.mark.asyncio async def test_ac18_cancel_instance(client: AsyncClient, db_session): """AC18: POST /api/v1/workflows/instances/{id}/cancel returns 200, status=cancelled.""" await seed_tenant_and_users(db_session) await login_client(client, "admin@tenanta.com") create_resp = await client.post( "/api/v1/workflows", json={"name": "Cancel WF", "steps": VALID_STEPS}, headers=ORIGIN_HEADER, ) wf_id = create_resp.json()["id"] inst_resp = await client.post( f"/api/v1/workflows/{wf_id}/instances", json={}, headers=ORIGIN_HEADER, ) inst_id = inst_resp.json()["id"] resp = await client.post( f"/api/v1/workflows/instances/{inst_id}/cancel", headers=ORIGIN_HEADER, ) assert resp.status_code == 200 data = resp.json() assert data["status"] == "cancelled" assert data["completed_at"] is not None # ─── Event-Triggered & Code-Engine (ACs 19-22) ─── @pytest.mark.asyncio async def test_ac19_event_triggered_workflow(db_session): """AC19: Event-triggered workflow — publish event → workflow instance auto-starts.""" seed = await seed_tenant_and_users(db_session) tenant_id = seed["tenant_a"].id admin_id = seed["admin_a"].id # Create a workflow with trigger_event await create_workflow( db_session, tenant_id, admin_id, { "name": "Event Triggered WF", "trigger_event": "company.created", "steps": VALID_STEPS, }, ) # Simulate event by calling start_instance_for_event instances = await start_instance_for_event( db_session, tenant_id, admin_id, "company.created", context={"company_id": "test"}, ) assert len(instances) >= 1 assert instances[0]["status"] == "pending" @pytest.mark.asyncio async def test_ac20_step_history_created(db_session): """AC20: workflow_step_history entry created on every step transition.""" seed = await seed_tenant_and_users(db_session) tenant_id = seed["tenant_a"].id admin_id = seed["admin_a"].id wf = await create_workflow( db_session, tenant_id, admin_id, {"name": "History WF", "steps": VALID_STEPS}, ) wf_id = wf["id"] inst = await create_instance(db_session, tenant_id, admin_id, wf_id) inst_id = inst["id"] # Advance to create more history entries await advance_instance(db_session, tenant_id, admin_id, inst_id, "approve") # Check step history result = await db_session.execute( select(WorkflowStepHistory).where(WorkflowStepHistory.instance_id == uuid.UUID(inst_id)) ) history = result.scalars().all() assert len(history) >= 3 # entered + approved + entered (next step) @pytest.mark.asyncio async def test_ac21_onboarding_workflow_on_user_creation(db_session): """AC21: Code-engine workflow — onboarding workflow runs on user creation.""" from app.workflows.code.onboarding import trigger_onboarding seed = await seed_tenant_and_users(db_session) tenant_id = seed["tenant_a"].id admin_id = seed["admin_a"].id new_user_id = seed["viewer_a"].id # Simulate existing user as "new" # Trigger onboarding instance = await trigger_onboarding(db_session, tenant_id, admin_id, new_user_id) assert instance is not None assert instance["status"] == "pending" assert instance["context"]["new_user_id"] == str(new_user_id) # Verify workflow was created with correct steps result = await db_session.execute( select(Workflow).where(Workflow.trigger_event == "user.created") ) workflows = result.scalars().all() assert len(workflows) >= 1 assert workflows[0].name == "User Onboarding" assert len(workflows[0].steps) == 3 @pytest.mark.asyncio async def test_ac22_approval_timeout_auto_reject(db_session): """AC22: Approval step timeout → auto-reject after configured hours (tested with mock timer).""" seed = await seed_tenant_and_users(db_session) tenant_id = seed["tenant_a"].id admin_id = seed["admin_a"].id # Create workflow with approval step wf = await create_workflow( db_session, tenant_id, admin_id, { "name": "Timeout WF", "steps": [ {"name": "Approval Step", "type": "approval", "config": {}}, ], }, ) # Create instance with timeout_hours=1 inst = await create_instance( db_session, tenant_id, admin_id, wf["id"], timeout_hours=1, ) inst_id = uuid.UUID(inst["id"]) # Manually set timeout_at to past to simulate timeout result = await db_session.execute( select(WorkflowInstance).where(WorkflowInstance.id == inst_id) ) instance = result.scalar_one() instance.timeout_at = datetime.now(UTC) - timedelta(hours=1) await db_session.flush() # Check timeout detection is_timed_out = await check_timeout(instance) assert is_timed_out is True # Auto-reject reject_result = await auto_reject_timeout(db_session, tenant_id, instance) assert reject_result["status"] == "rejected" assert reject_result["completed_at"] is not None # Check notification created result = await db_session.execute(select(Notification)) notifications = result.scalars().all() assert any(n.type == "workflow_timeout" for n in notifications) # ─── Edge Cases ─── @pytest.mark.asyncio async def test_workflow_not_found(client: AsyncClient, db_session): """Edge case: Get non-existent workflow returns 404.""" await seed_tenant_and_users(db_session) await login_client(client, "admin@tenanta.com") resp = await client.get(f"/api/v1/workflows/{uuid.uuid4()}") assert resp.status_code == 404 @pytest.mark.asyncio async def test_instance_not_found(client: AsyncClient, db_session): """Edge case: Get non-existent instance returns 404.""" await seed_tenant_and_users(db_session) await login_client(client, "admin@tenanta.com") resp = await client.get(f"/api/v1/workflows/instances/{uuid.uuid4()}") assert resp.status_code == 404 @pytest.mark.asyncio async def test_cancel_completed_instance_fails(client: AsyncClient, db_session): """Edge case: Cancelling a completed instance returns 400.""" await seed_tenant_and_users(db_session) await login_client(client, "admin@tenanta.com") create_resp = await client.post( "/api/v1/workflows", json={ "name": "Complete WF", "steps": [{"name": "Only step", "type": "action", "config": {"action_type": "noop"}}], }, headers=ORIGIN_HEADER, ) wf_id = create_resp.json()["id"] inst_resp = await client.post( f"/api/v1/workflows/{wf_id}/instances", json={}, headers=ORIGIN_HEADER, ) inst_id = inst_resp.json()["id"] # Advance to complete (single step) await client.post( f"/api/v1/workflows/instances/{inst_id}/advance", json={"decision": "approve"}, headers=ORIGIN_HEADER, ) # Try to cancel completed cancel_resp = await client.post( f"/api/v1/workflows/instances/{inst_id}/cancel", headers=ORIGIN_HEADER, ) assert cancel_resp.status_code == 400 @pytest.mark.asyncio async def test_workflow_tenant_isolation(client: AsyncClient, db_session): """Edge case: Workflow from tenant A not visible to tenant B.""" await seed_tenant_and_users(db_session) await login_client(client, "admin@tenanta.com") create_resp = await client.post( "/api/v1/workflows", json={"name": "Tenant A WF", "steps": VALID_STEPS}, headers=ORIGIN_HEADER, ) wf_id = create_resp.json()["id"] # Login as tenant B from httpx import ASGITransport from httpx import AsyncClient as AC # noqa: N817 import app.main app_instance = app.main.app async with AC(transport=ASGITransport(app=app_instance), base_url="http://test") as client_b: await login_client(client_b, "admin@tenantb.com") resp = await client_b.get(f"/api/v1/workflows/{wf_id}") assert resp.status_code == 404 @pytest.mark.asyncio async def test_workflow_create_rbac_viewer_blocked(client: AsyncClient, db_session): """Edge case: Viewer cannot create workflows (403).""" await seed_tenant_and_users(db_session) await login_client(client, "viewer@tenanta.com") resp = await client.post( "/api/v1/workflows", json={"name": "Viewer WF", "steps": VALID_STEPS}, headers=ORIGIN_HEADER, ) assert resp.status_code == 403 # ─── WorkflowEngine Unit Tests (engine.py direct coverage) ─── @pytest.mark.asyncio async def test_engine_process_action_step_noop(db_session): """Engine: action step with noop config advances to next step.""" from app.workflows.engine import WorkflowEngine seed = await seed_tenant_and_users(db_session) tenant_id = seed["tenant_a"].id admin_id = seed["admin_a"].id wf = await create_workflow( db_session, tenant_id, admin_id, { "name": "Engine Noop WF", "steps": [ {"name": "Action", "type": "action", "config": {"action_type": "noop"}}, {"name": "Approval", "type": "approval", "config": {}}, ], }, ) inst = await create_instance(db_session, tenant_id, admin_id, wf["id"]) # Fetch the instance object result = await db_session.execute( select(WorkflowInstance).where(WorkflowInstance.id == uuid.UUID(inst["id"])) ) instance = result.scalar_one() engine = WorkflowEngine(db_session, tenant_id) result_dict = await engine.process_step(instance) assert result_dict["current_step_index"] == 1 assert result_dict["status"] == "in_progress" @pytest.mark.asyncio async def test_engine_process_action_step_create_notification(db_session): """Engine: action step with create_notification creates a Notification.""" from app.workflows.engine import WorkflowEngine seed = await seed_tenant_and_users(db_session) tenant_id = seed["tenant_a"].id admin_id = seed["admin_a"].id wf = await create_workflow( db_session, tenant_id, admin_id, { "name": "Engine Notif WF", "steps": [ { "name": "Create Notif", "type": "action", "config": { "action_type": "create_notification", "user_id": str(admin_id), "title": "Test Notif", "body": "Hello from engine", "notification_type": "workflow_action", }, }, ], }, ) inst = await create_instance(db_session, tenant_id, admin_id, wf["id"]) result = await db_session.execute( select(WorkflowInstance).where(WorkflowInstance.id == uuid.UUID(inst["id"])) ) instance = result.scalar_one() engine = WorkflowEngine(db_session, tenant_id) await engine.process_step(instance) # Check notification created notif_result = await db_session.execute(select(Notification)) notifications = notif_result.scalars().all() assert len(notifications) >= 1 assert any(n.title == "Test Notif" for n in notifications) @pytest.mark.asyncio async def test_engine_process_action_completes_last_step(db_session): """Engine: action step on last step completes the instance.""" from app.workflows.engine import WorkflowEngine seed = await seed_tenant_and_users(db_session) tenant_id = seed["tenant_a"].id admin_id = seed["admin_a"].id wf = await create_workflow( db_session, tenant_id, admin_id, { "name": "Single Step WF", "steps": [{"name": "Only", "type": "action", "config": {"action_type": "noop"}}], }, ) inst = await create_instance(db_session, tenant_id, admin_id, wf["id"]) result = await db_session.execute( select(WorkflowInstance).where(WorkflowInstance.id == uuid.UUID(inst["id"])) ) instance = result.scalar_one() engine = WorkflowEngine(db_session, tenant_id) result_dict = await engine.process_step(instance) assert result_dict["status"] == "completed" assert result_dict["completed_at"] is not None @pytest.mark.asyncio async def test_engine_process_approval_step(db_session): """Engine: approval step pauses and sets status to in_progress.""" from app.workflows.engine import WorkflowEngine seed = await seed_tenant_and_users(db_session) tenant_id = seed["tenant_a"].id admin_id = seed["admin_a"].id wf = await create_workflow( db_session, tenant_id, admin_id, { "name": "Approval WF", "steps": [{"name": "Approval", "type": "approval", "config": {}}], }, ) inst = await create_instance(db_session, tenant_id, admin_id, wf["id"]) result = await db_session.execute( select(WorkflowInstance).where(WorkflowInstance.id == uuid.UUID(inst["id"])) ) instance = result.scalar_one() assert instance.status == "pending" engine = WorkflowEngine(db_session, tenant_id) result_dict = await engine.process_step(instance) assert result_dict["status"] == "in_progress" @pytest.mark.asyncio async def test_engine_process_notification_step(db_session): """Engine: notification step sends notification and advances.""" from app.workflows.engine import WorkflowEngine seed = await seed_tenant_and_users(db_session) tenant_id = seed["tenant_a"].id admin_id = seed["admin_a"].id wf = await create_workflow( db_session, tenant_id, admin_id, { "name": "Notif Step WF", "steps": [ { "name": "Notify", "type": "notification", "config": { "user_id": str(admin_id), "title": "Step Notif", "body": "You have a notification", "notification_type": "wf_notif", }, }, ], }, ) inst = await create_instance(db_session, tenant_id, admin_id, wf["id"]) result = await db_session.execute( select(WorkflowInstance).where(WorkflowInstance.id == uuid.UUID(inst["id"])) ) instance = result.scalar_one() engine = WorkflowEngine(db_session, tenant_id) result_dict = await engine.process_step(instance) assert result_dict["status"] == "completed" notif_result = await db_session.execute(select(Notification)) notifications = notif_result.scalars().all() assert any(n.title == "Step Notif" for n in notifications) @pytest.mark.asyncio async def test_engine_process_condition_eq_true(db_session): """Engine: condition step with eq operator, condition met, branches to on_true_step.""" from app.workflows.engine import WorkflowEngine seed = await seed_tenant_and_users(db_session) tenant_id = seed["tenant_a"].id admin_id = seed["admin_a"].id wf = await create_workflow( db_session, tenant_id, admin_id, { "name": "Cond Eq WF", "steps": [ { "name": "Check Status", "type": "condition", "config": { "field": "status", "operator": "eq", "value": "active", "on_true_step": 2, "on_false_step": 1, }, }, {"name": "False Step", "type": "action", "config": {"action_type": "noop"}}, {"name": "True Step", "type": "action", "config": {"action_type": "noop"}}, ], }, ) inst = await create_instance( db_session, tenant_id, admin_id, wf["id"], context={"status": "active"}, ) result = await db_session.execute( select(WorkflowInstance).where(WorkflowInstance.id == uuid.UUID(inst["id"])) ) instance = result.scalar_one() engine = WorkflowEngine(db_session, tenant_id) result_dict = await engine.process_step(instance) assert result_dict["current_step_index"] == 2 # on_true_step @pytest.mark.asyncio async def test_engine_process_condition_eq_false(db_session): """Engine: condition step with eq operator, condition not met, branches to on_false_step.""" from app.workflows.engine import WorkflowEngine seed = await seed_tenant_and_users(db_session) tenant_id = seed["tenant_a"].id admin_id = seed["admin_a"].id wf = await create_workflow( db_session, tenant_id, admin_id, { "name": "Cond Eq False WF", "steps": [ { "name": "Check Status", "type": "condition", "config": { "field": "status", "operator": "eq", "value": "active", "on_true_step": 2, "on_false_step": 1, }, }, {"name": "False Step", "type": "action", "config": {"action_type": "noop"}}, {"name": "True Step", "type": "action", "config": {"action_type": "noop"}}, ], }, ) inst = await create_instance( db_session, tenant_id, admin_id, wf["id"], context={"status": "inactive"}, ) result = await db_session.execute( select(WorkflowInstance).where(WorkflowInstance.id == uuid.UUID(inst["id"])) ) instance = result.scalar_one() engine = WorkflowEngine(db_session, tenant_id) result_dict = await engine.process_step(instance) assert result_dict["current_step_index"] == 1 # on_false_step @pytest.mark.asyncio async def test_engine_process_condition_ne(db_session): """Engine: condition step with ne (not equal) operator.""" from app.workflows.engine import WorkflowEngine seed = await seed_tenant_and_users(db_session) tenant_id = seed["tenant_a"].id admin_id = seed["admin_a"].id wf = await create_workflow( db_session, tenant_id, admin_id, { "name": "Cond NE WF", "steps": [ { "name": "Check", "type": "condition", "config": { "field": "level", "operator": "ne", "value": "low", }, }, {"name": "Next", "type": "action", "config": {"action_type": "noop"}}, ], }, ) inst = await create_instance( db_session, tenant_id, admin_id, wf["id"], context={"level": "high"}, ) result = await db_session.execute( select(WorkflowInstance).where(WorkflowInstance.id == uuid.UUID(inst["id"])) ) instance = result.scalar_one() engine = WorkflowEngine(db_session, tenant_id) result_dict = await engine.process_step(instance) assert result_dict["current_step_index"] == 1 # condition met, no branch → advance @pytest.mark.asyncio async def test_engine_process_condition_gt(db_session): """Engine: condition step with gt (greater than) operator.""" from app.workflows.engine import WorkflowEngine seed = await seed_tenant_and_users(db_session) tenant_id = seed["tenant_a"].id admin_id = seed["admin_a"].id wf = await create_workflow( db_session, tenant_id, admin_id, { "name": "Cond GT WF", "steps": [ { "name": "Check Count", "type": "condition", "config": {"field": "count", "operator": "gt", "value": 5}, }, {"name": "Next", "type": "action", "config": {"action_type": "noop"}}, ], }, ) inst = await create_instance( db_session, tenant_id, admin_id, wf["id"], context={"count": 10}, ) result = await db_session.execute( select(WorkflowInstance).where(WorkflowInstance.id == uuid.UUID(inst["id"])) ) instance = result.scalar_one() engine = WorkflowEngine(db_session, tenant_id) result_dict = await engine.process_step(instance) assert result_dict["current_step_index"] == 1 # 10 > 5, no branch → advance @pytest.mark.asyncio async def test_engine_process_condition_lt(db_session): """Engine: condition step with lt (less than) operator.""" from app.workflows.engine import WorkflowEngine seed = await seed_tenant_and_users(db_session) tenant_id = seed["tenant_a"].id admin_id = seed["admin_a"].id wf = await create_workflow( db_session, tenant_id, admin_id, { "name": "Cond LT WF", "steps": [ { "name": "Check Count", "type": "condition", "config": {"field": "count", "operator": "lt", "value": 100}, }, {"name": "Next", "type": "action", "config": {"action_type": "noop"}}, ], }, ) inst = await create_instance( db_session, tenant_id, admin_id, wf["id"], context={"count": 50}, ) result = await db_session.execute( select(WorkflowInstance).where(WorkflowInstance.id == uuid.UUID(inst["id"])) ) instance = result.scalar_one() engine = WorkflowEngine(db_session, tenant_id) result_dict = await engine.process_step(instance) assert result_dict["current_step_index"] == 1 # 50 < 100, no branch → advance @pytest.mark.asyncio async def test_engine_process_condition_contains(db_session): """Engine: condition step with contains operator on string.""" from app.workflows.engine import WorkflowEngine seed = await seed_tenant_and_users(db_session) tenant_id = seed["tenant_a"].id admin_id = seed["admin_a"].id wf = await create_workflow( db_session, tenant_id, admin_id, { "name": "Cond Contains WF", "steps": [ { "name": "Check Tags", "type": "condition", "config": {"field": "tags", "operator": "contains", "value": "vip"}, }, {"name": "Next", "type": "action", "config": {"action_type": "noop"}}, ], }, ) inst = await create_instance( db_session, tenant_id, admin_id, wf["id"], context={"tags": "vip,premium"}, ) result = await db_session.execute( select(WorkflowInstance).where(WorkflowInstance.id == uuid.UUID(inst["id"])) ) instance = result.scalar_one() engine = WorkflowEngine(db_session, tenant_id) result_dict = await engine.process_step(instance) assert result_dict["current_step_index"] == 1 # contains vip, no branch → advance @pytest.mark.asyncio async def test_engine_process_condition_contains_list(db_session): """Engine: condition step with contains operator on list.""" from app.workflows.engine import WorkflowEngine seed = await seed_tenant_and_users(db_session) tenant_id = seed["tenant_a"].id admin_id = seed["admin_a"].id wf = await create_workflow( db_session, tenant_id, admin_id, { "name": "Cond List WF", "steps": [ { "name": "Check List", "type": "condition", "config": {"field": "roles", "operator": "contains", "value": "admin"}, }, {"name": "Next", "type": "action", "config": {"action_type": "noop"}}, ], }, ) inst = await create_instance( db_session, tenant_id, admin_id, wf["id"], context={"roles": ["admin", "editor"]}, ) result = await db_session.execute( select(WorkflowInstance).where(WorkflowInstance.id == uuid.UUID(inst["id"])) ) instance = result.scalar_one() engine = WorkflowEngine(db_session, tenant_id) result_dict = await engine.process_step(instance) assert result_dict["current_step_index"] == 1 @pytest.mark.asyncio async def test_engine_process_condition_no_match_advances(db_session): """Engine: condition not met and no on_false_step → advances to next step.""" from app.workflows.engine import WorkflowEngine seed = await seed_tenant_and_users(db_session) tenant_id = seed["tenant_a"].id admin_id = seed["admin_a"].id wf = await create_workflow( db_session, tenant_id, admin_id, { "name": "Cond No Branch WF", "steps": [ { "name": "Check", "type": "condition", "config": {"field": "status", "operator": "eq", "value": "active"}, }, {"name": "Next", "type": "action", "config": {"action_type": "noop"}}, ], }, ) inst = await create_instance( db_session, tenant_id, admin_id, wf["id"], context={"status": "inactive"}, ) result = await db_session.execute( select(WorkflowInstance).where(WorkflowInstance.id == uuid.UUID(inst["id"])) ) instance = result.scalar_one() engine = WorkflowEngine(db_session, tenant_id) result_dict = await engine.process_step(instance) assert result_dict["current_step_index"] == 1 # no branch, advance @pytest.mark.asyncio async def test_engine_process_condition_completes_last_step(db_session): """Engine: condition on last step with no branch → completes.""" from app.workflows.engine import WorkflowEngine seed = await seed_tenant_and_users(db_session) tenant_id = seed["tenant_a"].id admin_id = seed["admin_a"].id wf = await create_workflow( db_session, tenant_id, admin_id, { "name": "Cond Last WF", "steps": [ { "name": "Check", "type": "condition", "config": {"field": "status", "operator": "eq", "value": "active"}, }, ], }, ) inst = await create_instance( db_session, tenant_id, admin_id, wf["id"], context={"status": "active"}, ) result = await db_session.execute( select(WorkflowInstance).where(WorkflowInstance.id == uuid.UUID(inst["id"])) ) instance = result.scalar_one() engine = WorkflowEngine(db_session, tenant_id) result_dict = await engine.process_step(instance) assert result_dict["status"] == "completed" @pytest.mark.asyncio async def test_engine_unknown_step_type(db_session): """Engine: unknown step type returns error with status_code 400.""" from app.workflows.engine import WorkflowEngine seed = await seed_tenant_and_users(db_session) tenant_id = seed["tenant_a"].id admin_id = seed["admin_a"].id wf = await create_workflow( db_session, tenant_id, admin_id, { "name": "Unknown Type WF", "steps": [{"name": "Bad", "type": "unknown_type", "config": {}}], }, ) inst = await create_instance(db_session, tenant_id, admin_id, wf["id"]) result = await db_session.execute( select(WorkflowInstance).where(WorkflowInstance.id == uuid.UUID(inst["id"])) ) instance = result.scalar_one() engine = WorkflowEngine(db_session, tenant_id) result_dict = await engine.process_step(instance) assert "error" in result_dict assert result_dict["status_code"] == 400 @pytest.mark.asyncio async def test_engine_workflow_not_found(db_session): """Engine: process_step with non-existent workflow returns 404 error.""" from unittest.mock import MagicMock, patch from app.models.workflow import WorkflowInstance from app.workflows.engine import WorkflowEngine seed = await seed_tenant_and_users(db_session) tenant_id = seed["tenant_a"].id admin_id = seed["admin_a"].id # Create a real workflow and instance wf = await create_workflow( db_session, tenant_id, admin_id, { "name": "Temp WF for Not Found Test", "steps": [{"name": "Step", "type": "action", "config": {"action_type": "noop"}}], }, ) inst = await create_instance(db_session, tenant_id, admin_id, wf["id"]) result = await db_session.execute( select(WorkflowInstance).where(WorkflowInstance.id == uuid.UUID(inst["id"])) ) instance = result.scalar_one() # Mock the db.execute to return None for the workflow lookup original_execute = db_session.execute async def mock_execute(stmt, *args, **kwargs): # Check if this is a Workflow select query stmt_str = str(stmt) if "FROM workflows" in stmt_str or "workflow" in stmt_str.lower(): mock_result = MagicMock() mock_result.scalar_one_or_none = MagicMock(return_value=None) return mock_result return await original_execute(stmt, *args, **kwargs) engine = WorkflowEngine(db_session, tenant_id) with patch.object(db_session, "execute", side_effect=mock_execute): result_dict = await engine.process_step(instance) assert "error" in result_dict assert result_dict["status_code"] == 404 @pytest.mark.asyncio async def test_engine_step_index_beyond_steps_completes(db_session): """Engine: current_step_index >= len(steps) marks instance as completed.""" from app.workflows.engine import WorkflowEngine seed = await seed_tenant_and_users(db_session) tenant_id = seed["tenant_a"].id admin_id = seed["admin_a"].id wf = await create_workflow( db_session, tenant_id, admin_id, { "name": "Beyond Steps WF", "steps": [{"name": "Step", "type": "action", "config": {"action_type": "noop"}}], }, ) inst = await create_instance(db_session, tenant_id, admin_id, wf["id"]) result = await db_session.execute( select(WorkflowInstance).where(WorkflowInstance.id == uuid.UUID(inst["id"])) ) instance = result.scalar_one() instance.current_step_index = 5 # Beyond steps length await db_session.flush() engine = WorkflowEngine(db_session, tenant_id) result_dict = await engine.process_step(instance) assert result_dict["status"] == "completed" @pytest.mark.asyncio async def test_engine_handle_event(db_session): """Engine: handle_event creates instances for matching workflows.""" from app.workflows.engine import handle_event seed = await seed_tenant_and_users(db_session) tenant_id = seed["tenant_a"].id admin_id = seed["admin_a"].id await create_workflow( db_session, tenant_id, admin_id, { "name": "Event Handler WF", "trigger_event": "company.created", "steps": VALID_STEPS, }, ) instances = await handle_event( db_session, tenant_id, "company.created", {"user_id": str(admin_id), "company_id": "test"}, ) assert len(instances) >= 1 assert instances[0]["status"] == "pending" @pytest.mark.asyncio async def test_engine_handle_event_no_match(db_session): """Engine: handle_event returns empty list when no workflows match.""" from app.workflows.engine import handle_event seed = await seed_tenant_and_users(db_session) tenant_id = seed["tenant_a"].id instances = await handle_event( db_session, tenant_id, "nonexistent.event", {"user_id": str(seed["admin_a"].id)}, ) assert instances == [] def test_engine_register_workflow_event_handlers(): """Engine: register_workflow_event_handlers subscribes to event bus.""" from collections import defaultdict from app.core.event_bus import get_event_bus from app.workflows.engine import register_workflow_event_handlers event_bus = get_event_bus() # Reset handlers to a fresh defaultdict for test isolation event_bus._handlers = defaultdict(list) register_workflow_event_handlers() assert "user.created" in event_bus._handlers assert "company.created" in event_bus._handlers assert "contact.created" in event_bus._handlers assert len(event_bus._handlers["user.created"]) >= 1 # ─── Route Error Path Tests ─── @pytest.mark.asyncio async def test_route_update_workflow_not_found(client: AsyncClient, db_session): """Route: PATCH non-existent workflow returns 404.""" await seed_tenant_and_users(db_session) await login_client(client, "admin@tenanta.com") resp = await client.patch( f"/api/v1/workflows/{uuid.uuid4()}", json={"name": "Updated"}, headers=ORIGIN_HEADER, ) assert resp.status_code == 404 @pytest.mark.asyncio async def test_route_update_workflow_rbac_blocked(client: AsyncClient, db_session): """Route: PATCH workflow as viewer returns 403.""" await seed_tenant_and_users(db_session) await login_client(client, "admin@tenanta.com") create_resp = await client.post( "/api/v1/workflows", json={"name": "RBAC Update WF", "steps": VALID_STEPS}, headers=ORIGIN_HEADER, ) wf_id = create_resp.json()["id"] # Login as viewer from httpx import ASGITransport from httpx import AsyncClient as AC # noqa: N817 import app.main app_instance = app.main.app async with AC( transport=ASGITransport(app=app_instance), base_url="http://test" ) as viewer_client: await login_client(viewer_client, "viewer@tenanta.com") resp = await viewer_client.patch( f"/api/v1/workflows/{wf_id}", json={"name": "Viewer Updated"}, headers=ORIGIN_HEADER, ) assert resp.status_code == 403 @pytest.mark.asyncio async def test_route_delete_workflow_not_found(client: AsyncClient, db_session): """Route: DELETE non-existent workflow returns 404.""" await seed_tenant_and_users(db_session) await login_client(client, "admin@tenanta.com") resp = await client.delete(f"/api/v1/workflows/{uuid.uuid4()}", headers=ORIGIN_HEADER) assert resp.status_code == 404 @pytest.mark.asyncio async def test_route_delete_workflow_rbac_blocked(client: AsyncClient, db_session): """Route: DELETE workflow as viewer returns 403.""" await seed_tenant_and_users(db_session) await login_client(client, "admin@tenanta.com") create_resp = await client.post( "/api/v1/workflows", json={"name": "RBAC Delete WF", "steps": VALID_STEPS}, headers=ORIGIN_HEADER, ) wf_id = create_resp.json()["id"] from httpx import ASGITransport from httpx import AsyncClient as AC # noqa: N817 import app.main app_instance = app.main.app async with AC( transport=ASGITransport(app=app_instance), base_url="http://test" ) as viewer_client: await login_client(viewer_client, "viewer@tenanta.com") resp = await viewer_client.delete(f"/api/v1/workflows/{wf_id}", headers=ORIGIN_HEADER) assert resp.status_code == 403 @pytest.mark.asyncio async def test_route_create_instance_workflow_not_found(client: AsyncClient, db_session): """Route: POST instance for non-existent workflow returns 404.""" await seed_tenant_and_users(db_session) await login_client(client, "admin@tenanta.com") resp = await client.post( f"/api/v1/workflows/{uuid.uuid4()}/instances", json={"context": {}}, headers=ORIGIN_HEADER, ) assert resp.status_code == 404 @pytest.mark.asyncio async def test_route_advance_instance_not_found(client: AsyncClient, db_session): """Route: POST advance for non-existent instance returns 404.""" await seed_tenant_and_users(db_session) await login_client(client, "admin@tenanta.com") resp = await client.post( f"/api/v1/workflows/instances/{uuid.uuid4()}/advance", json={"decision": "approve"}, headers=ORIGIN_HEADER, ) assert resp.status_code == 404 @pytest.mark.asyncio async def test_route_cancel_instance_not_found(client: AsyncClient, db_session): """Route: POST cancel for non-existent instance returns 404.""" await seed_tenant_and_users(db_session) await login_client(client, "admin@tenanta.com") resp = await client.post( f"/api/v1/workflows/instances/{uuid.uuid4()}/cancel", headers=ORIGIN_HEADER, ) assert resp.status_code == 404 @pytest.mark.asyncio async def test_route_advance_completed_instance_returns_400(client: AsyncClient, db_session): """Route: POST advance on completed instance returns 400.""" await seed_tenant_and_users(db_session) await login_client(client, "admin@tenanta.com") create_resp = await client.post( "/api/v1/workflows", json={ "name": "Single Step", "steps": [{"name": "Only", "type": "action", "config": {"action_type": "noop"}}], }, headers=ORIGIN_HEADER, ) wf_id = create_resp.json()["id"] inst_resp = await client.post( f"/api/v1/workflows/{wf_id}/instances", json={}, headers=ORIGIN_HEADER, ) inst_id = inst_resp.json()["id"] # Advance to complete await client.post( f"/api/v1/workflows/instances/{inst_id}/advance", json={"decision": "approve"}, headers=ORIGIN_HEADER, ) # Try to advance again — should get 400 resp = await client.post( f"/api/v1/workflows/instances/{inst_id}/advance", json={"decision": "approve"}, headers=ORIGIN_HEADER, ) assert resp.status_code == 400 @pytest.mark.asyncio async def test_route_cancel_rejected_instance_returns_400(client: AsyncClient, db_session): """Route: POST cancel on rejected instance returns 400.""" await seed_tenant_and_users(db_session) await login_client(client, "admin@tenanta.com") create_resp = await client.post( "/api/v1/workflows", json={"name": "Reject Cancel WF", "steps": VALID_STEPS}, headers=ORIGIN_HEADER, ) wf_id = create_resp.json()["id"] inst_resp = await client.post( f"/api/v1/workflows/{wf_id}/instances", json={}, headers=ORIGIN_HEADER, ) inst_id = inst_resp.json()["id"] # Reject first await client.post( f"/api/v1/workflows/instances/{inst_id}/advance", json={"decision": "reject"}, headers=ORIGIN_HEADER, ) # Try to cancel rejected instance resp = await client.post( f"/api/v1/workflows/instances/{inst_id}/cancel", headers=ORIGIN_HEADER, ) assert resp.status_code == 400 # ─── Service Edge Case Tests ─── @pytest.mark.asyncio async def test_service_check_timeout_no_timeout_at(db_session): """Service: check_timeout returns False when timeout_at is None.""" seed = await seed_tenant_and_users(db_session) tenant_id = seed["tenant_a"].id admin_id = seed["admin_a"].id wf = await create_workflow( db_session, tenant_id, admin_id, { "name": "No Timeout WF", "steps": VALID_STEPS, }, ) inst = await create_instance(db_session, tenant_id, admin_id, wf["id"]) result = await db_session.execute( select(WorkflowInstance).where(WorkflowInstance.id == uuid.UUID(inst["id"])) ) instance = result.scalar_one() is_timed_out = await check_timeout(instance) assert is_timed_out is False @pytest.mark.asyncio async def test_service_check_timeout_completed_status(db_session): """Service: check_timeout returns False when status is completed.""" seed = await seed_tenant_and_users(db_session) tenant_id = seed["tenant_a"].id admin_id = seed["admin_a"].id wf = await create_workflow( db_session, tenant_id, admin_id, { "name": "Completed Timeout WF", "steps": VALID_STEPS, }, ) inst = await create_instance( db_session, tenant_id, admin_id, wf["id"], timeout_hours=1, ) result = await db_session.execute( select(WorkflowInstance).where(WorkflowInstance.id == uuid.UUID(inst["id"])) ) instance = result.scalar_one() instance.status = "completed" instance.timeout_at = datetime.now(UTC) - timedelta(hours=1) await db_session.flush() is_timed_out = await check_timeout(instance) assert is_timed_out is False @pytest.mark.asyncio async def test_service_auto_reject_no_initiated_by(db_session): """Service: auto_reject_timeout works when initiated_by is None.""" seed = await seed_tenant_and_users(db_session) tenant_id = seed["tenant_a"].id admin_id = seed["admin_a"].id wf = await create_workflow( db_session, tenant_id, admin_id, { "name": "No Initiator WF", "steps": VALID_STEPS, }, ) inst = await create_instance(db_session, tenant_id, admin_id, wf["id"], timeout_hours=1) result = await db_session.execute( select(WorkflowInstance).where(WorkflowInstance.id == uuid.UUID(inst["id"])) ) instance = result.scalar_one() instance.initiated_by = None instance.timeout_at = datetime.now(UTC) - timedelta(hours=1) await db_session.flush() reject_result = await auto_reject_timeout(db_session, tenant_id, instance) assert reject_result["status"] == "rejected" @pytest.mark.asyncio async def test_service_find_workflows_no_match(db_session): """Service: find_workflows_for_event returns empty list when no match.""" seed = await seed_tenant_and_users(db_session) tenant_id = seed["tenant_a"].id workflows = await find_workflows_for_event(db_session, tenant_id, "nonexistent.event") assert workflows == [] @pytest.mark.asyncio async def test_service_start_instance_no_match(db_session): """Service: start_instance_for_event returns empty list when no matching workflows.""" seed = await seed_tenant_and_users(db_session) tenant_id = seed["tenant_a"].id admin_id = seed["admin_a"].id instances = await start_instance_for_event( db_session, tenant_id, admin_id, "nonexistent.event", context={}, ) assert instances == [] @pytest.mark.asyncio async def test_service_list_workflows_with_active_filter(db_session): """Service: list_workflows with is_active=True filter.""" seed = await seed_tenant_and_users(db_session) tenant_id = seed["tenant_a"].id admin_id = seed["admin_a"].id # Create active workflow await create_workflow( db_session, tenant_id, admin_id, { "name": "Active WF", "steps": VALID_STEPS, "is_active": True, }, ) # Create inactive workflow await create_workflow( db_session, tenant_id, admin_id, { "name": "Inactive WF", "steps": VALID_STEPS, "is_active": False, }, ) result = await list_workflows(db_session, tenant_id, is_active=True) assert result["total"] == 1 assert result["items"][0]["name"] == "Active WF" @pytest.mark.asyncio async def test_service_update_workflow_with_steps(db_session): """Service: update_workflow with steps update replaces steps array.""" seed = await seed_tenant_and_users(db_session) tenant_id = seed["tenant_a"].id admin_id = seed["admin_a"].id wf = await create_workflow( db_session, tenant_id, admin_id, { "name": "Steps Update WF", "steps": VALID_STEPS, }, ) new_steps = [{"name": "New Step", "type": "action", "config": {"action_type": "noop"}}] result = await update_workflow( db_session, tenant_id, admin_id, wf["id"], { "steps": new_steps, "is_active": False, }, ) assert result is not None assert len(result["steps"]) == 1 assert result["steps"][0]["name"] == "New Step" assert result["is_active"] is False @pytest.mark.asyncio async def test_service_update_workflow_not_found(db_session): """Service: update_workflow returns None for non-existent workflow.""" seed = await seed_tenant_and_users(db_session) tenant_id = seed["tenant_a"].id admin_id = seed["admin_a"].id result = await update_workflow( db_session, tenant_id, admin_id, str(uuid.uuid4()), {"name": "X"} ) assert result is None @pytest.mark.asyncio async def test_service_delete_workflow_not_found(db_session): """Service: delete_workflow returns False for non-existent workflow.""" seed = await seed_tenant_and_users(db_session) tenant_id = seed["tenant_a"].id admin_id = seed["admin_a"].id result = await delete_workflow(db_session, tenant_id, admin_id, str(uuid.uuid4())) assert result is False @pytest.mark.asyncio async def test_service_get_instance_not_found(db_session): """Service: get_instance returns None for non-existent instance.""" seed = await seed_tenant_and_users(db_session) tenant_id = seed["tenant_a"].id result = await get_instance(db_session, tenant_id, str(uuid.uuid4())) assert result is None