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(
|
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
@@ -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:
|
||||||
|
|||||||
Reference in New Issue
Block a user