# Outbox

> Events written in the same transaction as the change, read by named consumers with their own cursor and relayed as CloudEvents.

Source: https://bettersupabase.com/docs/blocks/outbox

The `outbox` [SQL module](/docs/blocks/sql) stores events in a table in
the same transaction as the change that caused them. An event exists only
if that transaction commits, so a rolled-back write never notifies anyone
and a committed one is never lost. Named consumers read the events in order,
each with its own cursor, and `createOutbox` from `better-supabase/blocks/outbox`
relays them to any `EventSink` (an HTTP endpoint, a queue, a bus).

```bash
pnpm better-supabase sql add outbox
```

Once the module is installed, the other SQL modules write their events to
it: `organizations` emits `organization.created`, `organization.switched` and the member
events, and `support-sessions` emits `support.started` and
`support.ended`. Set
`sql.modules.<module>.events` to `false` to keep one module out of the outbox.

## Emitting from SQL [#emitting-from-sql]

`emit_event(type, payload, subject, tenant, key, source)` appends an event
and returns its id (a `uuid` on a managed table, as text). Call it from your
own functions and triggers:

```sql
perform better_supabase.emit_event(
  'invoice.paid',
  jsonb_build_object('invoice_id', new.id, 'amount', new.amount),
  subject => 'invoices/' || new.id,
  tenant => new.organization_id::text,
  key => 'invoice-paid-' || new.id
);
```

| Argument  | Meaning                                                                                  |
| --------- | ---------------------------------------------------------------------------------------- |
| `type`    | Any string, usually `noun.verb`. Consumers filter on it                                  |
| `payload` | The event data. `subject` defaults to `payload ->> 'subject'`                            |
| `tenant`  | The organization the event belongs to, relayed as the `partitionkey` extension           |
| `key`     | With a key, emitting the same key again for the same tenant returns the first event's id |
| `source`  | The module or function that wrote it, relayed as the `producer` extension                |

The signed-in user (`auth.uid()`) is stored as the actor. Only
`service_role` can call `emit_event` directly; the module's own functions are
`security definer` and call it for the user. Add roles with
`sql.modules.outbox.options.emitRoles`.

`track_events(table, type_prefix, tenant_column)` adds a row trigger that
emits `<prefix>.created`, `.updated` and `.deleted` with the row as payload
and `<table>/<id>` as subject. An update that changes no column emits
nothing:

```sql
select better_supabase.track_events('public.invoices', 'invoice', 'organization_id');
```

## Relaying events [#relaying-events]

Create the outbox client on a server connection, register each consumer
once, then relay it from a cron route or a worker:

```ts title="lib/outbox.ts"
import { createOutbox } from "better-supabase/blocks/outbox";

export const outbox = createOutbox(postgres.admin, {
  source: "https://crm.example.com",
  typePrefix: "com.example.crm",
});

await outbox.register("billing", {
  types: ["invoice.*", "organization.created"],
});
```

```ts title="app/api/outbox/route.ts"
import "server-only";

import { httpSink } from "better-supabase/events";

import { outbox } from "@/lib/outbox";

export const GET = outbox.relayRoute({
  secret: process.env.CRON_SECRET,
  consumers: {
    billing: httpSink(process.env.BILLING_EVENTS_URL!),
  },
  onError: (error, consumer) => console.error(consumer, error.message),
});
```

`relay(consumer, sink)` claims up to `batch` events (100 by default), sends
them as one `sink.send` call and moves the cursor past them, until none are
left or `budgetMs` runs out. When the sink throws, the cursor stays where it
was and the consumer backs off (2 seconds, doubling up to `maxBackoff`, 10
minutes by default), so delivery is at least once. The CloudEvent `id` is the
row id, so a receiver can drop repeats.

After a failure the consumer claims one event at a time, so one bad event
can't hold back a whole batch. When that single event fails `maxAttempts`
times (an option of `relay` and `consume`, 10 by default), it moves to the consumer's dead letters and the cursor
passes it; the result's `deadLettered` counts them. A success resets the
attempts.

```ts
const dead = await outbox.deadLetters("search").orThrow();
// [{ consumer, event, attempts, error, deadAt }]
```

Dead letters keep the event as it was, so you can fix the cause and emit it
again. `purge` deletes the ones older than its `olderThan`.

Each type gets `typePrefix` (`dev.better-supabase` by default) in front, the
payload is the `data` and the tenant is the `partitionkey`. The actor is not
sent: context attributes must not carry personal data, so put it in the
payload when receivers need it.

A claim leases the consumer to one worker for `lease` (one minute by
default), so two cron runs never send the same batch at the same time.
`relayRoute` checks the bearer secret and relays every consumer at the same
time within one budget, since each has its own cursor and lease.

### Consumers inside the app [#consumers-inside-the-app]

`consume(consumer, handler)` claims and acknowledges the same way, but hands
the handler the outbox rows instead of CloudEvents: `type`, `payload`,
`tenant`, `subject`, `key`, `actorId`, `position` and `createdAt`. Use it for
consumers in your own app, such as search indexing or notifications, that need
the actor. A handler that throws leaves the batch for the next run.

```ts
await outbox.consume("search", async (events) => {
  for (const event of events) {
    await index.update(event.subject, {
      by: event.actorId,
      at: event.createdAt,
    });
  }
});
```

`relayRoute` takes a handler in place of a sink too:
`consumers: { search: (events) => indexEvents(events) }`.

### Ordering and open transactions [#ordering-and-open-transactions]

Positions come from an identity column, which hands out numbers before
commit, so a transaction that started first can commit after a later one.
The claim holds back every event whose transaction id (`xid`) is not older
than the oldest running transaction (`pg_snapshot_xmin`), and orders and
tracks events by `(xid, position)`. An event that commits late therefore
sorts after the consumer's cursor and is delayed, never skipped. A
long-running transaction delays all consumers until it ends.

When none of the claimed events match a consumer's `types`, its cursor
still moves past the settled events, so a consumer of rare types doesn't
rescan the whole table. `outbox.unregister(consumer)` removes a consumer
you no longer relay, so purges stop waiting for it. `purge` with
`{ ignoreIdle: "7 days" }` also stops waiting for a consumer that hasn't
claimed for that long, without removing it.

### Into a queue [#into-a-queue]

To process events with retries, relay them into a [job queue](/docs/blocks/jobs)
with the event id as the dedupe key:

```ts
const toJobs = {
  send: async (events: readonly CloudEvent[]) => {
    for (const event of events) {
      await jobs.enqueue("events", event, { dedupeKey: event.id }).orThrow();
    }
  },
};
await outbox.relay("jobs", toJobs);
```

## History and retention [#history-and-retention]

`outbox.history({ subject: 'invoices/42' })` returns the kept events for one
subject, oldest first, for an activity feed or a debugging view. Filter by
`type` and page with `cursor` (a position) and `limit`.

`outbox.purge(olderThan, batch)` (SQL `purge_outbox(older_than, batch)`)
deletes up to `batch` events (10,000 by default) older than
`sql.modules.outbox.options.retention` (`30 days` by default), but never one that a
registered consumer hasn't passed (see `ignoreIdle` above). It returns how many it deleted; run it
again while that equals `batch`. Schedule it with the jobs schedules or
pg\_cron.

## Existing tables [#existing-tables]

Adopt an events table you already have and rename its columns. Adopt mode
adds the `xid` column (and its index) when the table lacks it, so `sql sync`
and `sql upgrade` don't ask you to declare it first. Map `xid` to
`null` to keep the table as it is: the claim then orders by position alone
and waits `settle` (5 seconds by default, and never zero) after an event's
`created_at`, which skips an event whose transaction commits later than
that. Other columns you don't have map to `null`.

```ts title="better-supabase.config.ts"
export default defineConfig({
  sql: {
    modules: {
      outbox: {
        mode: "adopt",
        tables: { events: "public.domain_events" },
        idType: "uuid",
        columns: {
          events: { type: "kind", position: "seq", xid: null, key: null },
        },
        options: { settle: "2 seconds" },
      },
    },
  },
});
```

`position` must be a `bigint` that only grows. A managed table has a `uuid`
`id` and a separate `position` identity column; map both when you adopt a
table, or map them to the same column when its ids are a `bigint` sequence.

Options marked migration-only match an existing schema. The config accepts
them in `mode: "adopt"` only, and doctor warns about them (BS314) until you
remove them.

| Option          | Default                    | Meaning                                                                    |
| --------------- | -------------------------- | -------------------------------------------------------------------------- |
| `retention`     | `30 days`                  | Default age for `purge_outbox`                                             |
| `emitRoles`     | `["service_role"]`         | Roles allowed to call `emit_event` directly                                |
| `settle`        | `5 seconds`                | The wait before a claim reads an event, without `xid`; never zero          |
| `defaultSource` | none                       | The source `emit_event` stores when the caller passes none; migration-only |
| `blockSource`   | `better-supabase/{module}` | The source of events from SQL modules; migration-only                      |

## Error codes [#error-codes]

| `hint`                    | When                                          |
| ------------------------- | --------------------------------------------- |
| `OUTBOX_TYPE_REQUIRED`    | `emit_event` without a type                   |
| `OUTBOX_UNKNOWN_CONSUMER` | Claiming for a consumer that isn't registered |