5dc6f29ac1
Neues automation Plugin (app/plugins/builtins/automation/): - 7 DB-Modelle: AgentDefinition, AgentVersion, AutomationDefinition, AutomationVersion, AutomationCronJob, AgentRun, AutomationRun - Migration 0001_initial.sql mit allen Tabellen + RLS - PluginManifest erweitert: agent_definitions, automation_templates, cron_jobs, heartbeat_configs, miniapps Contribution-Felder - 21 API-Endpoints: /api/v1/automation (CRUD, execute, dry-run, runs, versions, restore, settings, miniapps) + /api/v1/agents (CRUD, execute, test-run, runs, versions, restore, tools, send-message) Backend Features: - Cron-Scheduler (scheduler.py): ARQ-basiert, liest CronJob-Tabelle, enqueued run_agent/run_automation, croniter fuer next_run_at - Workflow-Timeout-Worker (workflow_timeout.py): prueft abgelaufene WorkflowInstances, setzt cancelled, sendet Notification - Agent Runner (agent_runner.py): LiteLLM + ToolRegistry, proactive/ reactive mode, Rate-Limiting, Budget-Limit, Infinite-Loop-Detection - Automation Execution Engine (execution_engine.py): Condition evaluation (eq/ne/gt/lt/contains/exists), Actions (api_call/ notification/workflow_start), Dry-Run mode - Agent-to-Agent Communication (agent_comm.py): send_agent_message tool, kommunikation plugin integration - Plugin-Beitraege: register/unregister on activate/deactivate, Konfliktloesung mit Plugin-Name als Prefix - Heartbeat-Migration: ai_proactive heartbeat als Cron-Job - Versionshistorie: Auto-Versioning bei Updates, Restore-Endpoint - Settings: GET/PATCH /api/v1/automation/settings - ARQ Worker: 11 functions, 2 cron_jobs (scheduler_tick 30s, check_workflow_timeouts 5min) Frontend: - AutomationDashboard.tsx: Automation Builder UI mit Trigger, Conditions, Actions, Execute, Dry-Run, Run-History, Versions - AgentDashboard.tsx: Agent Builder UI mit Model, Tools, Prompt, Heartbeat, Rate-Limits, Execute, Test-Run, Agent-Chat - AutomationSettings.tsx: Settings + MiniApp-Builder - automation.ts: 24 React Query Hooks - automation.ts types: TypeScript Interfaces - routes/index.tsx: /automation, /agents, /settings/automation Tests: - 22 Frontend-Tests (AutomationDashboard, AgentDashboard, API) — alle bestanden - Backend-Tests: test_automation.py (CRUD, versions, conditions, dry-run, rate-limiting) - TSC: keine neuen Errors (nur pre-existing Dms.tsx) - croniter dependency installiert
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)
|