Skip to content

Outbox ​

An outbox stores the integration events of a change in the same transaction as the change itself. A relay then reads them and hands them to an event publisher, retrying until it succeeds: an event is never lost, even if the broker is down when the change is saved.

ts
await this.outbox.add(order.pullDomainEvents().map((event) => this.translator.translate(event, { correlationId })));

await new OutboxRelay(outbox, publisher).relay();
one unit of workcommand handlerplace the orderorderssave(order)outboxadd(events)OutboxRelayrelay()pendingEventPublisherpublish(events)other contextsbroker, consumers

When to use ​

Whenever an event must not be lost: other bounded contexts, other systems or an audit trail depend on it. Add the events in the unit of work of the change, and run the relay in the background.

Usage ​

Add the events of a change ​

In the command handler, inside the unit of work, after saving the aggregate, translated by an event translator.

src/ordering/application/commands/place-order.command.ts
ts
return this.unitOfWork.run(async () => {
	await this.orders.save(order);
	await this.outbox.add(order.pullDomainEvents().map((event) => this.translator.translate(event, { correlationId: orderId })));
	return ok();
});

Relay the events ​

OutboxRelay reads a batch of pending events, publishes them and marks them as published. It returns how many it published; if publishing fails, it throws and the events stay pending for the next call. Build it in the composition root with a batch size, and call it on a schedule from a driving adapter: a timer, a cron job, a worker.

src/ordering/ordering.module.ts
ts
const relay = new OutboxRelay(outbox, publisher, 100);
src/ordering/driving/timer/jobs/outbox-relay.job.ts
ts
import type { OutboxRelay } from "@alveolus/core";

export class OutboxRelayJob {
	constructor(private readonly relay: OutboxRelay) {}

	start(): NodeJS.Timeout {
		return setInterval(() => this.tick(), 1000);
	}

	private async tick(): Promise<void> {
		try {
			await this.relay.relay();
		} catch (error) {
			console.error("The outbox relay failed; the events stay pending for the next tick.", error);
		}
	}
}

Implement it ​

Extend Outbox in a driven adapter, typically over a table written through the transaction of the unit of work. Integration events are JSON: store them as they are, for instance in a jsonb column, and return them unchanged from pending, oldest first.

src/shared-kernel/driven/pg/adapters/pg-outbox.adapter.ts
ts
import { type AnyIntegrationEvent, Outbox } from "@alveolus/core";
import type { Pool } from "pg";

export class PgOutbox extends Outbox {
	constructor(private readonly db: Pool) {
		super();
	}

	async add(events: readonly AnyIntegrationEvent[]): Promise<void> {
		for (const event of events) {
			await this.db.query("INSERT INTO outbox (id, event) VALUES ($1, $2)", [event.id, event]);
		}
	}

	async pending(limit: number): Promise<readonly AnyIntegrationEvent[]> {
		const { rows } = await this.db.query("SELECT event FROM outbox WHERE published_at IS NULL ORDER BY created_at LIMIT $1", [limit]);
		return rows.map((row) => row.event);
	}

	async markPublished(ids: readonly string[]): Promise<void> {
		await this.db.query("UPDATE outbox SET published_at = now() WHERE id = ANY($1)", [ids]);
	}
}

Do not publish directly ​

❌ Avoid
ts
await this.orders.save(order);
await this.publisher.publish(events);
✅ Prefer
ts
await this.unitOfWork.run(async () => {
	await this.orders.save(order);
	await this.outbox.add(events);
	return ok();
});
Why?

If the process stops between save and publish, the order is placed and nobody hears about it. If the broker is down, the command fails although the order is valid. Writing the events in the same transaction as the change, then relaying them, removes both cases.

Reference ​

ts
abstract class Outbox extends Port {
	abstract add(events: readonly AnyIntegrationEvent[]): Promise<void>;
	abstract pending(limit: number): Promise<readonly AnyIntegrationEvent[]>;
	abstract markPublished(ids: readonly string[]): Promise<void>;
}

class OutboxRelay {
	constructor(outbox: Outbox, publisher: EventPublisher, batchSize?: number);
	relay(): Promise<number>;
}
MemberTypeDescription
add(events)Promise<void>Stores the events, in the current transaction.
pending(limit)Promise<readonly AnyIntegrationEvent[]>Returns up to limit unpublished events, oldest first.
markPublished(ids)Promise<void>Marks the events as published.
new OutboxRelay(outbox, publisher, batchSize = 100)Relays from the outbox to the publisher.
relay()Promise<number>Publishes one batch and returns how many events it published; 0 when none is pending.

Caveats

  • Delivery is at least once: if the relay stops between publishing and marking, the events are published again. Consumers ignore duplicates by id.
  • Several relays running at once may publish the same batch; lock the rows in pending (for instance FOR UPDATE SKIP LOCKED) if you run more than one.

Import from @alveolus/core or @alveolus/core/outbox.

See also ​

Released under the MIT License.