gate: worker skips plugin activation, only registers event handlers
Check Cross-Plugin Imports / check (push) Has been cancelled
Check Cross-Plugin Imports / check (push) Has been cancelled
This commit is contained in:
+18
-32
@@ -120,45 +120,31 @@ async def on_startup(ctx: dict[str, Any]) -> None:
|
|||||||
async_session = async_sessionmaker(worker_engine, expire_on_commit=False)
|
async_session = async_sessionmaker(worker_engine, expire_on_commit=False)
|
||||||
|
|
||||||
# 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.
|
||||||
|
# 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.models.tenant import Tenant as TenantModel
|
||||||
from app.core.db import set_tenant_context
|
from app.core.db import set_tenant_context
|
||||||
|
|
||||||
async with async_session() as db:
|
async with async_session() as db:
|
||||||
# Load all tenant IDs for per-tenant plugin activation
|
# Load all tenant IDs for per-tenant event handler registration
|
||||||
tenant_result = await db.execute(sa_select(TenantModel.id))
|
tenant_result = await db.execute(sa_select(TenantModel.id))
|
||||||
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 for plugin activation")
|
logger.info(f"Worker: loaded {len(all_tenant_ids)} tenants")
|
||||||
|
|
||||||
async with async_session() as db:
|
# Register event handlers only (no DB writes, no cron job registration)
|
||||||
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
|
||||||
result = await db.execute(
|
try:
|
||||||
sa_select(PluginModel).where(PluginModel.name == name)
|
# Just register event handlers, skip DB-writing on_activate
|
||||||
)
|
if hasattr(plugin, 'register_event_handlers'):
|
||||||
plugin_record = result.scalar_one_or_none()
|
await plugin.register_event_handlers(event_bus)
|
||||||
if plugin_record is None or not plugin_record.active:
|
logger.info(f"Worker: registered event handlers for {name}")
|
||||||
continue
|
except Exception as exc:
|
||||||
# Activate plugin per-tenant with tenant context set
|
logger.warning(f"Worker: failed to register event handlers for {name}: {exc}")
|
||||||
activation_failed = False
|
|
||||||
for tenant_id in all_tenant_ids:
|
|
||||||
try:
|
|
||||||
await set_tenant_context(db, tenant_id)
|
|
||||||
await plugin.on_activate(db, container, event_bus)
|
|
||||||
# Flush to detect any RLS errors that were swallowed by the plugin
|
|
||||||
await db.flush()
|
|
||||||
except Exception as exc:
|
|
||||||
# RLS may block duplicate cron job inserts — rollback and continue
|
|
||||||
# The cron jobs are already registered from the API container startup
|
|
||||||
logger.warning(f"Worker: plugin {name} activation issue for tenant {tenant_id}: {exc}")
|
|
||||||
await db.rollback()
|
|
||||||
activation_failed = True
|
|
||||||
break
|
|
||||||
if not activation_failed:
|
|
||||||
logger.info(f"Worker: activated plugin {name}")
|
|
||||||
await db.commit()
|
|
||||||
|
|
||||||
# Register webhook dispatcher on the event bus
|
# Register webhook dispatcher on the event bus
|
||||||
register_webhook_event_handlers(event_bus)
|
register_webhook_event_handlers(event_bus)
|
||||||
|
|||||||
Reference in New Issue
Block a user