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:
+128
-104
@@ -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(
|
||||
db: AsyncSession,
|
||||
redis: aioredis.Redis | None = None,
|
||||
batch_size: int = 50,
|
||||
tenant_ids: list[uuid.UUID] | None = None,
|
||||
) -> int:
|
||||
"""Process one batch of pending outbox events.
|
||||
|
||||
1. Claim up to *batch_size* pending events using ``FOR UPDATE SKIP LOCKED``
|
||||
so multiple workers don't interfere.
|
||||
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.
|
||||
Iterates over all tenants, setting tenant context for RLS before
|
||||
claiming and processing events for each tenant.
|
||||
|
||||
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
|
||||
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:
|
||||
Number of events successfully published.
|
||||
@@ -167,106 +265,32 @@ async def process_outbox_batch(
|
||||
event_bus = get_event_bus()
|
||||
published_count = 0
|
||||
|
||||
# Claim a batch of pending events
|
||||
rows = (
|
||||
await db.execute(_CLAIM_SQL, {"batch_size": batch_size})
|
||||
).fetchall()
|
||||
# Load tenant IDs if not provided
|
||||
if tenant_ids is None:
|
||||
result = await db.execute(text("SELECT id FROM tenants"))
|
||||
tenant_ids = [row[0] for row in result]
|
||||
|
||||
if not rows:
|
||||
return 0
|
||||
for tenant_id in tenant_ids:
|
||||
# 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:
|
||||
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
|
||||
# Claim a batch of pending events for this tenant
|
||||
rows = (
|
||||
await db.execute(_CLAIM_SQL, {"batch_size": batch_size})
|
||||
).fetchall()
|
||||
|
||||
# payload comes back as a dict from JSONB
|
||||
if isinstance(payload, str):
|
||||
import json
|
||||
payload_dict = json.loads(payload)
|
||||
else:
|
||||
payload_dict = payload
|
||||
if not rows:
|
||||
continue
|
||||
|
||||
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)})
|
||||
for row in rows:
|
||||
published = await _process_single_outbox_event(db, event_bus, row)
|
||||
if published:
|
||||
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
|
||||
|
||||
+32
-11
@@ -103,7 +103,6 @@ async def on_startup(ctx: dict[str, Any]) -> None:
|
||||
|
||||
# Initialize plugin registry and discover built-in plugins
|
||||
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.webhook_dispatcher import register_webhook_event_handlers
|
||||
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
|
||||
|
||||
registry = get_registry()
|
||||
from app.core.db import get_worker_engine
|
||||
worker_engine = get_worker_engine()
|
||||
registry.initialize(worker_engine, app=None)
|
||||
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()
|
||||
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)
|
||||
# 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]
|
||||
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():
|
||||
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'):
|
||||
@@ -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
|
||||
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:
|
||||
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:
|
||||
logger.info("Outbox: published %d events", count)
|
||||
except Exception as exc:
|
||||
|
||||
Reference in New Issue
Block a user