---
url: /core/application/outbox.md
---
# Outbox

An outbox stores the [integration events](./integration-events.md) of a change in the same
transaction as the change itself. A relay then reads them and hands them to an
[event publisher](./event-publishers.md), 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();
```

## 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](./unit-of-work.md) of the change, and run the relay in
the background.

## Usage

### Add the events of a change

In the [command handler](./command-handlers.md), inside the unit of work, after saving the
aggregate, translated by an [event translator](./event-translators.md).

```ts [src/ordering/application/commands/place-order.command.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.

```ts [src/ordering/ordering.module.ts]
const relay = new OutboxRelay(outbox, publisher, 100);
```

```ts [src/ordering/driving/timer/jobs/outbox-relay.job.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.

```ts [src/shared-kernel/driven/pg/adapters/pg-outbox.adapter.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

```ts [❌ Avoid]
await this.orders.save(order);
await this.publisher.publish(events);
```

```ts [✅ Prefer]
await this.unitOfWork.run(async () => {
	await this.orders.save(order);
	await this.outbox.add(events);
	return ok();
});
```

::: details 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>;
}
```

| Member | Type | Description |
| --- | --- | --- |
| `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

* [Unit of Work](./unit-of-work.md), the transaction the outbox is written in
* [Event translators](./event-translators.md), which build what the outbox stores
* [Event publishers](./event-publishers.md), what the relay calls
