Skip to content

Event bus

The event bus is a small in-process publish/subscribe layer (msflib.eventbus). Modules use it to react to each other without importing each other: the workspaces module, for example, listens for account creation and sets up a default workspace, and the account module never knows. It is also how a host app hooks its own behavior into module lifecycles.

It is synchronous and in-process. It is not a message queue: nothing is persisted, retried or delivered to other processes. For background work see long-running jobs.

Concepts

  • An event is a string name (or a BaseEnum/Enum member, normalized to its value) such as "order-placed" or "account-create-pre-commit".
  • A listener is a function called with whatever positional and keyword arguments the emitter passed. By convention module events pass (instance, options), where options is a dict that may carry the DB session.
  • An Emitter holds the listeners. Each FastAPI app gets its own, wrapped in an AppEmitter, so apps and tests do not share listeners.
  • The module-level emitter is a proxy to the emitter active in the current context. Module code emits through it; the middleware installed by bind_app_emitter makes it point at your app's emitter during each request.

Wiring it into your app

Bind one emitter per app at creation, then pass the AppEmitter to module hook registrations and register your own listeners.

from fastapi import FastAPI
from msflib.eventbus import bind_app_emitter

app = FastAPI()
app_emitter = bind_app_emitter(app)


@app_emitter.on("order-placed")
def on_order_placed(order, options):
    print("order", order["id"], "placed")

bind_app_emitter(app) stores the emitter on app.state, adds a middleware that activates it for each HTTP and WebSocket request, and is safe to call again (get_app_emitter(app) returns the existing one). @app_emitter.on(...) registers the function as a listener on that emitter.

The app emitter is active for requests, for FastAPI BackgroundTasks, and for call_later and create_task calls started from a request. Code in your web process that runs outside the request context (scripts, threads, lifespan code, scheduler jobs) uses the default emitter instead; wrap it in with use_app_emitter(app): (from msflib.eventbus) to make the app emitter current. A worker process has no FastAPI app: register its listeners on the worker's emitter, AppEmitter(get_emitter()) (both from msflib.eventbus).

Module hooks are registered by calling the module's register_event_hooks, at import time or in a lifespan handler:

from msflib.workspaces.eventbus import register_event_hooks

register_event_hooks(
    workspace_action=workspace_action,
    user_action=user_action,
    account_action=account_action,
    emitter=app_emitter,
    tenancy_settings=settings.scope("TENANCY"),
)

Modules that ship hooks, with their eventbus packages, include account, auth, conversation, documents (msflib.documents.drivelink), knowledge and workspaces. Each exports its event names and, where defined, an EVENT_CONTRACTS map from event name to the options type. The core msflib.eventbus also exports EventName (the generic model lifecycle names), EventOptions and merge_event_contracts.

A TestClient that is not used as a context manager skips the lifespan, so hooks registered there are absent in such tests; use with TestClient(app) or register the hooks at import time.

The testsite in the repository (testsite/app/main.py) shows how the documents and knowledge hooks are registered in a lifespan handler.

Emitting

Choose the method by what should happen when nobody listens or a listener fails.

Method No listener A listener raises
emit_optional returns False logged as a warning, remaining listeners still run, returns False
emit_required raises EventBusRequiredListenerError raises EventBusListenerExecutionError (original error as __cause__)
emit_strict returns False re-raises the original exception (or returns it with raise_error=False)

emit_optional returns True only when there was at least one listener and all of them succeeded. Use it for notifications that must not break the caller. Use emit_required when a listener's work is part of the operation, for instance work that has to happen inside the same transaction.

Each has an _async variant (emit_optional_async, emit_required_async, emit_strict_async) that awaits async listeners. The synchronous emit_required and emit_strict reject async listeners with an error, and emit_optional schedules coroutine listeners on the running loop instead of waiting for them. Keep listeners for transactional events synchronous. emit_if_listeners still exists but is deprecated in favor of emit_optional.

from msflib.eventbus import emitter

delivered = emitter.emit_optional("order-placed", {"id": 1}, {})

Model lifecycle events

Every ModelAction create, update and delete emits events automatically, so you can hook any model without subclassing its action. For a model named Customer:

  • customer-create-pre-commit, customer-update-pre-commit, customer-delete-pre-commit, and the generic model-create-pre-commit and so on. Pre-commit events are emitted with emit_required after the row is flushed and before commit, so a listener runs in the same transaction and a listener error aborts the operation. They are only emitted when something is listening. The listener receives (model, payload) where payload includes the session.
  • The matching -post-commit events fire after the outermost commit with emit_optional, so listener errors are logged and do not fail the request. Post-commit listeners get a plain snapshot of the row rather than the live ORM object, which may be expired by then.

Model names are the lowercase class name, so a listener for the pre-commit event of Customer is registered on "customer-create-pre-commit".

Pitfalls

  • Register listeners on the app emitter with @app_emitter.on(...).
  • emitter.on(...) on the module-level emitter object registers on whichever emitter is current when it runs. At import time or in a lifespan handler that is the default emitter, which an app bound with bind_app_emitter(app) never sees, so a required event then fails with EventBusRequiredListenerError. Inside with use_app_emitter(app): or a request it is the app emitter. Prefer @app_emitter.on.
  • FastAPI BackgroundTasks still run inside the request context, so the module-level emitter is your app's emitter there. Code in your web process that runs outside the request context (a threading.Thread, a scheduler job, a script, a lifespan task) sees the default emitter; wrap it in with use_app_emitter(app):. A worker process has no app, so register its listeners on AppEmitter(get_emitter()).
  • Calling a module's register_event_hooks twice on the same emitter would double-register. The modules guard this with HookRegistry: register each listener with the emitter first, then call registry.bind(emitter=..., listeners=[(event, listener), ...]) with those functions, so a second call replaces the earlier bindings rather than duplicating them. bind itself does not register listeners. Use reset in tests.
  • Pre-commit listeners that raise abort the operation. Do slow or external work (email, HTTP calls) in post-commit listeners or a job queue.
  • EventBus logging is silent by default (a NullHandler). Call configure_eventbus_logger(handler=..., level=...) to see warnings from failed optional listeners.