Files
Agent Zero abbe7a18fc fix(audit): P0-P3 audit fixes — 838 ruff errors → 0, 30 F821 bugs fixed, 118 files changed
- P0: hooks.py 3-tuple fix, trigger_dispatcher Contract, contacts/plugin unregister_actions_by_owner
- P0: 5 test files — check_permission mocks removed, hardcoded DB credential → env var
- P1: attachment_service DmsFile via Contract helper, restore_registry/history_hooks dedup
- P1: mail/plugin restore unregister, mcp_client datetime.now(UTC), saved_views/filters patterns
- P1: ProtectedRoute fail-closed, 13 test assertion fixes (bcrypt, DB-URLs, SECRET_KEYs)
- P2: deprecated notifications → post_system_message (3 files), forgejo Base, report_generator lazy import
- P2: webhooks permissions, deps.py/roles.py plugin perms removed, import_export default
- P2: address/tags/entity_links patterns removed, worker.py Contract-Umgehungen fixed
- P2: 28 frontend TODOs (hardcoded constants, deprecated notification API)
- P3: dead code, duplicates, deprecated imports, private attr, __import__ inline
- P3: 8 frontend TODOs (LucideIcons, inline styles, XSS, i18n)
- ruff: 838 → 0 (612 auto-fix + 246 manual + 27 F821 regression fix)
- F821: 30 → 0 (AutomationDefinition, DmsFile, user_id, Path, Any, String)
- Contract-Umgehungen: 2 neue gefunden (worker.py:169, worker.py:280) und gefixt
2026-08-16 01:17:18 +02:00

89 lines
3.1 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="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)
# Register all job functions with the job registry
from app.core.job_registry import register_job # noqa: E402
register_job("check_workflow_timeouts", check_workflow_timeouts)