Outbox
Events written in the same transaction as the change, read by named consumers with their own cursor and relayed as CloudEvents.
The outbox SQL module 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).
pnpm better-supabase sql add outboxOnce 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
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:
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:
select better_supabase.track_events('public.invoices', 'invoice', 'organization_id');Relaying events
Create the outbox client on a server connection, register each consumer once, then relay it from a cron route or a worker:
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"],
});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.
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
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.
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
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
To process events with retries, relay them into a job queue with the event id as the dedupe key:
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
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
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.
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
hint | When |
|---|---|
OUTBOX_TYPE_REQUIRED | emit_event without a type |
OUTBOX_UNKNOWN_CONSUMER | Claiming for a consumer that isn't registered |
Last updated on
Jobs, idempotency and webhook inbox
Background jobs on Supabase Queues, safe retries and exactly-once webhook handling on the Postgres you already have.
Workflows
A run registry for any workflow engine that members read through RLS, cron schedules, counting semaphores, admission control for starts, and useWorkflowRuns over Realtime.