Most video webhook receivers do the work inline. Take the event, call the provider's API back for the playback ID and the duration, write the row, answer. On one upload against a warm function that is fine, and pgmq is overkill.
The budget is smaller than it looks. Shopify gives a receiver five seconds, retries eight times over four hours, and after eight consecutive failures deletes the subscription, if it was configured using the Admin API. Stripe retries for up to three days with exponential backoff before giving up. A cold start plus one API round trip can eat most of a five-second budget, and the sender can't tell a slow write from a dead endpoint, so it retries and you process the event twice.
So the work has to leave the receiver: verify, record, enqueue, answer, and let something else drain the queue on a timer. Supabase documents the basic wiring of a cron job into a consuming edge function. We go straight at the operational half: the retry budget, the dead letter rule, the backoff that keeps a failing message off a rate-limited API, and a record that tells an empty queue from an empty inbox.
What the webhook handler should do
We give the handler three jobs and nothing else: verify the signature, record that the event arrived, and push the work onto a queue. All three are local, so your response time stops depending on how the provider's API is feeling.
select pgmq.create('webhook_jobs');Supabase exposes the queue functions through a pgmq_public wrapper schema, so an edge function reaches them with the same client it uses for tables.
import { createClient } from "jsr:@supabase/supabase-js@2";
const db = createClient(
Deno.env.get("SUPABASE_URL")!,
Deno.env.get("SUPABASE_SERVICE_ROLE_KEY")!,
);
Deno.serve(async (req) => {
const raw = await req.text(); // raw bytes, never req.json()
if (!verifySignature(raw, req.headers.get("x-signature"))) {
return new Response("bad signature", { status: 401 });
}
const event = JSON.parse(raw);
// insert ... on conflict ("eventId") do nothing
const { data: inserted } = await db
.schema("app")
.from("webhook_events")
.upsert(
{ eventId: event.id, type: event.type, payload: event },
{ onConflict: "eventId", ignoreDuplicates: true },
)
.select("id");
if (!inserted?.length) {
return new Response("duplicate", { status: 200 }); // already queued
}
await db.schema("pgmq_public").rpc("send", {
queue_name: "webhook_jobs",
message: event,
});
return new Response("ok", { status: 200 }); // answer fast, work later
});Read the body as text and verify against those exact bytes, because req.json() parses and reserializes, and what you then hash is no longer what the provider signed. The insert lands before the enqueue, and the conflict clause turns a redelivery into a cheap 200 rather than a second copy of the job.
We ship this as fastpix-webhook, running with verify_jwt set to false because FastPix signs its own requests rather than sending a Supabase JWT. Our signing secret is base64, so a re-typed or truncated value fails verification quietly: no rows, no error.
Wire a consumer that runs itself
Nothing in Postgres calls pgmq.read for you. It's an ordinary function, so something has to call it on a timer, and on Supabase that something is three parts: pg_cron for the schedule, pg_net for the outbound call, and Vault for the URL and key it needs.
Store the two values the job needs in Vault, so the SQL never carries a service key in plain text.
select vault.create_secret(
'http://host.docker.internal:54321/functions/v1', 'fastpix_functions_url');
select vault.create_secret('<SERVICE_ROLE_KEY>', 'fastpix_service_role_key');On a local stack the host matters, because pg_net runs inside the Postgres container where localhost means the container itself. Get either name wrong and everything still looks healthy from outside, which is why we put the Vault step ahead of the webhook step in our own docs.
select cron.schedule(
'drain_webhook_jobs',
'10 seconds',
$$
select net.http_post(
url := (select decrypted_secret from vault.decrypted_secrets
where name = 'fastpix_functions_url') || '/webhook-worker',
headers := jsonb_build_object(
'Content-Type', 'application/json',
'Authorization', 'Bearer ' || (select decrypted_secret
from vault.decrypted_secrets where name = 'fastpix_service_role_key')
),
body := '{}'::jsonb
);
$$
);Now the worker, which claims a bounded batch and archives only what succeeded.
const BATCH = Number(Deno.env.get("BATCH_SIZE") ?? 10);
const VT_SECONDS = 60;
Deno.serve(async () => {
const { data: messages } = await db.schema("pgmq_public").rpc("read", {
queue_name: "webhook_jobs",
sleep_seconds: VT_SECONDS, // how long this batch stays invisible
n: BATCH,
});
for (const m of messages ?? []) {
try {
await handleEvent(m.message);
await db.schema("pgmq_public").rpc("archive", {
queue_name: "webhook_jobs",
message_id: m.msg_id,
});
} catch (err) {
await recordFailure(m.msg_id, err); // do not delete, do not archive
}
}
return new Response("ok");
});The catch branch is the retry and it works by doing nothing: a message you neither delete nor archive becomes visible again when its visibility timeout expires, and the next run picks it up with read_ct one higher. We drain every ten seconds into fastpix-worker, and that interval carries the one caveat we had to design around. Edge functions open a direct, non-pooled Postgres connection, so a small instance can approach its limit. Point them at Supavisor, or slow our drain in 0003_fastpix_setup_cron_job.sql.
When you don't need a queue
A queue isn't always the answer, and two alternatives come up before it.
Do the work in the function, after responding. Send the 200, then carry on in the same invocation. That works while the follow-up is fast and cheap enough to lose, and it stops working the moment the work can be slow or can fail, because a serverless function can be torn down the instant it responds.
Put a hosted queue in front. A managed queue does all of this and more, and the cost is a second system: more credentials, another failure mode, another bill. If the queue would be the only thing outside Postgres, pgmq keeps the enqueue in the same transaction as the write that caused it.
We reach for a queue when the work can fail, which is exactly this case: a receiver answering inside a stranger's timeout, doing work that calls back out to that same stranger.
How many messages per pgmq read?
Pick the largest batch whose worst-case processing time still fits inside both the function's wall clock and the visibility timeout you passed to read. The batch argument is n through pgmq_public and qty on pgmq.read itself, and it moves throughput, cost and failure blast radius at once, which is why leaving it at 1 is a decision nobody makes on purpose.
| Constraint | The question to answer | What it caps |
|---|---|---|
| Function wall clock | How long does the slowest message take end to end? | Batch size times worst-case time must finish before the function is killed |
| Visibility timeout | What vt did you pass to read? | The batch must finish before vt expires, or message one reappears while you are on message five |
| Downstream rate limit | How many API calls does one message make? | Batch size times runs per minute must stay under the provider's ceiling |
| Blast radius | What happens if the function dies mid-batch? | Every message in the batch gets its read_ct incremented, including ones you never touched |
The blast radius row surprises people: a crash on message two still costs every message in the batch an attempt it never used. Our own worker reads a fixed ten per run, which is a number in the code rather than a setting, and 5 to 10 suits most handlers making one or two outbound calls per message. Raise the batch before you shorten the interval, because a shorter interval multiplies connections and a bigger batch does not.
Dead-letter with pgmq read_ct
pgmq increments read_ct every time it hands a message to a consumer, and the number we want you to hold on to is that a fresh message reaches your worker already reading 1 rather than 0. A message at 5 has been taken five times and deleted zero times, which isn't slow but unprocessable, and it will cycle until somebody stops it.
select msg_id, read_ct, enqueued_at, message->>'type' as event_type
from pgmq.q_webhook_jobs
where read_ct > 3
order by read_ct desc;This is the threshold our own worker runs on, under the name FASTPIX_MAX_READ_CT, defaulting to seven: once a message has been handed out that many times it gets archived rather than cycled again. Dead-lettering is then a second queue and a guard at the top of the loop:
select pgmq.create('webhook_jobs_dead');const MAX_ATTEMPTS = 5;
for (const m of messages ?? []) {
if (m.read_ct > MAX_ATTEMPTS) { // strictly greater: read_ct starts at 1
await db.schema("pgmq_public").rpc("send", {
queue_name: "webhook_jobs_dead",
message: { payload: m.message, attempts: m.read_ct },
});
await db.schema("pgmq_public").rpc("archive", {
queue_name: "webhook_jobs",
message_id: m.msg_id,
});
continue;
}
// ... normal handling
}The comparison operator is the whole trap, because read_ct is already 1 on the first read, so a >= against a budget of five retires the message on its fifth appearance and therefore after four real attempts. That is a fifth of the budget lost to an off-by-one nobody catches in review. Put the check above the try as well, because if processing the message is what kills the worker then code after the attempt never runs.
The default retry is worse than it looks too. A failed message comes back after exactly the visibility timeout you passed to read, every time, so an API that just rate-limited you gets hit again in sixty seconds, and again sixty seconds later. We push the message further out instead, using read_ct as the exponent:
create or replace function app.backoff(p_msg_id bigint, p_read_ct int)
returns void language sql security definer as $$
select pgmq.set_vt('webhook_jobs', p_msg_id,
least(30 * power(2, p_read_ct)::int, 3600));
$$;The pgmq_public wrapper doesn't expose set_vt, which is why we reach it through a small function of our own. Call it from the catch branch and the gaps grow from half a minute to an hour while the budget stays at five. Alert on max(read_ct) rather than depth: depth says the queue is busy, read_ct says it's stuck.
Keep a record beside the queue
A queue isn't a log. Delete a message and the row is gone; archive it and the row moves to the archive table, which holds the payload but not the outcome. Either way the queue reaches zero, and zero has two meanings: everything was handled, or nothing ever arrived.
create table if not exists app.webhook_events (
id uuid primary key default gen_random_uuid(),
"eventId" text unique not null,
"type" text not null,
"receivedAt" timestamptz not null default now(),
"processStatus" text not null default 'received',
"lastError" text,
payload jsonb not null
);The receiver inserts a row with status received before it enqueues, and the worker updates that row to processed or failed with the error text. The unique constraint on "eventId" makes deduplication possible, but on its own it doesn't give you idempotency: a plain insert hitting a unique violation raises, and a raised error in a webhook receiver becomes a 500 that tells the provider to deliver the event again, forever. That is why we swallow the conflict deliberately and answer 200 on the duplicate path.
This is fastpix.webhook_events in our integration, and its column names match the FastPix API, so they are camelCase and need double quotes in SQL.
select "type", "processStatus", "lastError" from fastpix.webhook_events
order by "receivedAt" desc limit 10;| What you see | What it means |
|---|---|
| No rows at all | FastPix is not reaching the webhook, or the signature is failing |
| Rows stuck at received | Nothing is draining the queue. Recheck the Vault secrets |
| Rows at failed | Read lastError. Usually a wrong token ID or secret |
The middle row is the failure mode with no error attached. Miss the Vault secrets and events arrive, land in the queue and sit there while the cron job fires on time and its HTTP call fails silently. Without this table you are guessing, and the troubleshooting page maps the other symptoms to their causes.
Re-fetch instead of trusting the payload
Re-fetch the resource by its ID rather than writing what the payload says. A webhook payload is a snapshot from the instant the event fired, and your worker reads it seconds or minutes later, possibly out of order after a retry. Writing that snapshot into your tables means writing state that may already be wrong.
Consider an upload that moves from created to processing to ready in a few seconds, firing three events. Retried at different times, they can drain in any order, and trusting the payloads leaves your row saying processing for an asset that is already playable.
Re-fetching by ID removes the ordering problem, because each event becomes a trigger meaning "this resource changed, go look" and the API answers with current state. Our worker takes the event as a signal and writes fastpix.media, fastpix.live_streams or fastpix.uploads from that response, and npx @fastpix/supabase reconcile sweeps the last 24 hours for events that never arrived at all.
Generate the pgmq queue with one command
Everything above is a few hundred lines of SQL and TypeScript, and if the events you're queuing are video events you need write none of it. One command provisions the queue, both cron jobs and all four edge functions, with the audit table and the drain wired:
npx @fastpix/supabase initRun it against a local stack, add the two Vault secrets, upload one video, and watch a row appear in fastpix.media within seconds. The integration docs carry the full sequence, and you can run the queue against a real project on the free plan, ten videos and no card, enough to watch it move and to break it on purpose.
Frequently Asked Questions (FAQs)
What does read_ct mean in pgmq?
read_ct is the number of times pgmq has handed that message to a consumer, and it is already 1 the first time your worker sees it rather than 0. A message sitting at 5 has been attempted five times and deleted zero times, which makes it the counter to dead-letter on.
How do you build a dead letter queue with pgmq?
Create a second queue, then check read_ct at the top of the consumer loop before attempting work. Write the comparison as read_ct > budget rather than read_ct >= budget, because read_ct starts at 1 and the greater-or-equal form quietly costs you one attempt.
What actually calls pgmq.read on a schedule in Supabase?
Nothing does until you wire it, because Postgres has no background consumer. On Supabase the usual arrangement is a pg_cron job that uses pg_net to POST to an edge function, and that function reads, processes and archives. Supabase documents the pattern, and pg_cron accepts schedules like '10 seconds'.
How many messages should a pgmq consumer read per run?
Pick the largest batch whose worst-case processing time fits inside both the function's wall-clock limit and the visibility timeout passed to read. A handler making one or two outbound calls per message sits at 5 to 10. Our own worker reads a fixed ten per run. Don't confuse that with FASTPIX_MAX_READ_CT, which despite sitting nearby is a retry cap rather than a batch size.
How do you add a retry backoff to a pgmq consumer?
Call pgmq.set_vt in the catch branch so the delay grows with read_ct instead of staying at the timeout you passed to read. Without it a failing message returns on a fixed interval, which is a hot loop pointed at whatever just rejected you. The pgmq_public wrapper does not expose set_vt, so wrap it yourself.
Why is my pgmq queue filling up with nothing draining?
Almost always the scheduled consumer is not reaching the worker. Check the cron job exists and is running, then check its secrets, because a job built from Vault secrets that are missing or wrong fails its HTTP call quietly. On a local stack the URL must use host.docker.internal.
Should a worker trust the webhook payload or re-fetch the resource?
Re-fetch the resource by its ID rather than trusting the payload. A payload is a snapshot from the moment the event fired, and your worker reads it later and possibly out of order, so writing it can overwrite newer state. Re-fetching is idempotent.






