83 lines
2.9 KiB
Python
83 lines
2.9 KiB
Python
|
|
"""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)
|