"""In-process event bus for publish/subscribe.""" from __future__ import annotations import asyncio from collections import defaultdict from collections.abc import Callable, Coroutine from typing import Any EventHandler = Callable[[dict[str, Any]], Coroutine[Any, Any, None]] class EventBus: """Simple async event bus for in-process pub/sub.""" def __init__(self) -> None: self._handlers: dict[str, list[EventHandler]] = defaultdict(list) def subscribe(self, event_name: str, handler: EventHandler) -> None: """Subscribe a handler to an event.""" self._handlers[event_name].append(handler) def unsubscribe(self, event_name: str, handler: EventHandler) -> None: """Unsubscribe a handler from an event.""" if event_name in self._handlers: self._handlers[event_name] = [h for h in self._handlers[event_name] if h is not handler] async def publish(self, event_name: str, payload: dict[str, Any]) -> None: """Publish an event to all subscribers.""" handlers = self._handlers.get(event_name, []) tasks = [asyncio.create_task(h(payload)) for h in handlers] if tasks: await asyncio.gather(*tasks, return_exceptions=True) # Global event bus instance _event_bus = EventBus() def get_event_bus() -> EventBus: """Get the global event bus.""" return _event_bus def register_workflow_event_handlers() -> None: """Register workflow event handlers on the global event bus. Subscribes to events that can trigger workflows (user.created, etc.). Should be called during application startup. """ from app.workflows.engine import register_workflow_event_handlers as _register _register()