"""Workflow timeout checker — cancels timed-out workflow instances.""" from __future__ import annotations import logging from datetime import UTC, datetime from typing import Any from sqlalchemy import select from app.core.db import get_session_factory from app.models.notification import Notification from app.models.workflow import Workflow, WorkflowInstance logger = logging.getLogger(__name__) async def check_workflow_timeouts(ctx: dict[str, Any]) -> None: """Run every 5 min. Queries WorkflowInstance where timeout_at <= now AND status in (pending, in_progress). Sets status='cancelled', completed_at=now. Sends notification to initiator.""" now = datetime.now(UTC) factory = get_session_factory() async with factory() as db: result = await db.execute( select(WorkflowInstance).where( WorkflowInstance.timeout_at <= now, WorkflowInstance.status.in_(["pending", "in_progress"]), ) ) instances = list(result.scalars().all()) if not instances: logger.debug("check_workflow_timeouts: no timed-out instances") return logger.info("check_workflow_timeouts: %d instance(s) timed out", len(instances)) for inst in instances: try: async with factory() as db: # Re-fetch to ensure fresh state result = await db.execute( select(WorkflowInstance).where(WorkflowInstance.id == inst.id) ) instance = result.scalar_one_or_none() if instance is None: continue if instance.status not in ("pending", "in_progress"): continue instance.status = "cancelled" instance.completed_at = now # Get workflow name for notification wf_result = await db.execute( select(Workflow).where(Workflow.id == instance.workflow_id) ) workflow = wf_result.scalar_one_or_none() workflow_name = workflow.name if workflow else "Unknown" # Send notification to initiator if instance.initiated_by: notification = Notification( tenant_id=instance.tenant_id, user_id=instance.initiated_by, type="workflow_timeout", title=f"Workflow '{workflow_name}' cancelled due to timeout", body=f"The workflow instance timed out and was automatically cancelled.", ) db.add(notification) await db.flush() logger.info( "Cancelled timed-out workflow instance %s (workflow=%s)", inst.id, workflow_name, ) except Exception: logger.exception("Failed to process timeout for instance %s", inst.id)