Jobs, idempotency and webhook inbox
Background jobs on Supabase Queues, safe retries and exactly-once webhook handling on the Postgres you already have.
better-supabase/blocks/jobs sits on top of the jobs, idempotency and
webhook-inbox SQL modules. It needs a service-role
connection, such as createPostgres(...).admin, because the functions are
closed to anon and authenticated. Every call returns a Result.
better-supabase sql add jobs idempotency webhook-inboxJobs
Jobs run on Supabase Queues
(pgmq) by default, so queues and messages show up in the dashboard. The
jobs module installs pgmq and adds thin better_supabase.* functions for
what pgmq doesn't do itself: lease-safe completion, retries with backoff,
dead letters, deduplication keys and schedules.
Two module options change where jobs and schedules live. Both keep the same functions, so the TypeScript side doesn't change:
| Option | Values | What it does |
|---|---|---|
sql.modules.jobs.options.backend | pgmq (default), table | table stores jobs in better_supabase.job_messages and claims them with for update skip locked, for projects without pgmq |
sql.modules.jobs.options.scheduler | pg_cron (default), drain | drain stores schedules in better_supabase.job_schedules, with a time zone each, and the drain route enqueues them; pg_cron isn't needed |
export default defineConfig({
sql: {
modules: { jobs: { options: { backend: "table", scheduler: "drain" } } },
},
});On the pgmq backend, the module's data file also runs create extension if not exists pgmq, so better-supabase sql data puts the extension in a
migration even when the schema diff leaves it out (pg-delta can, because
pgmq owns its own schema). Doctor warns about an extension no migration
creates (BS321).
Run better-supabase sql sync after changing either option. On the table
backend, completed and dead jobs stay in job_messages with archived_at
set. better_supabase.purge_job_archive(queue, older_than, batch, dead_older_than) deletes completed jobs after older_than (7 days) and
dead letters after dead_older_than (30 days) on both backends, so dead
letters stay around long enough to replay.
Declare queues with a Standard Schema for their payload. enqueue
validates the payload, so a wrong payload is a type error at compile time
and a validation error at runtime. Queue names follow pgmq's rule:
lowercase letters, digits and underscores.
import { createJobs } from "better-supabase/blocks/jobs";
import * as v from "valibot";
export const jobs = createJobs(postgres.admin, {
send_email: v.object({
to: v.pipe(v.string(), v.email()),
template: v.string(),
}),
});
await jobs.enqueue(
"send_email",
{ to, template: "welcome" },
{ delay: 300, dedupeKey: `welcome:${userId}` },
);delay is in seconds. To run at a fixed time, pass runAt as a
Temporal.Instant. Jobs and inbox messages report enqueuedAt,
visibleUntil and receivedAt as Temporal.Instant too (see
Temporal). The third argument takes temporal, for
runtimes without a global Temporal, and now, the clock for runAt delays
and schedule runs. Pass the same now you give defineSupabase so jobs and
repositories share one clock in tests.
Any number of workers can run side by side: pgmq hides a claimed message for the length of its lease.
const controller = new AbortController();
await jobs.work(
"send_email",
async (payload, job, signal) => {
await sendEmail(payload, { signal });
},
{ concurrency: 4, lease: 60, signal: controller.signal },
);- A job is done when its handler returns; it is archived (
pgmq.a_<queue>). If the handler throws or returns a failedResult, the job is retried with exponential backoff and full jitter: a random wait up to 10s, 20s, 40s and so on, at most an hour, so failed jobs don't retry in bursts.last_erroris kept on the message. AftermaxAttempts(5 by default) it is archived withdead: true. A job whose worker died on its last attempt is archived as dead by the next claim instead of running again. jobs.replay(queue, id)(SQLbetter_supabase.replay_dead_job) enqueues a dead letter again with its payload, attempts and dedupe key, and removes it from the archive.- While a batch runs, a heartbeat keeps extending the lease of every job in
it that hasn't finished, including the ones waiting their turn. If a worker
crashes, the lease ends and another worker picks the job up. On pgmq the
lease is extended with
pgmq.set_vt.job.attemptsis pgmq'sread_ct; a worker whose lease was taken over getsfalsefromcompleteinstead of finishing someone else's job. - An idle worker polls after
pollInterval(1 second), then doubles the wait after each empty poll up tomaxPollInterval(30 seconds). The next claimed job resets it. - A claim that fails with a
network,timeout,serializationorrate_limitederror is retried up to 3 times, waitingpollIntervaland then twice as long each time. Any other claim error, or a fourth in a row, stops every lane, andworkordrainrejects with aTypeErrorwhosecauseis theDbError. - A
dedupeKeykeeps one waiting or running job per key. Withdedupe: "waiting"it coalesces only onto a job no worker has claimed yet, so a change made while the job runs queues one follow-up instead of being dropped: a debounce per entity. In SQL this isenqueue_job(..., dedupe_key => key, dedupe_running => false). drain(queue, handler)works the queue until it is empty, which suits cron and Edge Functions.claim,complete,failandextendare available when you want to drive the queue yourself.
Queue health
An admin page reads each queue's state with stats, pages through dead
letters with listDead, and puts them back with retryDead:
const stats = await jobs.stats().orThrow();
// { send_email: { ready, inFlight, delayed, dead, oldestAgeSeconds } }
const dead = await jobs.listDead("send_email", { limit: 50 }).orThrow();
// [{ id, payload, context, attempts, maxAttempts, lastError, enqueuedAt, diedAt }]
const next = await jobs
.listDead("send_email", { limit: 50, cursor: dead.at(-1)?.id })
.orThrow();
await jobs.retryDead("send_email", { ids: [dead[0]!.id] }).orThrow();
await jobs.retryDead("send_email").orThrow(); // the newest 1000| Field | Counts |
|---|---|
ready | messages the next claim takes |
inFlight | messages a worker holds, or that wait out a retry backoff |
delayed | messages enqueued for later that no worker claimed yet |
dead | dead letters in the archive, until purge_job_archive removes them |
oldestAgeSeconds | the age of the oldest message that isn't done, null when none |
stats takes a list of queue names and defaults to every queue
createJobs declared. listDead returns stored payloads without
validating them, since the schema may have changed since. All three call
better_supabase.job_queue_stats, list_dead_jobs and retry_dead_jobs,
so they need a SQL connection.
Actor and tenant
The queue itself runs on a service-role connection, but the work inside a handler should run as the user who asked for it, so your RLS policies still decide what it can read and write.
Pass the request context to enqueue and the job records its actor and
tenant next to the payload. The tenant is context.tenant, the one
tenant() resolved for the connection, or the tenant_id claim (then
app_metadata.tenant_id). The handler gets them back as job.context, and
bs.forContext turns that into
repositories running as that user, in that tenant, over direct Postgres:
await jobs.enqueue("send_invoice", { invoiceId }, { context: db.$context });
await jobs.work("send_invoice", async (payload, job) => {
const db = await bs.forContext(job.context).orThrow();
const invoice = await db.invoices.findById(payload.invoiceId).orThrow();
// ...
});A job enqueued in a support session or an impersonated session also records
the session's act claim. forContext keeps it, so a job enqueued in a
read-only support session runs read-only, as the same support session.
forContext fails with forbidden when the job recorded no user, for
example one enqueued by a cron schedule, and never falls back to the service
role. Work that no user owns uses bs.admin(job.context) on purpose: it
stamps createdBy and applies the tenant() filter, but bypasses RLS, so
keep it visible in code review. admin.$with(job.context) is the same
service-role connection and does not apply RLS either. See
Jobs, webhooks and agents without a session.
A job enqueued without a tenant runs with a context that has none, so
tenant() applies its onMissing (a forbidden error by default). Workers
that handle cross-tenant jobs on purpose pass allTenants: true to work or
drain; jobs that recorded no tenant then run without the tenant scope, and
jobs that did keep theirs. schedule takes the same { context } as its last
argument.
Schedules
schedule enqueues a payload on a cron schedule. Scheduling the same name
again replaces it. A schedule is five cron fields (lists, ranges, steps, and
names such as MON-FRI), a macro (@hourly, @daily, @weekly,
@monthly, @yearly) or an interval (30 seconds, 5 minutes). An
invalid schedule returns an error before anything is written.
await jobs.schedule("nightly-digest", "0 3 * * *", "send_email", {
to: "team@example.com",
template: "digest",
});
await jobs.unschedule("nightly-digest");With the default pg_cron scheduler, enable the
pg_cron extension first. pg_cron
reads every schedule in cron.timezone (UTC on Supabase), so a timeZone
other than UTC returns an error.
With scheduler: "drain", each schedule has its own time zone, and the
drain route enqueues a run when one is due:
await jobs.schedule(
"weekday-digest",
"0 9 * * MON-FRI",
"send_email",
{ to: "team@example.com", template: "digest" },
{ timeZone: "Europe/Amsterdam" },
);The next run is computed in that zone, so 09:00 stays 09:00 across daylight
saving changes. A time that a change skips (02:30 on the night clocks move
forward) runs right after the gap. When the drain route didn't run for a
while, missed runs collapse into one. runSchedules() does the same work
without the route, and nextCronRun(cron, timeZone, after) is exported for
your own checks.
Under the drain scheduler, a schedule also belongs to a tenant: tenant in
the options, or the tenant of context. Each named schedule keeps its own
cron, time zone and tenant, and its runs record the context, so a handler
runs for that tenant. A product page reads schedule state with
listSchedules, and tenant deletion removes them all with unscheduleAll:
await jobs.schedule(
`workflow:${definitionId}`,
"0 8 * * MON",
"run_workflow",
{ definitionId },
{ timeZone: "Europe/Amsterdam", context: { tenant: organizationId } },
);
const schedules = await jobs
.listSchedules({ prefix: "workflow:", tenant: organizationId })
.orThrow();
// [{ name, cron, timeZone, queue, tenant, nextRun, lastRun, leasedUntil, createdAt }]
await jobs.unscheduleAll({ tenant: organizationId }).orThrow();Under pg_cron, listSchedules returns pg_cron's jobs with only the name and
cron, and a schedule with a tenant returns an error.
ensureSchedules keeps a whole set in step with a source of truth, such as
the app's built-in schedules on deploy or a tenant's workflow settings after
an edit. It writes every definition, removes the schedules under prefix
(and tenant, when given) that the set no longer names, and returns both
lists. Every name must start with the prefix, and every payload and cron is
checked before anything is written. A definition whose cron and time zone
didn't change keeps its next run, so a run that is due but not yet drained
still happens and running it twice changes nothing.
const result = await jobs
.ensureSchedules(
workflows.map((workflow) => ({
name: `workflow:${workflow.id}`,
cron: workflow.cron,
queue: "run_workflow",
payload: { definitionId: workflow.id },
timeZone: workflow.timeZone,
})),
{ prefix: "workflow:", tenant: organizationId },
)
.orThrow();
// { scheduled: ["workflow:..."], removed: ["workflow:..."] }Definitions take the same timeZone, tenant and context as schedule;
one without a tenant gets the tenant option.
Schedules from SQL
Under the drain scheduler, a database trigger or function can schedule a job
itself with better_supabase.schedule_job. Leave out next_run: the
schedule is stored without one, and the next drain computes its first run
in the schedule's time zone, counted from when it was written, so a run
that falls between the write and the drain still happens. The app needs no
job that copies rows into schedules: a trigger on the table that holds the
cron keeps the schedule current, and the drain route runs it.
create function public.schedule_reminder()
returns trigger
language plpgsql
security definer
set search_path = ''
as $$
begin
perform better_supabase.schedule_job(
'reminder:' || new.id,
new.cron,
'send_reminder',
jsonb_build_object('id', new.id),
new.time_zone,
null,
new.organization_id::text
);
return null;
end;
$$;
create trigger schedule_reminder
after insert or update of cron, time_zone on public.reminders
for each row execute function public.schedule_reminder();Make the trigger function security definer, so a signed-in user's insert
can write the schedule without access to better_supabase. Writing a
schedule again with the same cron and time zone keeps its next run; a new
cron or time zone clears it, and the next drain computes the first run of
the new one. Call better_supabase.unschedule_job('reminder:' || old.id)
from a delete trigger to remove it.
schedule_job rejects a schedule that isn't five cron fields, a macro or an
interval. A cron that passes that check but that the drain can't read (a
field out of range) stays without a next run, and listSchedules shows its
nextRun as null.
External schedulers
The drain scheduler needs no database extension: any scheduler that calls the drain route on a fixed cadence runs every schedule, whatever their own cadences. Call it at least as often as your most frequent schedule, every minute for minute-level schedules. Vercel Cron, a GitHub Actions workflow, a Kubernetes CronJob or a Supabase Edge Function on a timer all work; each call leases the due schedules, so overlapping calls never enqueue a run twice.
Three scheduling layers
Three blocks take a cron, and each one is for a different owner. All three
parse the cron and compute the next run with the same helpers
(assertCron and nextCronRun), so a cron and a time zone mean the same
thing in each.
| Layer | Who sets it | What a run does |
|---|---|---|
| Job schedules (this block) | The app, in code or SQL | Enqueues a job with a fixed payload |
| Workflow schedules | The app, or members with workflow.admin | Starts a workflow run, with the idempotency key schedule:<id>:<fire time> |
| AI tasks | Each user, for their own prompt | Runs the prompt and logs the chat each run wrote to |
Pick the lowest layer that does the job. A nightly cleanup is a job schedule. A recurring start of a durable workflow that members can see and pause is a workflow schedule. A prompt a user writes and schedules in the product is an AI task, which the jobs block then runs.
Drain route
drainRoute turns the queues into an HTTP endpoint for a cron caller such as
Vercel Cron. Each call checks the bearer
secret, enqueues due schedules, then drains each queue in handlers until
they are empty or budgetMs (50 seconds by default) is spent. Jobs it
already claimed still finish.
import "server-only";
import { jobs } from "@/lib/jobs";
export const GET = jobs.drainRoute({
secret: process.env.CRON_SECRET,
budgetMs: 50_000,
handlers: {
send_email: (payload) => sendEmail(payload),
},
onError: (error, job) => console.error(error.message, job?.id),
});{
"crons": [{ "path": "/api/jobs/drain", "schedule": "* * * * *" }]
}Vercel sends Authorization: Bearer <CRON_SECRET>; a request without it gets
401. The response is JSON: schedules (runs enqueued), queues (the
succeeded and failed counts per queue), budgetExhausted and errors (the
schedule and queue runs that failed as a whole). Keep budgetMs below the
function's maximum duration. drain takes the same budgetMs when you call
it yourself.
Nothing calls the route under next dev or another local server, so
outbox relays and queued jobs wait until a deploy. devDrain calls it on an
interval (every minute by default) in development, and devDrainSecret
generates the route's secret when none is set. Call both before the route
reads the secret, for example in Next.js instrumentation.ts:
import { devDrain, devDrainSecret } from "better-supabase/blocks/jobs";
export function register() {
if (process.env.NODE_ENV !== "development") return;
if (process.env.NEXT_RUNTIME !== "nodejs") return;
devDrain({
url: `http://localhost:${process.env.PORT ?? "3000"}/api/jobs/drain`,
secret: devDrainSecret(process.env),
});
}devDrain sends GET with Authorization: Bearer <secret>, skips a call
while the previous one still runs, and reports a failed call or an answer
other than 2xx to onError (console.warn by default) without stopping.
every sets the interval in milliseconds; stop() or an aborted signal
ends it, and tick() calls the route once. devDrainSecret(env, name)
returns env[name] (CRON_SECRET by default) or writes a random one there.
monitor hooks around each authorized drain report to a cron monitor
without wrapping the route. onStart returns a value, such as a check-in
id, that onFinish receives with the result. A hook that throws goes to
onError and never fails the drain.
export const GET = jobs.drainRoute({
secret: process.env.CRON_SECRET,
handlers,
monitor: {
onStart: () =>
Sentry.captureCheckIn({
monitorSlug: "jobs-drain",
status: "in_progress",
}),
onFinish: (result, checkInId) =>
Sentry.captureCheckIn({
checkInId: String(checkInId),
monitorSlug: "jobs-drain",
status: result.errors > 0 ? "error" : "ok",
}),
},
});Queue backends
createJobs accepts a SQL connection, a Supabase client, or any
QueueBackend (API version 1). sqlQueueBackend(sql) and
pgmqPublicBackend(client) are the two built in; implement the interface to
put jobs somewhere else, and prove it with testQueueBackend from
conformance kits.
import { createJobs, type QueueBackend } from "better-supabase/blocks/jobs";
const backend: QueueBackend = myQueue;
const jobs = createJobs(backend, { send_email: schema });Without a database connection
Where there's no direct connection, such as an Edge Function without a pooler
URL, pass a service-role Supabase client instead. Jobs then go through the
pgmq_public RPCs, which you turn on in the dashboard under Integrations, Queues, "Expose
Queues via PostgREST".
const jobs = createJobs(createClient(url, secretKey), { send_email: schema });Over PostgREST you get enqueue, claim, complete, fail, drain and
work, with these differences:
- there is no lease check on
complete; - a failed job reappears when its lease ends (
retryInis ignored), and there's no heartbeat; dedupeKey,schedule,stats,listDeadandretryDeadreturninvalid_request.
Idempotency keys
handle wraps a handler so that a retried request with the same
Idempotency-Key gets the first response back instead of running twice:
import { createIdempotency } from "better-supabase/blocks/jobs";
const idempotency = createIdempotency(postgres.admin, { ttl: "24 hours" });
export const POST = (request: Request) =>
idempotency.handle(request, async () =>
Response.json(await createOrder(request), { status: 201 }),
);| Situation | Response |
|---|---|
| First request | The handler's response, stored |
| Same key and same body | The stored response, with idempotency-replayed: true |
| Same key while the first request is still running | 409 with Retry-After |
| Same key and a different body | 422 with code IDEMPOTENCY_KEY_REUSED |
| Handler throws or returns a 5xx | The key is released, so the client can retry |
The fingerprint is a SHA-256 of the method, path and body. Set
required: true to reject requests that have no key. The error responses
are Problem Details; pass problem to answer in your app's error format
instead (see Problem Details).
Outside HTTP, begin(key, fingerprint, scope) returns a holder token with
the started state, and complete(key, holder, status, body, scope) and
release(key, holder, scope) need it. A caller whose lock ran out may find
that another caller took the key over; its complete and release then
return false and change nothing, so it can't store its late result over
the new holder's or free a key someone else holds:
const started = await idempotency.begin(key, fingerprint).orThrow();
if (started.state === "started") {
const stored = await idempotency
.complete(key, started.holder!, 200, result)
.orThrow();
if (!stored) console.warn("lost the key to another caller");
}Leases
withLease(sql, key, fn) runs fn while it holds the lease on key, so
work that must not run twice at once (one automation per conversation, one
mutation per agent run) runs once, and other callers get
{ acquired: false }. The lease holds for seconds (60 by default) and is
released when fn ends or throws; for long work, lease.extend() moves
the expiry and returns false once the lease was lost:
import { withLease } from "better-supabase/blocks/jobs";
const outcome = await withLease(
postgres.admin,
`conversation:${conversationId}`,
async (lease) => {
await step1();
await lease.extend();
return step2();
},
{ seconds: 30 },
).orThrow();
if (!outcome.acquired) return; // another worker has itIn SQL the same lease is better_supabase.acquire_lease(key, seconds, scope), which returns a holder token or null, extend_lease(key, holder, seconds, scope) and release_lease(key, holder, scope), for the service
role. An expired lease goes to the next caller, and the former holder's
extend_lease and release_lease return false.
Webhook inbox
The inbox verifies a webhook, stores it once, and returns quickly. A worker then processes it, with the same retries as jobs. This way your provider never times out, and a redelivered webhook is only handled once.
import { createWebhookInbox } from "better-supabase/blocks/jobs";
const inbox = createWebhookInbox(postgres.admin, {
source: "billing",
secrets: process.env.BILLING_WEBHOOK_SECRET!,
});
export const POST = (request: Request) => inbox.receive(request);
await inbox.process(async (message) => {
if (message.type === "invoice.paid") await markPaid(message.payload);
});createWebhookInbox was called createInbox before 0.6, and its types were
Inbox, InboxOptions and so on. 0.6 removes the old names;
createInbox from better-supabase/blocks/inbox is the conversation inbox
block.
secrets verifies Standard Webhooks signatures. Pass
verify to use a provider's own scheme instead. receive answers 202 for
a new message and 200 for a duplicate. A message id is unique per source
and tenant, so two tenants of one provider can send the same id. A failed
verification gets a Problem Details response, and a body over maxBodyBytes
(1 MiB by default) gets 413 before it is read in full.
A failed message is retried with exponential backoff until maxAttempts
(8 by default, set per source on createWebhookInbox), then marked dead. Each
message stores the limit it arrived with. The handler sees attempts (1 on
the first) and maxAttempts, so it can tell its last attempt, for example
to notify someone before the message is given up:
const crm = createWebhookInbox(postgres.admin, {
source: "crm",
secrets,
maxAttempts: 3,
});
await crm.process(async (message) => {
const result = await syncContact(message.payload);
if (!result.ok && message.attempts === message.maxAttempts) {
await alertOwner(message, result.error);
}
return result;
});A message whose worker died during its last attempt is marked dead by the
next process instead of running again.
process runs until no message is ready. In a serverless function, pass
budgetMs so a backlog can't outrun the function's maximum duration: it
stops claiming once the budget is spent, finishes the messages it already
claimed, and the next call picks up the rest. batch (10) and lease (300
seconds) set how many messages one claim takes and how long they stay with
the worker. While a handler runs, a heartbeat renews its message's lease
every half lease. A worker whose lease ran out and was taken by another
claim can't complete, fail or checkpoint that message: the inbox checks the
attempt it claimed, and skips a message whose lease it can't renew.
export const GET = async () =>
Response.json(await inbox.process(handleMessage, { budgetMs: 50_000 }));A webhook carries no session, so pick the identity in the handler. Look up
the user and tenant the event belongs to (a Stripe customer id maps to an
organization and its owner), then run the work with
bs.actingAs(userId, { tenant_id }) so RLS applies:
await inbox.process(async (message) => {
const owner = await ownerOfCustomer(message.payload.customer);
const db = bs.actingAs(owner.userId, { tenant_id: owner.organizationId });
await db.invoices
.update(message.payload.invoice, { status: "paid" })
.orThrow();
});Use bs.admin() only for events no user owns, such as a provider-wide
status change. See
Jobs, webhooks and agents without a session.
Per-tenant integrations
Most inboxes serve integrations a tenant connected, so a message can belong to
a tenant: tenantOf(payload) reads it from the payload, and verify can
return it next to the id and payload. list({ tenant }) shows a tenant's
messages newest first (an integration's delivery log), and
purge({ tenant, olderThan: 0 }) removes them when the tenant goes.
const inbox = createWebhookInbox(postgres.admin, {
source: "chat-provider",
secrets: process.env.CHAT_WEBHOOK_SECRET!,
tenantOf: (payload) => connectionTenant(payload),
});
const deliveries = await inbox
.list({ tenant: organizationId, status: "dead" })
.orThrow();A provider that pages its deliveries needs progress that survives a retry.
message.checkpoint(fields) merges fields into the stored progress while
the worker holds the message, and the next attempt reads it from
message.progress:
await inbox.process(async (message) => {
let cursor = message.progress.cursor as string | undefined;
do {
const page = await provider.fetchPage(message.payload, cursor);
await apply(page.items);
cursor = page.next;
await message.checkpoint({ cursor });
} while (cursor);
});inbox.store({ id, payload, type, tenant }) stores an event your code
already verified, for providers whose SDK verifies and parses in one call:
export async function POST(request: Request) {
const event = await stripe.webhooks.constructEventAsync(
await request.text(),
request.headers.get("stripe-signature")!,
process.env.STRIPE_WEBHOOK_SECRET!,
);
const stored = await inbox.store({
id: event.id,
type: event.type,
payload: event,
});
return stored.ok
? new Response(null, { status: 202 })
: new Response(null, { status: 500 });
}A source that only store fills, such as a chat SDK that verifies its own
requests, needs neither secrets nor verify. Its process, list and
purge work as usual, and receive throws a TypeError, so an unverified
request is never stored:
const chat = createWebhookInbox(postgres.admin, { source: "chat" });
await chat.store({ id: event.id, type: event.type, payload: event });Last updated on