fix: Gate 5 — worker event handlers and per-tenant outbox processing

Worker fixes:
- registry.initialize uses get_migration_engine() for DDL (not worker_engine)
- Worker session uses get_worker_session_factory() (crm_worker, not crm_api)
- Event handlers only registered for active plugins (is_active check)
- Outbox processing per-tenant with set_config(app.current_tenant_id)
- process_outbox_job uses get_worker_session_factory() and loads tenant_ids
- Removed unused get_engine import

Outbox fixes:
- process_outbox_batch iterates over tenants, sets RLS context per tenant
- _process_single_outbox_event extracted for clarity
- Events claimed per-tenant (RLS-compatible, no BYPASSRLS needed)
- Commit after each tenant to release locks

Gate 5 requirements met:
- Plugin event handlers registered for active plugins only
- No plugin routers registered in worker
- Outbox events without handlers marked as no_handlers
- Failed consumers trigger retry with exponential backoff
- Processing is idempotent (consumer_inbox check)
- Every worker DB access sets app.current_tenant_id
- Worker cannot read/write other tenant data (RLS enforced)
This commit is contained in:
Agent Zero
2026-07-31 23:09:25 +02:00
parent 89fe7a4750
commit cea21ff576
2 changed files with 160 additions and 115 deletions
+128 -104
View File
@@ -139,25 +139,123 @@ async def enqueue_outbox_event(
) )
async def _process_single_outbox_event(
db: AsyncSession,
event_bus,
row: tuple,
) -> bool:
"""Process a single outbox event. Returns True if published successfully."""
event_id = row[0]
tenant_id = row[1]
event_name = row[2]
payload = row[3]
attempts = row[4]
max_attempts = row[5]
aggregate_type = row[6] if len(row) > 6 else None
aggregate_id = row[7] if len(row) > 7 else None
occurred_at = row[8] if len(row) > 8 else None
correlation_id = row[9] if len(row) > 9 else None
schema_version = row[10] if len(row) > 10 else 1
# payload comes back as a dict from JSONB
if isinstance(payload, str):
import json
payload_dict = json.loads(payload)
else:
payload_dict = payload
try:
# Enrich payload with standardized event envelope metadata
payload_dict.setdefault("_event_id", str(event_id))
payload_dict.setdefault("_event_name", event_name)
payload_dict.setdefault("_event_timestamp", datetime.now(timezone.utc).isoformat())
payload_dict.setdefault("_tenant_id", str(tenant_id))
payload_dict.setdefault("_aggregate_type", aggregate_type)
payload_dict.setdefault("_aggregate_id", str(aggregate_id) if aggregate_id else None)
payload_dict.setdefault("_occurred_at", occurred_at.isoformat() if occurred_at else None)
payload_dict.setdefault("_correlation_id", str(correlation_id) if correlation_id else None)
payload_dict.setdefault("_schema_version", schema_version)
# Idempotency check: has this event already been processed? (P1.5 fix)
already_processed = await db.execute(
text("SELECT 1 FROM consumer_inbox WHERE event_id = :eid AND status = 'processed' LIMIT 1"),
{"eid": str(event_id)},
)
if already_processed.first():
# Event was already processed by all consumers — mark as published
await db.execute(_MARK_PUBLISHED_SQL, {"id": str(event_id)})
logger.debug("Outbox event %s already processed, marking as published", event_id)
return True
results = await event_bus.publish_with_results(event_name, payload_dict)
# Check if any handlers were registered at all
handler_count = len(results)
# If any handler raised, treat as failure
handler_errors = [r for r in results if r is not None]
if handler_errors:
raise handler_errors[0]
if handler_count == 0:
# No handlers registered — mark as 'no_handlers' not 'published'
await db.execute(
text("UPDATE event_outbox SET status = 'no_handlers', published_at = now() WHERE id = :id"),
{"id": str(event_id)},
)
logger.warning("Outbox event %s (%s) had no handlers registered", event_id, event_name)
else:
# Record in consumer_inbox for idempotency (P1.5 fix)
await db.execute(
text("INSERT INTO consumer_inbox (event_id, consumer_name, status, processed_at) VALUES (:eid, :name, 'processed', now()) ON CONFLICT DO NOTHING"),
{"eid": str(event_id), "name": event_name},
)
await db.execute(_MARK_PUBLISHED_SQL, {"id": str(event_id)})
return True
except Exception as exc:
logger.error(
"Failed to publish outbox event %s (%s): %s",
event_id, event_name, exc,
exc_info=True,
)
new_attempts = attempts + 1
if new_attempts >= max_attempts:
await db.execute(_FAIL_SQL, {"id": str(event_id)})
logger.warning(
"Outbox event %s marked as failed after %d attempts",
event_id, new_attempts,
)
else:
backoff = timedelta(seconds=(2 ** new_attempts) * 10)
next_retry = datetime.now(timezone.utc) + backoff
await db.execute(
_RETRY_SQL,
{
"id": str(event_id),
"attempts": new_attempts,
"next_retry_at": next_retry,
},
)
return False
async def process_outbox_batch( async def process_outbox_batch(
db: AsyncSession, db: AsyncSession,
redis: aioredis.Redis | None = None, redis: aioredis.Redis | None = None,
batch_size: int = 50, batch_size: int = 50,
tenant_ids: list[uuid.UUID] | None = None,
) -> int: ) -> int:
"""Process one batch of pending outbox events. """Process one batch of pending outbox events.
1. Claim up to *batch_size* pending events using ``FOR UPDATE SKIP LOCKED`` Iterates over all tenants, setting tenant context for RLS before
so multiple workers don't interfere. claiming and processing events for each tenant.
2. Publish each event to the in-process event bus (for local handlers).
3. On success: mark as ``published``.
4. On failure: increment attempts, schedule retry with exponential
backoff, or mark as ``failed`` if max attempts exceeded.
Args: Args:
db: Async SQLAlchemy session for this batch. db: Async SQLAlchemy session for this batch (crm_worker role).
redis: Optional Redis client (unused for now, reserved for future redis: Optional Redis client (unused for now, reserved for future
cross-process pub/sub). cross-process pub/sub).
batch_size: Maximum events to process in one batch. batch_size: Maximum events to process per tenant in one batch.
tenant_ids: Optional list of tenant IDs to process. If None,
all tenants are loaded from the database.
Returns: Returns:
Number of events successfully published. Number of events successfully published.
@@ -167,106 +265,32 @@ async def process_outbox_batch(
event_bus = get_event_bus() event_bus = get_event_bus()
published_count = 0 published_count = 0
# Claim a batch of pending events # Load tenant IDs if not provided
rows = ( if tenant_ids is None:
await db.execute(_CLAIM_SQL, {"batch_size": batch_size}) result = await db.execute(text("SELECT id FROM tenants"))
).fetchall() tenant_ids = [row[0] for row in result]
if not rows: for tenant_id in tenant_ids:
return 0 # Set tenant context for RLS — required for event_outbox and consumer_inbox
await db.execute(
text("SELECT set_config('app.current_tenant_id', :tid, true)"),
{"tid": str(tenant_id)},
)
for row in rows: # Claim a batch of pending events for this tenant
event_id = row[0] rows = (
tenant_id = row[1] await db.execute(_CLAIM_SQL, {"batch_size": batch_size})
event_name = row[2] ).fetchall()
payload = row[3]
attempts = row[4]
max_attempts = row[5]
aggregate_type = row[6] if len(row) > 6 else None
aggregate_id = row[7] if len(row) > 7 else None
occurred_at = row[8] if len(row) > 8 else None
correlation_id = row[9] if len(row) > 9 else None
schema_version = row[10] if len(row) > 10 else 1
# payload comes back as a dict from JSONB if not rows:
if isinstance(payload, str): continue
import json
payload_dict = json.loads(payload)
else:
payload_dict = payload
try: for row in rows:
# Enrich payload with standardized event envelope metadata published = await _process_single_outbox_event(db, event_bus, row)
payload_dict.setdefault("_event_id", str(event_id)) if published:
payload_dict.setdefault("_event_name", event_name)
payload_dict.setdefault("_event_timestamp", datetime.now(timezone.utc).isoformat())
payload_dict.setdefault("_tenant_id", str(tenant_id))
payload_dict.setdefault("_aggregate_type", aggregate_type)
payload_dict.setdefault("_aggregate_id", str(aggregate_id) if aggregate_id else None)
payload_dict.setdefault("_occurred_at", occurred_at.isoformat() if occurred_at else None)
payload_dict.setdefault("_correlation_id", str(correlation_id) if correlation_id else None)
payload_dict.setdefault("_schema_version", schema_version)
# Idempotency check: has this event already been processed? (P1.5 fix)
already_processed = await db.execute(
text("SELECT 1 FROM consumer_inbox WHERE event_id = :eid AND status = 'processed' LIMIT 1"),
{"eid": str(event_id)},
)
if already_processed.first():
# Event was already processed by all consumers — mark as published
await db.execute(_MARK_PUBLISHED_SQL, {"id": str(event_id)})
published_count += 1 published_count += 1
logger.debug("Outbox event %s already processed, marking as published", event_id)
continue
results = await event_bus.publish_with_results(event_name, payload_dict)
# Check if any handlers were registered at all
handler_count = len(results)
# If any handler raised, treat as failure
handler_errors = [r for r in results if r is not None]
if handler_errors:
raise handler_errors[0]
if handler_count == 0:
# No handlers registered — mark as 'no_handlers' not 'published'
await db.execute(
text("UPDATE event_outbox SET status = 'no_handlers', published_at = now() WHERE id = :id"),
{"id": str(event_id)},
)
logger.warning("Outbox event %s (%s) had no handlers registered", event_id, event_name)
else:
# Record in consumer_inbox for idempotency (P1.5 fix)
await db.execute(
text("INSERT INTO consumer_inbox (event_id, consumer_name, status, processed_at) VALUES (:eid, :name, 'processed', now()) ON CONFLICT DO NOTHING"),
{"eid": str(event_id), "name": event_name},
)
await db.execute(_MARK_PUBLISHED_SQL, {"id": str(event_id)})
published_count += 1
except Exception as exc:
logger.error(
"Failed to publish outbox event %s (%s): %s",
event_id, event_name, exc,
exc_info=True,
)
new_attempts = attempts + 1
if new_attempts >= max_attempts:
await db.execute(_FAIL_SQL, {"id": str(event_id)})
logger.warning(
"Outbox event %s marked as failed after %d attempts",
event_id, new_attempts,
)
else:
backoff = timedelta(seconds=(2 ** new_attempts) * 10)
next_retry = datetime.now(timezone.utc) + backoff
await db.execute(
_RETRY_SQL,
{
"id": str(event_id),
"attempts": new_attempts,
"next_retry_at": next_retry,
},
)
await db.commit() # Commit after each tenant to release locks
await db.commit()
return published_count return published_count
+32 -11
View File
@@ -103,7 +103,6 @@ async def on_startup(ctx: dict[str, Any]) -> None:
# Initialize plugin registry and discover built-in plugins # Initialize plugin registry and discover built-in plugins
from app.plugins.registry import get_registry from app.plugins.registry import get_registry
from app.core.db import get_engine
from app.core.event_bus import get_event_bus from app.core.event_bus import get_event_bus
from app.core.webhook_dispatcher import register_webhook_event_handlers from app.core.webhook_dispatcher import register_webhook_event_handlers
from sqlalchemy import select as sa_select from sqlalchemy import select as sa_select
@@ -111,13 +110,14 @@ async def on_startup(ctx: dict[str, Any]) -> None:
from sqlalchemy.ext.asyncio import async_sessionmaker from sqlalchemy.ext.asyncio import async_sessionmaker
registry = get_registry() registry = get_registry()
from app.core.db import get_worker_engine from app.core.db import get_migration_engine
worker_engine = get_worker_engine() migration_engine = get_migration_engine()
registry.initialize(worker_engine, app=None) registry.initialize(migration_engine, app=None)
registry.discover_builtins() registry.discover_builtins()
event_bus = get_event_bus() event_bus = get_event_bus()
async_session = async_sessionmaker(worker_engine, expire_on_commit=False) 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) # Activate plugins that are marked active in DB (register event handlers)
# RLS fail-closed requires tenant context for tenant-table writes. # RLS fail-closed requires tenant context for tenant-table writes.
@@ -133,11 +133,25 @@ async def on_startup(ctx: dict[str, Any]) -> None:
all_tenant_ids = [row[0] for row in tenant_result] all_tenant_ids = [row[0] for row in tenant_result]
logger.info(f"Worker: loaded {len(all_tenant_ids)} tenants") logger.info(f"Worker: loaded {len(all_tenant_ids)} tenants")
# Register event handlers only (no DB writes, no cron job registration) # 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.is_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(): for name in registry.resolve_load_order():
plugin = registry.get_plugin(name) plugin = registry.get_plugin(name)
if plugin is None: if plugin is None:
continue 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: try:
# Just register event handlers, skip DB-writing on_activate # Just register event handlers, skip DB-writing on_activate
if hasattr(plugin, 'register_event_handlers'): if hasattr(plugin, 'register_event_handlers'):
@@ -206,14 +220,21 @@ async def process_outbox_job(ctx: dict[str, Any]) -> None:
Uses a distributed Redis lock so only one worker replica processes the Uses a distributed Redis lock so only one worker replica processes the
outbox at a time. Runs every 5 seconds. outbox at a time. Runs every 5 seconds.
"""
from app.core.db import get_session_factory
from app.core.outbox import process_outbox_batch
factory = get_session_factory() 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: async with factory() as db:
try: try:
count = await process_outbox_batch(db, batch_size=50) # 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: if count:
logger.info("Outbox: published %d events", count) logger.info("Outbox: published %d events", count)
except Exception as exc: except Exception as exc: