Files
leocrm/app/core/worker.py
T
Agent Zero b5546ea7bd feat(B.13-B.17+B.16): Error-Handling-Infra, Observability, Graceful Shutdown, Cost-Cap, API Versioning
B.13 Error-Handling-Infrastruktur:
- ErrorCategory Enum (TRANSIENT/PERMANENT/PARTIAL), ApiError erweitert
- Einheitliches Error-Response-Format: {code, detail, field, trace_id, retryable, category}
- 3 FastAPI Exception-Handler (ApiError, HTTPException, unhandled)
- classify_exception() Helper, 6 neue Error-Codes
- 28 Tests in test_error_handling.py

B.14 Observability & trace_id-Korrelation:
- trace_id pro Request (UUID4 short) in structlog contextvars
- X-Trace-Id Response-Header
- Sensitive Fields structlog processor
- llm_complete()/llm_embed() akzeptieren trace_id kwarg
- 12 Tests in test_observability.py

B.15 Graceful Shutdown & Connection Draining:
- _shutdown_event + _inflight_requests Tracking in main.py
- drain_all_connections() in ws_helpers.py
- Worker on_shutdown pausiert WorkflowInstances (status=paused)
- 8 Tests in test_graceful_shutdown.py

B.16 API Versioning Strategie:
- Plugin-Dev-Guide Kapitel 30: URL-basiertes Versioning, Breaking Change Prozess

B.17 Cost Overrun Protection:
- llm_monthly_budget_usd + llm_hard_cutoff Settings
- _check_tenant_budget() vor jedem LLM-Call
- _track_tenant_cost() in Redis (INCRBYFLOAT)
- _check_cost_alerts() bei 50%/80%/100% -> post_system_message()
- 20 Tests in test_cost_protection.py

Total: 68 neue Tests, alle grün. Keine Regressionen.
2026-08-13 21:33:14 +02:00

354 lines
14 KiB
Python

"""ARQ worker configuration and entrypoint."""
from __future__ import annotations
import logging
import traceback
from typing import Any
from arq.connections import RedisSettings
from arq import cron
from app.config import get_settings
from app.core.job_registry import get_all_jobs, get_job, register_job
logger = logging.getLogger(__name__)
# ── Distributed lock helpers ─────────────────────────────────────────────────
# When multiple worker replicas run concurrently, cron jobs must not fire
# on every replica. We use a short-lived Redis SET NX lock per cron call
# so only one replica actually executes the job.
import uuid # noqa: E402
async def _acquire_cron_lock(job_name: str, ttl_seconds: int = 120) -> str | None:
"""Try to acquire a distributed lock for a cron job.
Returns a lock token (random UUID) if acquired, or None if another
replica already holds the lock. The lock auto-expires after
*ttl_seconds* to avoid deadlocks if a worker crashes mid-job.
"""
from app.core.auth import get_redis
client = get_redis()
token = str(uuid.uuid4())
lock_key = f"leocrm:cron_lock:{job_name}"
acquired = await client.set(lock_key, token, nx=True, ex=ttl_seconds)
return token if acquired else None
async def _release_cron_lock(job_name: str, token: str) -> None:
"""Release a previously acquired cron lock using a safe compare-and-delete."""
from app.core.auth import get_redis
client = get_redis()
lock_key = f"leocrm:cron_lock:{job_name}"
# Lua script ensures we only delete if the token matches (avoid
# releasing a lock that was already expired and re-acquired).
script = (
b"if redis.call('get', KEYS[1]) == ARGV[1] "
b"then return redis.call('del', KEYS[1]) "
b"else return 0 end"
)
await client.eval(script, 1, lock_key, token.encode())
def _wrap_cron_with_lock(job_name: str, func: Any, ttl_seconds: int = 120) -> Any:
"""Wrap a cron callable so it acquires a distributed lock first.
If the lock cannot be acquired (another replica is handling it), the
wrapped function is silently skipped.
"""
import functools
@functools.wraps(func)
async def _locked_wrapper(ctx: dict[str, Any], *args: Any, **kwargs: Any) -> Any:
token = await _acquire_cron_lock(job_name, ttl_seconds=ttl_seconds)
if token is None:
logger.debug("Cron job '%s' skipped — lock held by another replica", job_name)
return None
try:
return await func(ctx, *args, **kwargs)
finally:
await _release_cron_lock(job_name, token)
return _locked_wrapper
def _get_redis_settings() -> RedisSettings:
"""Get Redis settings from app config."""
settings = get_settings()
return RedisSettings.from_dsn(settings.redis_url)
async def on_startup(ctx: dict[str, Any]) -> None:
"""Called when worker starts."""
logger.info("ARQ worker starting...")
# Initialize Redis singleton (same as API lifespan)
from app.core.auth import init_redis
await init_redis()
# Initialize service container
from app.core.service_container import get_container
container = get_container()
await container.initialize()
# Initialize plugin registry and discover built-in plugins
from app.plugins.registry import get_registry
from app.core.event_bus import get_event_bus
from app.core.webhook_dispatcher import register_webhook_event_handlers
from sqlalchemy import select as sa_select
from app.models.plugin import Plugin as PluginModel
from sqlalchemy.ext.asyncio import async_sessionmaker
registry = get_registry()
from app.core.db import get_migration_engine
migration_engine = get_migration_engine()
registry.initialize(migration_engine, app=None)
registry.discover_builtins()
event_bus = get_event_bus()
from app.core.db import get_worker_session_factory
async_session = get_worker_session_factory()
# Activate plugins that are marked active in DB (register event handlers)
# RLS fail-closed requires tenant context for tenant-table writes.
# The worker skips plugin activation — cron jobs and contributions
# are registered by the API container's startup. The worker only
# needs event handlers and job processing.
from app.models.tenant import Tenant as TenantModel
from app.core.db import set_tenant_context
async with async_session() as db:
# Load all tenant IDs for per-tenant event handler registration
tenant_result = await db.execute(sa_select(TenantModel.id))
all_tenant_ids = [row[0] for row in tenant_result]
logger.info(f"Worker: loaded {len(all_tenant_ids)} tenants")
# Register event handlers only for active plugins (no DB writes, no cron job registration)
# Load active plugin names from DB (global + tenant-specific)
active_plugin_names: set[str] = set()
async with async_session() as db:
# Global plugins that are marked active
result = await db.execute(
sa_select(PluginModel.name).where(PluginModel.active == True)
)
active_plugin_names = {row[0] for row in result}
logger.info(f"Worker: {len(active_plugin_names)} active plugins: {active_plugin_names}")
for name in registry.resolve_load_order():
plugin = registry.get_plugin(name)
if plugin is None:
continue
# Only register event handlers for active plugins
if name not in active_plugin_names:
logger.debug(f"Worker: skipping event handlers for inactive plugin {name}")
continue
try:
# Just register event handlers, skip DB-writing on_activate
if hasattr(plugin, 'register_event_handlers'):
await plugin.register_event_handlers(event_bus)
logger.info(f"Worker: registered event handlers for {name}")
except Exception as exc:
logger.warning(f"Worker: failed to register event handlers for {name}: {exc}")
# Register webhook dispatcher on the event bus
register_webhook_event_handlers(event_bus)
logger.info("Worker: webhook event handlers registered")
# Register trigger dispatcher — generic event→automation bridge
from app.core.trigger_dispatcher import register_trigger_dispatcher
register_trigger_dispatcher(event_bus)
logger.info("Worker: trigger dispatcher registered")
# Register search providers (normally done by app startup)
try:
from app.plugins.builtins.unified_search.provider_registry import auto_register_providers
factory = async_session
async with factory() as db:
await auto_register_providers(db)
logger.info("Search providers registered for worker")
except Exception:
logger.warning("Failed to register search providers in worker", exc_info=True)
async def on_shutdown(ctx: dict[str, Any]) -> None:
"""Called when worker shuts down.
Pauses any running WorkflowRun instances so they can be resumed after
restart, then closes Redis.
"""
logger.info("ARQ worker shutting down...")
# Pause running workflow instances so they can be resumed after restart
try:
from app.core.db import get_worker_session_factory
from sqlalchemy import select as sa_select
from app.models.workflow import WorkflowInstance
session_factory = get_worker_session_factory()
async with session_factory() as db:
result = await db.execute(
sa_select(WorkflowInstance).where(
WorkflowInstance.status == "running"
)
)
running = result.scalars().all()
if running:
for wf in running:
wf.status = "paused"
await db.commit()
logger.info(f"Paused {len(running)} running workflow(s) for graceful shutdown")
except Exception as exc:
logger.warning(f"Failed to pause running workflows during shutdown: {exc}")
from app.core.auth import close_redis
await close_redis()
# ---------------------------------------------------------------------------
# Lazy-load plugin jobs via importlib so they register themselves with the
# job_registry. This avoids circular imports and keeps the worker decoupled
# from plugin internals.
# ---------------------------------------------------------------------------
def _lazy_register_plugin_jobs() -> None:
"""Import each plugin job module so its register_job() call fires."""
plugin_job_modules = [
"app.core.jobs",
"app.plugins.builtins.unified_search.jobs",
"app.plugins.builtins.ai_proactive.jobs",
"app.plugins.builtins.automation.scheduler",
"app.plugins.builtins.automation.workflow_timeout",
"app.plugins.builtins.automation.agent_runner",
"app.plugins.builtins.automation.execution_engine",
"app.plugins.builtins.tasks.jobs",
]
for mod_name in plugin_job_modules:
try:
import importlib
importlib.import_module(mod_name)
logger.debug("Lazy-loaded plugin jobs from %s", mod_name)
except Exception:
logger.warning("Failed to lazy-load plugin jobs from %s", mod_name, exc_info=True)
# Trigger lazy registration at module level so jobs are available when
# WorkerSettings.functions is evaluated.
_lazy_register_plugin_jobs()
# ── Outbox processor job ────────────────────────────────────────────────────
async def process_outbox_job(ctx: dict[str, Any]) -> None:
"""Poll the transactional outbox and publish pending events.
Uses a distributed Redis lock so only one worker replica processes the
outbox at a time. Runs every 5 seconds.
Processes events per-tenant by setting tenant context for RLS.
"""
from app.core.db import get_worker_session_factory
from app.core.outbox import process_outbox_batch
from sqlalchemy import text as sa_text
factory = get_worker_session_factory()
async with factory() as db:
try:
# Load all tenant IDs for per-tenant outbox processing
tenant_result = await db.execute(sa_text("SELECT id FROM tenants"))
tenant_ids = [row[0] for row in tenant_result]
count = await process_outbox_batch(db, batch_size=50, tenant_ids=tenant_ids)
if count:
logger.info("Outbox: published %d events", count)
except Exception as exc:
logger.error("Outbox processing failed", exc_info=True)
await db.rollback()
# Report to Forgejo
try:
from app.plugins.builtins.forgejo_error_reporter.service import report_error_to_forgejo
await report_error_to_forgejo({
"message": f"[Worker] Outbox processing failed: {exc}",
"stack": traceback.format_exc(),
"context": {"source": "worker_outbox_job"},
})
except Exception:
pass
# Register the outbox job so it appears in get_all_jobs()
register_job("process_outbox", process_outbox_job)
# ── Outbox retention cleanup job ─────────────────────────────────────────────
async def cleanup_outbox_job(ctx: dict[str, Any]) -> None:
"""Delete published outbox events older than 30 days.
Runs hourly to prevent the outbox table from growing indefinitely.
Iterates per-tenant for RLS compliance.
"""
from app.core.db import get_worker_session_factory
from app.core.outbox import cleanup_published_events
from sqlalchemy import text as sa_text
factory = get_worker_session_factory()
async with factory() as db:
try:
tenant_result = await db.execute(sa_text("SELECT id FROM tenants"))
tenant_ids = [row[0] for row in tenant_result]
total_deleted = 0
for tenant_id in tenant_ids:
await db.execute(
sa_text("SELECT set_config('app.current_tenant_id', :tid, true)"),
{"tid": str(tenant_id)},
)
deleted = await cleanup_published_events(db, retention_days=30)
total_deleted += deleted
await db.commit()
if total_deleted:
logger.info("Outbox retention: cleaned up %d published events", total_deleted)
except Exception:
logger.error("Outbox retention cleanup failed", exc_info=True)
await db.rollback()
register_job("cleanup_outbox", cleanup_outbox_job)
class WorkerSettings:
"""ARQ worker settings."""
functions = get_all_jobs()
redis_settings = _get_redis_settings()
on_startup = on_startup
on_shutdown = on_shutdown
max_jobs = 10
max_tries = 3
job_timeout = 300
queue_name = "arq:queue"
cron_jobs = [
cron(
_wrap_cron_with_lock("scheduler_tick", get_job("scheduler_tick")),
minute={0, 5, 10, 15, 20, 25, 30, 35, 40, 45, 50, 55},
),
cron(
_wrap_cron_with_lock("tasks_due_reminder", get_job("tasks_due_reminder")),
hour=8, minute=0,
),
# Outbox processor — every 5 seconds, guarded by distributed lock
cron(
_wrap_cron_with_lock("process_outbox", process_outbox_job, ttl_seconds=30),
second={0, 5, 10, 15, 20, 25, 30, 35, 40, 45, 50, 55},
),
# Outbox retention cleanup — hourly
cron(
_wrap_cron_with_lock("cleanup_outbox", cleanup_outbox_job, ttl_seconds=300),
minute=0,
),
]