Skip to content

Events — one fact, five verbs

A mutation in the record world announces itself as one business object event. There is one shape and five verbs — created, updated, deleted, archived, restored — one per edge of the status machine.

Two mutations announce nothing, and a listener must not assume otherwise:

  • The draft world is silent. DraftWorkspace.create, update and discard emit nothing. activate is the moment a record first speaks, and it emits created — see Lifecycle.
  • A no-op update. update diffs the record before and after; an empty diff emits nothing. “Someone called update” is not an event.

BusinessObjectEvent (@onerp/meta) is a discriminated union keyed on operation. Every verb carries record — the state the transition left behind, or for deleted the last state there was:

operation payload
created record — the born state
updated previous, record, and the changes diff between them
deleted record — the final state
archived, restored record — the state after the move

A listener acts on a transition, so the event hands it whole state: re-reading the record it was just given is a query it should never have to make. updated also carries previous, which is what lets a listener reconcile against prior intent instead of re-deriving it from whatever it wrote last time — ReservationMovementListener posts the difference between the hold the event found and the hold it left, and never reads the ledger to work out the first half.

changes is the observed Changes value computed from previous and record at the emit site: a cached derivation, not requested Update intent and not a second truth. It is the cheap way to ask did this field move ("baseUnit" in event.changes) without diffing two records yourself.

archived and restored carry no previous: the verb already names the status it came from, and the rest of the record is untouched by a lifecycle move.

Every event also carries the business object name, the record id, and a cause (below). The union is generic over the schema, so a listener subscribed to a single verb gets a fully-typed payload.

An audit row is not this shape. The two used to be one type; they are deliberately separate now, because an event is an in-flight signal that wants whole state, while an audit row is history that stores a diff. Welding them let the audit table’s storage budget decide what a listener was allowed to know.

Events dispatch on the topic business-object.{namespace}.{name}.{operation} — for example business-object.master.Product.created. The name is derived from the business object and operation (businessObjectEventName), not stored on the payload. Subscribe with NestJS’s @OnEvent:

@OnEvent("business-object.master.Product.updated", { suppressErrors: false })
onProductUpdated(event: BusinessObjectUpdatedEvent<typeof ProductSchema>) { … }

suppressErrors: false makes the listener part of the write: a throw rolls the whole unit of work back — the right default for an invariant a listener enforces, which should veto the write rather than fail quietly beside it. The emitter runs with wildcard: true, so a cross-cutting listener subscribes to every verb of every business object with one segment pattern — @OnEvent("business-object.**"), which is how AuditService listens.

Dispatch — before commit, before the transaction

Section titled “Dispatch — before commit, before the transaction”

BusinessObjectManager.create awaits EventBus.emit, which runs the listeners one at a time in registration order and awaits each, before it returns. One at a time is deliberate: two listeners staging writes to the same record would otherwise interleave their read and their save, and the last save would win. A listener throw propagates out of the mutation and rolls the unit of work back — that is what lets a listener veto a write, and why cross-object invariants live here: a listener gets full dependency injection where a determination gets none. A listener that keeps a field on another business object is a rollup, and gets its listener written for it.

No SQL transaction is open while listeners run. The mutation stages its writes in the per-unit-of-work repository buffer; the transaction opens later, in UnitOfWork.commit(), which flushes the staged repos. So at listener time:

a listener reading through sees
BusinessObjectManager, or uow.repo(Schema) the mutation — repo.get checks the staging buffer before the database
raw Kysely on uowStore.current().db the pre-mutation database — the row is not written yet

product-locked-fields.listener.ts is the live example of the second row, and it is correct there: it queries inventory_ledger, a table no Product mutation touches. Atomicity is unaffected — only when rows exist changes.

A listener that needs the actual Transaction registers a callback with uow.onCommit(cb). It runs inside the commit transaction, after the staged repos flush, and a throw from it aborts the transaction with everything else in it — this is how the stock ledger gets written. GoodsMovementListener turns a movement’s lines into postings and hands the batch to PostingEngine, which registers one callback for all of them:

async stage(entries: LedgerPosting[]): Promise<void> {
this.uow.current().onCommit((tx) => this.post(tx, entries));
}

Two siblings sit beside it. onAfterCommit(cb) runs after the transaction commits — for a side effect that must only happen if the write is durable (deleting an object from storage, publishing an outbox row); its errors are caught and logged, because the outcome is already sealed and a throw would shadow the real result. onRollback(cb) undoes an eager side effect when the unit of work fails.

Cause — attributing the immediate trigger

Section titled “Cause — attributing the immediate trigger”

cause names what triggered this mutation, one level up — not the originator of a long cascade. The bus stamps it from an AsyncLocalStorage frame that is replaced at every nesting level: in an A→B→C cascade, C’s cause is B. The correlationId is what groups a whole cascade — every audit row written in one unit of work shares it.

There are four kinds, plus null:

cause stamped around
{ kind: "event", operation, businessObject, businessObjectId } each in-process dispatch — so a cascade listener’s writes point back at the event that woke it
{ kind: "action", action } every action invocation, root and bound
{ kind: "job", queue, jobId } a background job body
{ kind: "conversation", conversationId } each tool call an agent makes — the thread the write came from; the actor says who
null a write with no frame above it — a plain HTTP create or update

An HTTP request that runs an action therefore does not produce cause: null; the action frame is above the write.

The same frame carries the cascade depth, and EventBus refuses to emit at MAX_DEPTH (10), so a listener that re-triggers its own event throws instead of spinning. Depth counts any cause frame — a write started by an action or a job is already at depth 1 — and because it lives in the per-async-context frame it measures true chain depth, never in-flight events across concurrent requests.

Alongside the in-process dispatch, EventBus.emit registers one onAfterCommit callback per event on the ambient unit of work. When (and only when) the transaction commits, each registered outbound handler receives the event; a rollback delivers nothing, because the unit drops its callbacks. The two paths are separate on purpose: in-process listeners are transactional and can veto, outbound delivery is post-commit and best-effort. There is no separate “integration event” shape — a transport decides its own payload at the bridge.

Two transports sit on the bridge: the workflow engine, below, and broadcast, which tells open browser tabs which record changed so they refetch it. Broadcast is volatile by design; it is not a way to deliver events to an integration.

A workflow is driven by the outbound bridge, not the in-process bus. WorkflowExplorer registers an outbound handler at boot and forwards each committed event to the workflow engine under an app-id-prefixed topic — onerp/business-object.sales.SalesOrder.created for the app configured as WorkflowModule.forRoot({ id: "onerp" }). Three consequences are worth knowing before choosing a workflow over a listener:

  • A workflow runs after the commit and cannot veto the write. An invariant that must be able to refuse belongs in a listener or a validator.
  • Forwarding is fire-and-forget: a failure is logged at warn and nothing retries it. The durable outbox will own retries when it lands.
  • With testing: true the explorer skips its wiring entirely.

businessObjectEvent(Schema, "created") builds the trigger a workflow subscribes with. A workflow is a decorated provider (see registration): the triggers go into the Workflow.on base, and run’s event type derives from them — the union, for a multi-trigger workflow — so triggering on archived while reading a created payload is a compile error, not a drift kept in sync by eye. The engine hands the handler its own envelope, so the business object event is event.data. Work that must survive a retry goes in step.run, which the explorer wraps in a tenant frame, a system actor and a unit of work of its own:

export class DemoOnOrderCreated extends Workflow.on(
[businessObjectEvent(SalesOrderSchema, "created")],
{ id: "sales/demo-on-order-created" }
) {
async run({ event, step }: WorkflowContext<
BusinessObjectCreatedEvent<typeof SalesOrderSchema>
>): Promise<void> {
const order = await step.run("fetch-order", () =>
this.bom.get(SalesOrderSchema, event.data.businessObjectId)
);
}
}

The event also carries the whole record, so a workflow that only hands it on reads what it was given rather than fetching it again. Starting an agent never goes in a step. AgentService.run opens a unit of work per tool call and one more to save the thread, and a step’s unit would nest and throw; the call sits directly in the workflow’s run. The agent acts as its own principal, not as the workflow’s system actor, and a retry mints a fresh thread:

export class OrderIntake extends Workflow.on(
[businessObjectEvent(SalesOrderSchema, "created")],
{ id: "sales/order-intake" }
) {
async run({ event }: WorkflowContext<
BusinessObjectCreatedEvent<typeof SalesOrderSchema>
>): Promise<void> {
const { documentNumber, id } = event.data.record;
await this.agent.run({
agent: orderIntake,
message: `Sales order ${documentNumber} (id ${id}) was just created. Review it.`,
});
}
}