Files
leocrm/app/core/worker.py
T
Agent Zero 74936b3972 Phase 5 (v2): Processing-Recovery, Retention-Cleanup, Replay-Delivery-Reset
- recover_stuck_events: Reset processing events stuck >120s back to pending
- cleanup_published_events: Delete published events older than 30 days
- Replay now resets outbox_deliveries for clean retry
- Worker: hourly retention cleanup cron job
- API: /recover-stuck and /cleanup-published endpoints
- process_outbox_batch: auto-recovery at start of each tenant iteration
- 23/23 tests passing (5 new tests)
2026-08-02 23:47:29 +02:00

327 lines
13 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 redis.asyncio as aioredis # noqa: E402
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.
"""
settings = get_settings()
client = aioredis.from_url(settings.redis_url)
token = str(uuid.uuid4())
lock_key = f"leocrm:cron_lock:{job_name}"
try:
acquired = await client.set(lock_key, token, nx=True, ex=ttl_seconds)
return token if acquired else None
finally:
await client.aclose()
async def _release_cron_lock(job_name: str, token: str) -> None:
"""Release a previously acquired cron lock using a safe compare-and-delete."""
settings = get_settings()
client = aioredis.from_url(settings.redis_url)
lock_key = f"leocrm:cron_lock:{job_name}"
try:
# 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())
finally:
await client.aclose()
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 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."""
logger.info("ARQ worker shutting down...")
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,
),
]