"""Workflow routes — CRUD, instance lifecycle, advance/cancel.""" from __future__ import annotations import uuid from fastapi import APIRouter, Depends, HTTPException, Query, Request, Response, status from sqlalchemy.ext.asyncio import AsyncSession from app.core.db import get_db from app.deps import require_permission from app.schemas.workflow import AdvanceRequest, InstanceCreate, WorkflowCreate, WorkflowUpdate from app.services import workflow_service router = APIRouter(prefix="/api/v1/workflows", tags=["workflows"]) # ─── Workflow CRUD ─── @router.get("") async def list_workflows( page: int = Query(1, ge=1), page_size: int = Query(20, ge=1, le=100), is_active: bool | None = Query(None), db: AsyncSession = Depends(get_db), current_user: dict = Depends(require_permission("workflows:read")), ): """List workflows with pagination.""" tenant_id = uuid.UUID(current_user["tenant_id"]) user_id = uuid.UUID(current_user["user_id"]) is_admin = current_user.get("is_system_admin", False) try: return await workflow_service.list_workflows( db, tenant_id, page=page, page_size=page_size, is_active=is_active, user_id=user_id, is_system_admin=is_admin, ) except PermissionError as e: raise HTTPException(status_code=403, detail=str(e)) from e @router.post("", status_code=status.HTTP_201_CREATED) async def create_workflow( body: WorkflowCreate, db: AsyncSession = Depends(get_db), current_user: dict = Depends(require_permission("workflows:write")), ): """Create a new workflow definition. Requires write permission.""" tenant_id = uuid.UUID(current_user["tenant_id"]) user_id = uuid.UUID(current_user["user_id"]) data = body.model_dump() try: return await workflow_service.create_workflow(db, tenant_id, user_id, data) except PermissionError as e: raise HTTPException(status_code=403, detail=str(e)) from e @router.get("/instances") async def list_instances( page: int = Query(1, ge=1), page_size: int = Query(20, ge=1, le=100), status: str | None = Query(None), db: AsyncSession = Depends(get_db), current_user: dict = Depends(require_permission("workflows:read")), ): """List workflow instances with optional status filter.""" tenant_id = uuid.UUID(current_user["tenant_id"]) user_id = uuid.UUID(current_user["user_id"]) is_admin = current_user.get("is_system_admin", False) try: return await workflow_service.list_instances( db, tenant_id, page=page, page_size=page_size, status_filter=status, user_id=user_id, is_system_admin=is_admin, ) except PermissionError as e: raise HTTPException(status_code=403, detail=str(e)) from e @router.get("/{workflow_id}") async def get_workflow( workflow_id: str, db: AsyncSession = Depends(get_db), current_user: dict = Depends(require_permission("workflows:read")), ): """Get a single workflow by ID.""" tenant_id = uuid.UUID(current_user["tenant_id"]) user_id = uuid.UUID(current_user["user_id"]) is_admin = current_user.get("is_system_admin", False) try: result = await workflow_service.get_workflow(db, tenant_id, workflow_id, user_id=user_id, is_system_admin=is_admin) except PermissionError as e: raise HTTPException(status_code=403, detail=str(e)) from e if result is None: raise HTTPException( status_code=status.HTTP_404_NOT_FOUND, detail={"detail": "Workflow not found", "code": "not_found"}, ) return result @router.patch("/{workflow_id}") async def update_workflow( workflow_id: str, body: WorkflowUpdate, db: AsyncSession = Depends(get_db), current_user: dict = Depends(require_permission("workflows:write")), ): """Update a workflow definition.""" tenant_id = uuid.UUID(current_user["tenant_id"]) user_id = uuid.UUID(current_user["user_id"]) is_admin = current_user.get("is_system_admin", False) data = body.model_dump(exclude_unset=True) try: result = await workflow_service.update_workflow(db, tenant_id, user_id, workflow_id, data, is_system_admin=is_admin) except PermissionError as e: raise HTTPException(status_code=403, detail=str(e)) from e if result is None: raise HTTPException( status_code=status.HTTP_404_NOT_FOUND, detail={"detail": "Workflow not found", "code": "not_found"}, ) return result @router.delete("/{workflow_id}", status_code=status.HTTP_204_NO_CONTENT) async def delete_workflow( workflow_id: str, db: AsyncSession = Depends(get_db), current_user: dict = Depends(require_permission("workflows:write")), ): """Delete a workflow definition.""" tenant_id = uuid.UUID(current_user["tenant_id"]) user_id = uuid.UUID(current_user["user_id"]) is_admin = current_user.get("is_system_admin", False) try: deleted = await workflow_service.delete_workflow(db, tenant_id, user_id, workflow_id, is_system_admin=is_admin) except PermissionError as e: raise HTTPException(status_code=403, detail=str(e)) from e if not deleted: raise HTTPException( status_code=status.HTTP_404_NOT_FOUND, detail={"detail": "Workflow not found", "code": "not_found"}, ) return Response(status_code=status.HTTP_204_NO_CONTENT) # ─── Instance endpoints ─── @router.post("/{workflow_id}/instances", status_code=status.HTTP_201_CREATED) async def create_instance( workflow_id: str, body: InstanceCreate, db: AsyncSession = Depends(get_db), current_user: dict = Depends(require_permission("workflows:write")), ): """Create a new workflow instance.""" tenant_id = uuid.UUID(current_user["tenant_id"]) user_id = uuid.UUID(current_user["user_id"]) is_admin = current_user.get("is_system_admin", False) try: result = await workflow_service.create_instance( db, tenant_id, user_id, workflow_id=workflow_id, context=body.context, timeout_hours=body.timeout_hours, is_system_admin=is_admin, ) except PermissionError as e: raise HTTPException(status_code=403, detail=str(e)) from e if result is None: raise HTTPException( status_code=status.HTTP_404_NOT_FOUND, detail={"detail": "Workflow not found", "code": "not_found"}, ) return result @router.get("/instances/{instance_id}") async def get_instance( instance_id: str, db: AsyncSession = Depends(get_db), current_user: dict = Depends(require_permission("workflows:read")), ): """Get a workflow instance with step history.""" tenant_id = uuid.UUID(current_user["tenant_id"]) user_id = uuid.UUID(current_user["user_id"]) is_admin = current_user.get("is_system_admin", False) try: result = await workflow_service.get_instance(db, tenant_id, instance_id, user_id=user_id, is_system_admin=is_admin) except PermissionError as e: raise HTTPException(status_code=403, detail=str(e)) from e if result is None: raise HTTPException( status_code=status.HTTP_404_NOT_FOUND, detail={"detail": "Instance not found", "code": "not_found"}, ) return result @router.post("/instances/{instance_id}/advance") async def advance_instance( instance_id: str, body: AdvanceRequest, db: AsyncSession = Depends(get_db), current_user: dict = Depends(require_permission("workflows:write")), ): """Advance or reject a workflow instance step. Body decision: "approve" or "reject". Returns 200 with updated instance. """ tenant_id = uuid.UUID(current_user["tenant_id"]) user_id = uuid.UUID(current_user["user_id"]) is_admin = current_user.get("is_system_admin", False) try: result = await workflow_service.advance_instance( db, tenant_id, user_id, instance_id=instance_id, decision=body.decision, comment=body.comment, is_system_admin=is_admin, ) except PermissionError as e: raise HTTPException(status_code=403, detail=str(e)) from e if result is None: raise HTTPException( status_code=status.HTTP_404_NOT_FOUND, detail={"detail": "Instance not found", "code": "not_found"}, ) if "error" in result: raise HTTPException( status_code=status.HTTP_400_BAD_REQUEST, detail={"detail": result["error"], "code": "invalid_state"}, ) return result @router.post("/instances/{instance_id}/cancel") async def cancel_instance( instance_id: str, db: AsyncSession = Depends(get_db), current_user: dict = Depends(require_permission("workflows:write")), ): """Cancel a workflow instance. Returns 200 with cancelled instance.""" tenant_id = uuid.UUID(current_user["tenant_id"]) user_id = uuid.UUID(current_user["user_id"]) is_admin = current_user.get("is_system_admin", False) try: result = await workflow_service.cancel_instance( db, tenant_id, user_id, instance_id=instance_id, is_system_admin=is_admin, ) except PermissionError as e: raise HTTPException(status_code=403, detail=str(e)) from e if result is None: raise HTTPException( status_code=status.HTTP_404_NOT_FOUND, detail={"detail": "Instance not found", "code": "not_found"}, ) if "error" in result: raise HTTPException( status_code=status.HTTP_400_BAD_REQUEST, detail={"detail": result["error"], "code": "invalid_state"}, ) return result # ─── G-RUN: Resume waiting workflow instance ────────────────────────────────── @router.post("/instances/{instance_id}/resume") async def resume_instance( instance_id: str, db: AsyncSession = Depends(get_db), current_user: dict = Depends(require_permission("workflows:write")), ): """Resume a waiting workflow instance (e.g. after wait timer expired).""" tenant_id = uuid.UUID(current_user["tenant_id"]) from app.workflows.engine import WorkflowEngine from app.models.workflow import WorkflowInstance from sqlalchemy import select result = await db.execute( select(WorkflowInstance).where( WorkflowInstance.id == uuid.UUID(instance_id), WorkflowInstance.tenant_id == tenant_id, WorkflowInstance.deleted_at.is_(None), ) ) instance = result.scalar_one_or_none() if instance is None: raise HTTPException( status_code=status.HTTP_404_NOT_FOUND, detail={"detail": "Instance not found", "code": "not_found"}, ) if instance.status != "waiting": raise HTTPException( status_code=status.HTTP_400_BAD_REQUEST, detail={"detail": "Instance is not waiting", "code": "invalid_state"}, ) engine = WorkflowEngine(db, tenant_id) return await engine.resume(instance) # ─── G-MAN: Manual trigger — start workflow from UI button ──────────────────── @router.post("/{workflow_id}/trigger", status_code=status.HTTP_201_CREATED) async def manual_trigger( workflow_id: str, request: Request, db: AsyncSession = Depends(get_db), current_user: dict = Depends(require_permission("workflows:write")), ): """Manually trigger a workflow — starts a new instance with optional context.""" tenant_id = uuid.UUID(current_user["tenant_id"]) user_id = uuid.UUID(current_user["user_id"]) try: body = await request.json() except Exception: body = {} result = await workflow_service.create_instance( db, tenant_id, user_id, workflow_id=workflow_id, context=body, ) if result is None: raise HTTPException( status_code=status.HTTP_404_NOT_FOUND, detail={"detail": "Workflow not found", "code": "not_found"}, ) return result # ─── G-WEB: Incoming webhook trigger ────────────────────────────────────────── @router.post("/webhook/{token}", status_code=status.HTTP_200_OK) async def webhook_trigger( token: str, request: Request, db: AsyncSession = Depends(get_db), ): """Incoming webhook trigger — starts a workflow via secure token.""" from app.models.webhook import Webhook from sqlalchemy import select result = await db.execute( select(Webhook).where( Webhook.token == token, Webhook.is_active.is_(True), ) ) webhook = result.scalar_one_or_none() if webhook is None: raise HTTPException( status_code=status.HTTP_404_NOT_FOUND, detail={"detail": "Webhook not found", "code": "not_found"}, ) try: payload = await request.json() except Exception: payload = {} result = await workflow_service.create_instance( db, webhook.tenant_id, None, workflow_id=str(webhook.workflow_id) if hasattr(webhook, "workflow_id") else str(webhook.entity_id), context={"webhook_payload": payload, "webhook_token": token}, ) if result is None: raise HTTPException( status_code=status.HTTP_404_NOT_FOUND, detail={"detail": "Workflow not found", "code": "not_found"}, ) return {"status": "triggered", "instance": result} # ─── G-LOG: Step history for a workflow instance ───────────────────────────── @router.get("/instances/{instance_id}/history") async def get_instance_history( instance_id: str, db: AsyncSession = Depends(get_db), current_user: dict = Depends(require_permission("workflows:read")), ): """Get step history for a workflow instance (execution log).""" tenant_id = uuid.UUID(current_user["tenant_id"]) from app.models.workflow import WorkflowStepHistory from sqlalchemy import select result = await db.execute( select(WorkflowStepHistory) .where( WorkflowStepHistory.tenant_id == tenant_id, WorkflowStepHistory.instance_id == uuid.UUID(instance_id), ) .order_by(WorkflowStepHistory.created_at.asc()) ) history = result.scalars().all() return { "items": [ { "id": str(h.id), "step_index": h.step_index, "step_type": h.step_type, "action": h.action, "actor_id": str(h.actor_id) if h.actor_id else None, "details": h.details, "created_at": h.created_at.isoformat() if h.created_at else None, } for h in history ], "total": len(history), } # ─── G-APPROVAL: Approval step using central ApprovalRequest ───────────────── @router.post("/instances/{instance_id}/approve") async def approve_workflow_step( instance_id: str, request: Request, db: AsyncSession = Depends(get_db), current_user: dict = Depends(require_permission("workflows:write")), ): """Approve the current approval step of a workflow instance.""" tenant_id = uuid.UUID(current_user["tenant_id"]) user_id = uuid.UUID(current_user["user_id"]) try: body = await request.json() except Exception: body = {} comment = body.get("comment", "") from app.core.approval import create_approval_request, decide_approval from app.models.workflow import WorkflowInstance from sqlalchemy import select result = await db.execute( select(WorkflowInstance).where( WorkflowInstance.id == uuid.UUID(instance_id), WorkflowInstance.tenant_id == tenant_id, WorkflowInstance.deleted_at.is_(None), ) ) instance = result.scalar_one_or_none() if instance is None: raise HTTPException( status_code=status.HTTP_404_NOT_FOUND, detail={"detail": "Instance not found", "code": "not_found"}, ) approval = await create_approval_request( db=db, tenant_id=tenant_id, entity_type="workflow_instance", entity_id=uuid.UUID(instance_id), action="workflow_step_approval", requested_by=user_id, requested_by_type="user", ) await decide_approval( db=db, approval_id=approval["id"], decision="approved", decided_by=user_id, comment=comment, ) return await workflow_service.advance_instance( db, tenant_id, user_id, instance_id=instance_id, decision="approved", comment=comment, is_system_admin=current_user.get("is_system_admin", False), ) @router.post("/instances/{instance_id}/reject") async def reject_workflow_step( instance_id: str, request: Request, db: AsyncSession = Depends(get_db), current_user: dict = Depends(require_permission("workflows:write")), ): """Reject the current approval step of a workflow instance.""" tenant_id = uuid.UUID(current_user["tenant_id"]) user_id = uuid.UUID(current_user["user_id"]) try: body = await request.json() except Exception: body = {} comment = body.get("comment", "") from app.core.approval import create_approval_request, decide_approval from app.models.workflow import WorkflowInstance from sqlalchemy import select result = await db.execute( select(WorkflowInstance).where( WorkflowInstance.id == uuid.UUID(instance_id), WorkflowInstance.tenant_id == tenant_id, WorkflowInstance.deleted_at.is_(None), ) ) instance = result.scalar_one_or_none() if instance is None: raise HTTPException( status_code=status.HTTP_404_NOT_FOUND, detail={"detail": "Instance not found", "code": "not_found"}, ) approval = await create_approval_request( db=db, tenant_id=tenant_id, entity_type="workflow_instance", entity_id=uuid.UUID(instance_id), action="workflow_step_approval", requested_by=user_id, requested_by_type="user", ) await decide_approval( db=db, approval_id=approval["id"], decision="rejected", decided_by=user_id, comment=comment, ) return await workflow_service.cancel_instance( db, tenant_id, user_id, instance_id=instance_id, is_system_admin=current_user.get("is_system_admin", False), ) # ─── G-UI-TEMPL: Template Gallery ───────────────────────────────────────────── WORKFLOW_TEMPLATES = [ { "id": "welcome_email", "name": "Welcome Email", "description": "Send a welcome email when a new contact is created", "trigger_event": "contact.after_create", "steps": [ {"name": "Wait 1 hour", "type": "wait", "config": {"duration_seconds": 3600}}, {"name": "Send Welcome", "type": "mail", "config": { "to": "{{context.email}}", "subject": "Welcome to our service!", "body": "Hello {{context.name}},\n\nWelcome aboard! We're excited to have you.\n\nBest regards,\nThe Team", }}, ], }, { "id": "contact_followup", "name": "Contact Follow-Up", "description": "Create a follow-up task 3 days after contact creation", "trigger_event": "contact.after_create", "steps": [ {"name": "Wait 3 days", "type": "wait", "config": {"duration_seconds": 259200}}, {"name": "Create Follow-Up Task", "type": "crm", "config": { "action": "create_contact", "data": {"title": "Follow up with {{context.name}}", "priority": "medium"}, }}, {"name": "Notify Owner", "type": "notification", "config": { "title": "Follow-up reminder", "body": "Time to follow up with {{context.name}}", }}, ], }, { "id": "approval_chain", "name": "Approval Chain", "description": "Two-step approval: manager approves, then sends notification", "trigger_event": "manual", "steps": [ {"name": "Manager Approval", "type": "approval", "config": {}}, {"name": "Send Result", "type": "mail", "config": { "to": "{{context.initiator_email}}", "subject": "Your request has been approved", "body": "Your request has been approved by management.", }}, {"name": "Log Completion", "type": "event", "config": { "event_name": "approval_chain.completed", "payload": {}, }}, ], }, ] @router.get("/templates") async def list_workflow_templates( current_user: dict = Depends(require_permission("workflows:read")), ): """List available workflow templates (G-UI-TEMPL).""" return {"items": WORKFLOW_TEMPLATES, "total": len(WORKFLOW_TEMPLATES)} @router.post("/templates/{template_id}/instantiate", status_code=status.HTTP_201_CREATED) async def instantiate_template( template_id: str, db: AsyncSession = Depends(get_db), current_user: dict = Depends(require_permission("workflows:write")), ): """Instantiate a workflow template — creates a workflow from the template.""" template = next((t for t in WORKFLOW_TEMPLATES if t["id"] == template_id), None) if template is None: raise HTTPException( status_code=status.HTTP_404_NOT_FOUND, detail={"detail": "Template not found", "code": "not_found"}, ) tenant_id = uuid.UUID(current_user["tenant_id"]) user_id = uuid.UUID(current_user["user_id"]) data = { "name": template["name"], "description": template["description"], "trigger_event": template["trigger_event"], "steps": template["steps"], "is_active": True, } return await workflow_service.create_workflow(db, tenant_id, user_id, data)