"""Workflow execution engine — step processing, conditions, approvals, wait/resume. Processes workflow instances by evaluating steps sequentially. Supports step types: action, approval, notification, condition, wait, http, mail, calendar, dms, search, agent, crm, event, webhook. G-RUN: Durable/resumable with resume_at, step_state, lock_owner. G-RETRY: Retry with backoff for failed steps. G-IDEMP: Idempotency key for side-effect steps. G-LOG: Execution log per step (input, output, duration, status). """ from __future__ import annotations import logging import time import uuid from datetime import UTC, datetime from typing import Any from sqlalchemy import select from sqlalchemy.ext.asyncio import AsyncSession from app.core.event_bus import get_event_bus from app.core.notifications import post_system_message from app.models.workflow import Workflow, WorkflowInstance from app.services.workflow_service import ( _instance_to_dict, _log_step_history, create_instance, find_workflows_for_event, ) from app.workflows.step_handlers import StepResult, get_step_handler logger = logging.getLogger(__name__) class WorkflowEngine: """Processes workflow instances through their defined steps. Step types (built-in): - action: Executes a configured action (e.g. create entity, send notification) - approval: Pauses and waits for user approve/reject via API - notification: Sends a notification to specified users - condition: Evaluates a condition and branches accordingly - wait: Pauses for a duration or until a specific time (G-WAIT) - http: Sends an HTTP request with SSRF protection (G-HTTP) - mail: Sends an email via the mail plugin (G-MAIL) - calendar: Creates/updates/deletes calendar events (G-CAL) - dms: Interacts with the document management system (G-DMS) - search: Runs a unified search query (G-SEARCH) - agent: Invokes an autonomous AI agent (G-AGENT) - crm: Creates/updates/deletes contacts and companies (G-CRM) - event: Publishes an event to the event bus (G-EVT) - webhook: Sends an outgoing webhook (G-WEB) """ def __init__(self, db: AsyncSession, tenant_id: uuid.UUID): self.db = db self.tenant_id = tenant_id async def process_step(self, instance: WorkflowInstance) -> dict[str, Any]: """Process the current step of a workflow instance. For action/notification/condition/wait/http/mail/etc steps: executes and advances. For approval steps: sets status to in_progress and waits. For wait steps: sets resume_at and pauses. Returns the updated instance dict. """ wf_result = await self.db.execute( select(Workflow).where(Workflow.id == instance.workflow_id) ) workflow = wf_result.scalar_one_or_none() if workflow is None: return {"error": "Workflow not found", "status_code": 404} steps = workflow.steps or [] if instance.current_step_index >= len(steps): instance.status = "completed" instance.completed_at = datetime.now(UTC) await self.db.flush() # Publish workflow.completed event event_bus = get_event_bus() await event_bus.publish('workflow.completed', { 'workflow_id': str(instance.workflow_id), 'instance_id': str(instance.id), 'tenant_id': str(self.tenant_id), 'initiated_by': str(instance.initiated_by) if instance.initiated_by else None, }) # Post workflow completion to Communication (G-WORK) try: from app.plugins.builtins.contracts import get_contract_registry komm = get_contract_registry().get("kommunikation") if komm and instance.initiated_by: room_title = f"Workflow: {workflow.name if hasattr(workflow, 'name') else str(instance.workflow_id)}" conv_id = await komm.find_locked_room_id( db=self.db, tenant_id=self.tenant_id, plugin_name="workflow", title=room_title, ) if not conv_id: room = await komm.create_plugin_room( db=self.db, tenant_id=self.tenant_id, user_id=instance.initiated_by, plugin_name="workflow", title=room_title, participant_type="workflow", ) conv_id = uuid.UUID(room["conversation_id"]) await komm.send_message( db=self.db, tenant_id=self.tenant_id, conversation_id=conv_id, sender_id=instance.workflow_id, sender_type="system", content=f"Workflow completed: {instance.status}", content_format="text", blocks=[ { "block_type": "action_card", "block_data": { "title": f"Workflow Result: {instance.status}", "description": f"Workflow instance {str(instance.id)[:8]} completed successfully", "actions": [ {"label": "View Details", "action": "view_workflow_instance", "data": {"instance_id": str(instance.id)}}, ], }, "sort_order": 0, } ], metadata={"workflow_id": str(instance.workflow_id), "instance_id": str(instance.id), "status": instance.status}, ) await self.db.flush() except Exception: logger.warning("Failed to post workflow result to communication", exc_info=True) return _instance_to_dict(instance) step = steps[instance.current_step_index] step_type = step.get("type", "action") # Log step entry (G-LOG) step_start_time = time.monotonic() await _log_step_history( self.db, self.tenant_id, instance.id, step_index=instance.current_step_index, step_type=step_type, action="processing", details={"step_name": step.get("name"), "config": step.get("config", {})}, ) # Approval steps pause and wait for user input if step_type == "approval": if instance.status == "pending": instance.status = "in_progress" instance.resume_reason = "approval" await self.db.flush() return _instance_to_dict(instance) # Try registered step handlers first (new G step types) handler = get_step_handler(step_type) if handler is not None: return await self._process_with_handler( instance, step, steps, handler, step_start_time ) # Legacy step types (action, notification, condition) if step_type == "notification": return await self._process_notification(instance, step) elif step_type == "condition": return await self._process_condition(instance, step, steps) elif step_type == "action": return await self._process_action(instance, step, steps) else: return {"error": f"Unknown step type: {step_type}", "status_code": 400} async def _process_with_handler( self, instance: WorkflowInstance, step: dict, steps: list, handler: Any, step_start_time: float, ) -> dict[str, Any]: """Execute a step using a registered step handler (G step types).""" step_type = step.get("type", "action") # ── Decision Guard: check if action requires human review (Punkt 9) ── from app.workflows.decision_guard import check_decision_guard step_config = step.get("config", {}) action_name = step_config.get("action", step_type) guard_result = await check_decision_guard( db=self.db, tenant_id=self.tenant_id, instance_id=instance.id, step_config=step_config, action=action_name, ) if not guard_result["allowed"]: # Guard blocks — create ApprovalRequest and pause workflow try: from app.core.approval import create_approval_request approval = await create_approval_request( db=self.db, tenant_id=self.tenant_id, entity_type="workflow_instance", entity_id=instance.id, action=f"decision_guard:{action_name}", requested_by=instance.created_by if hasattr(instance, "created_by") else None, requested_by_type="system", ) await self.db.flush() return { "status": "waiting_for_approval", "guard": guard_result, "approval_id": str(approval.id), "step_index": instance.current_step_index, "message": guard_result.get("reason", "Human review required"), } except Exception as e: logger.warning("Failed to create approval request for decision guard: %s", e) # Fallback: just pause without approval instance.status = "in_progress" instance.resume_reason = "decision_guard" await self.db.flush() return { "status": "waiting_for_approval", "guard": guard_result, "step_index": instance.current_step_index, "message": guard_result.get("reason", "Human review required"), } try: result: StepResult = await handler( self.db, self.tenant_id, instance, step, ) except Exception as e: logger.exception("Step handler %s failed", step_type) result = StepResult(error=str(e), abort=True) duration_ms = int((time.monotonic() - step_start_time) * 1000) # Handle abort if result.abort: instance.status = "failed" instance.error_message = result.error or "Unknown error" await _log_step_history( self.db, self.tenant_id, instance.id, step_index=instance.current_step_index, step_type=step_type, action="failed", details={ "error": result.error, "duration_ms": duration_ms, "output": result.output, }, ) await self.db.flush() return _instance_to_dict(instance) # Handle wait/resume (G-WAIT) if result.wait_until is not None: instance.status = "waiting" instance.resume_at = result.wait_until instance.resume_reason = result.wait_reason or "wait" await _log_step_history( self.db, self.tenant_id, instance.id, step_index=instance.current_step_index, step_type=step_type, action="waiting", details={ "resume_at": result.wait_until.isoformat(), "reason": result.wait_reason, "duration_ms": duration_ms, }, ) await self.db.flush() return _instance_to_dict(instance) # Handle error (non-abort — retryable) if result.error: instance.retry_count += 1 if instance.retry_count >= instance.max_retries: instance.status = "failed" instance.error_message = result.error await _log_step_history( self.db, self.tenant_id, instance.id, step_index=instance.current_step_index, step_type=step_type, action="failed", details={ "error": result.error, "duration_ms": duration_ms, "retry_count": instance.retry_count, }, ) await self.db.flush() return _instance_to_dict(instance) else: # Retry: stay on same step, set resume_at with backoff import asyncio backoff = min(2 ** instance.retry_count, 60) from datetime import timedelta instance.resume_at = datetime.now(UTC) + timedelta(seconds=backoff) instance.resume_reason = "retry" instance.status = "waiting" await _log_step_history( self.db, self.tenant_id, instance.id, step_index=instance.current_step_index, step_type=step_type, action="retry", details={ "error": result.error, "retry_count": instance.retry_count, "resume_at": instance.resume_at.isoformat(), "duration_ms": duration_ms, }, ) await self.db.flush() return _instance_to_dict(instance) # Handle branch (next_index override) if result.next_index is not None: instance.current_step_index = result.next_index instance.status = "in_progress" elif result.advance: # Advance to next step next_idx = instance.current_step_index + 1 if next_idx >= len(steps): instance.status = "completed" instance.completed_at = datetime.now(UTC) else: instance.current_step_index = next_idx instance.status = "in_progress" # Store step output in step_state (G-CTX) step_key = f"step_{instance.current_step_index}_output" instance.step_state[step_key] = result.output # Log success (G-LOG) await _log_step_history( self.db, self.tenant_id, instance.id, step_index=instance.current_step_index, step_type=step_type, action="executed", details={ "output": result.output, "duration_ms": duration_ms, }, ) await self.db.flush() return _instance_to_dict(instance) async def _process_action( self, instance: WorkflowInstance, step: dict, steps: list ) -> dict[str, Any]: """Process an action step — executes the configured action and advances.""" config = step.get("config", {}) action_type = config.get("action_type", "noop") # Execute action based on type if action_type == "create_notification": user_id = config.get("user_id") or ( str(instance.initiated_by) if instance.initiated_by else None ) if user_id: await post_system_message( self.db, self.tenant_id, uuid.UUID(user_id), config.get("notification_type", "info"), config.get("notification_title", ""), config.get("notification_body", ""), ) elif action_type == "noop": pass # No operation — just advance # Log action executed await _log_step_history( self.db, self.tenant_id, instance.id, step_index=instance.current_step_index, step_type="action", action="executed", details={"action_type": action_type}, ) # Advance to next step next_idx = instance.current_step_index + 1 if next_idx >= len(steps): instance.status = "completed" instance.completed_at = datetime.now(UTC) else: instance.current_step_index = next_idx instance.status = "in_progress" await self.db.flush() return _instance_to_dict(instance) async def _process_notification(self, instance: WorkflowInstance, step: dict) -> dict[str, Any]: """Process a notification step — sends notification and advances.""" config = step.get("config", {}) user_id = config.get("user_id") or ( str(instance.initiated_by) if instance.initiated_by else None ) if user_id: await post_system_message( self.db, self.tenant_id, uuid.UUID(user_id), config.get("notification_type", "workflow_notification"), config.get("title", "Workflow notification"), config.get("body", ""), ) await _log_step_history( self.db, self.tenant_id, instance.id, step_index=instance.current_step_index, step_type="notification", action="sent", details={"user_id": user_id}, ) # Advance — notification steps auto-advance wf_result = await self.db.execute( select(Workflow).where(Workflow.id == instance.workflow_id) ) workflow = wf_result.scalar_one_or_none() steps = workflow.steps if workflow else [] next_idx = instance.current_step_index + 1 if next_idx >= len(steps): instance.status = "completed" instance.completed_at = datetime.now(UTC) else: instance.current_step_index = next_idx await self.db.flush() return _instance_to_dict(instance) async def _process_condition( self, instance: WorkflowInstance, step: dict, steps: list ) -> dict[str, Any]: """Process a condition step — evaluates condition and branches. Config format: { "field": "context_key", "operator": "eq|ne|gt|lt|contains", "value": "expected_value", "on_true_step": optional_index, "on_false_step": optional_index } """ config = step.get("config", {}) field = config.get("field", "") operator = config.get("operator", "eq") expected = config.get("value") actual = instance.context.get(field) condition_met = False if operator == "eq": condition_met = actual == expected elif operator == "ne": condition_met = actual != expected elif operator == "gt": condition_met = actual is not None and expected is not None and actual > expected elif operator == "lt": condition_met = actual is not None and expected is not None and actual < expected elif operator == "contains": condition_met = ( actual is not None and expected in actual if isinstance(actual, str | list) else False ) await _log_step_history( self.db, self.tenant_id, instance.id, step_index=instance.current_step_index, step_type="condition", action="evaluated", details={"field": field, "operator": operator, "condition_met": condition_met}, ) # Branch or advance if condition_met and "on_true_step" in config: instance.current_step_index = config["on_true_step"] elif not condition_met and "on_false_step" in config: instance.current_step_index = config["on_false_step"] else: next_idx = instance.current_step_index + 1 if next_idx >= len(steps): instance.status = "completed" instance.completed_at = datetime.now(UTC) else: instance.current_step_index = next_idx await self.db.flush() return _instance_to_dict(instance) async def resume(self, instance: WorkflowInstance) -> dict[str, Any]: """Resume a waiting workflow instance. Called when resume_at has passed, an approval is decided, or an event/webhook triggers a resume. Clears resume_at/resume_reason and processes the current step. """ if instance.status != "waiting": return _instance_to_dict(instance) # Clear resume state instance.resume_at = None instance.resume_reason = None instance.status = "in_progress" # For wait steps, advance to next step after resume wf_result = await self.db.execute( select(Workflow).where(Workflow.id == instance.workflow_id) ) workflow = wf_result.scalar_one_or_none() if workflow is None: return {"error": "Workflow not found", "status_code": 404} steps = workflow.steps or [] if instance.current_step_index < len(steps): step = steps[instance.current_step_index] step_type = step.get("type", "action") # If it was a wait step, advance to next if step_type == "wait": next_idx = instance.current_step_index + 1 if next_idx >= len(steps): instance.status = "completed" instance.completed_at = datetime.now(UTC) else: instance.current_step_index = next_idx await self.db.flush() # Process the next step return await self.process_step(instance) async def acquire_lock( self, instance: WorkflowInstance, owner: str, ttl_seconds: int = 300 ) -> bool: """Acquire a Redis lock for a workflow instance (G-RUN concurrency). Prevents two workers from processing the same instance simultaneously. """ from app.core.redis import get_redis import redis.asyncio as aioredis try: r = await get_redis() lock_key = f"workflow_lock:{instance.id}" acquired = await r.set( lock_key, owner, nx=True, ex=ttl_seconds, ) if acquired: instance.lock_owner = owner instance.lock_expires_at = datetime.now(UTC) + timedelta(seconds=ttl_seconds) await self.db.flush() return True return False except Exception as e: logger.warning("Failed to acquire workflow lock: %s", e) return True # Fail open — allow processing without lock async def release_lock(self, instance: WorkflowInstance) -> None: """Release the Redis lock for a workflow instance.""" from app.core.redis import get_redis try: r = await get_redis() lock_key = f"workflow_lock:{instance.id}" await r.delete(lock_key) instance.lock_owner = None instance.lock_expires_at = None await self.db.flush() except Exception as e: logger.warning("Failed to release workflow lock: %s", e) async def handle_event( db: AsyncSession, tenant_id: uuid.UUID, event_name: str, payload: dict[str, Any], ) -> list[dict[str, Any]]: """Handle an event by starting matching workflow instances. Called by the event bus integration. Finds all active workflows with trigger_event matching event_name and creates instances. """ workflows = await find_workflows_for_event(db, tenant_id, event_name) instances: list[dict[str, Any]] = [] for wf in workflows: inst = await create_instance( db, tenant_id, uuid.UUID(payload.get("user_id", str(uuid.uuid4()))) if payload.get("user_id") else None or uuid.uuid4(), str(wf.id), context=payload, ) if inst: instances.append(inst) return instances def register_workflow_event_handlers() -> None: """Register event bus handlers for workflow triggers. Subscribes to the event bus to auto-start workflows when events fire. Uses a wildcard '*' subscription to catch ALL events and dynamically check which workflows have a matching trigger_event. Should be called during application startup. """ event_bus = get_event_bus() async def _workflow_event_handler(payload: dict[str, Any]) -> None: """Handle events that may trigger workflows.""" from app.core.db import create_db_session tenant_id_str = payload.get("tenant_id") event_name = payload.get("event", "") if not tenant_id_str or not event_name: return tenant_id = uuid.UUID(tenant_id_str) async with create_db_session(tenant_id) as db: await handle_event(db, tenant_id, event_name, payload) # Subscribe to ALL events via wildcard '*' — the handler dynamically # queries for workflows whose trigger_event matches the published event. event_bus.subscribe('*', _workflow_event_handler) async def find_resumable_workflows( db: AsyncSession, tenant_id: uuid.UUID, ) -> list[WorkflowInstance]: """Find workflow instances that are waiting and their resume_at has passed. Called by the ARQ cron job to resume waiting workflows. """ now = datetime.now(UTC) result = await db.execute( select(WorkflowInstance).where( WorkflowInstance.tenant_id == tenant_id, WorkflowInstance.status == "waiting", WorkflowInstance.resume_at.is_not(None), WorkflowInstance.resume_at <= now, ) ) return list(result.scalars().all())