Ingesting Meta webhooks with a Postgres outbox
Durable rows, safe retries and real-time updates for Instagram and Facebook DMs, without a message broker.
tialuk, my personal project, is a B2B inbox for sales conversations that happen in Instagram and Facebook DMs. Every inbound message, read receipt, reaction or edit reaches the backend (Bun, Elysia, PostgreSQL, MikroORM) as a Meta webhook. Here's how the ingestion pipeline works, why it runs on a Postgres table, and where it still falls short.
What Meta expects
- Acknowledge fast. Meta expects an answer within seconds. Database work and AI calls inline risk a timeout, and a timeout means redelivery.
- Redelivery means duplicates. Processing has to be idempotent.
- Verify the sender. Each POST carries
X-Hub-Signature-256, an HMAC-SHA256 of the raw body keyed with the app secret. - Never lose an event. Returning 200 before the event is durable loses it.
Options I considered
My original architecture notes said: publish each webhook to GCP Pub/Sub, return 200, consume it from the same server. That is still the long-term target, but it depended on an infrastructure-as-code track I hadn't built, and I didn't want ingestion blocked on Terraform. Processing inline was out because of timeouts. An in-memory queue would acknowledge fast but lose everything on a restart or deploy.
I chose a table in the database I already had. The route verifies the signature, does one INSERT with the raw payload and returns 200; a worker loop in the same process does the rest. The table survives restarts, and debugging is one query: SELECT * FROM webhook_inbox WHERE status = 'FAILED'. (The codebase calls it an outbox; strictly it's closer to an inbox: the request only writes the event.)
I deliberately skipped a transport abstraction. My first generic WebhookTransport turned out to be just a consumer of this table, and its semantics (batch claims, expiring leases, terminal vs. retryable failures) don't map one-to-one onto Pub/Sub's ack/nack/dead-letter model. It waits for a second real implementation.
The pipeline, step by step
1. Verify the raw bytes, store, acknowledge
Elysia parses bodies according to Content-Type, and any re-serialization breaks the HMAC, so the route disables parsing:
// route option: parse: "none"
const rawBody = await request.text();
const signature = request.headers.get("x-hub-signature-256");
if (!verifySignature(rawBody, signature, appSecret)) return status(401);
let payload;
try { payload = JSON.parse(rawBody); } catch { return status(400); }
const row = await enqueue(payload, "META", signature);
if (!row) return status(500); // not stored: let Meta retry
return new Response("", { status: 200 });
function verifySignature(rawBody: string, header: string | null, secret: string) {
if (!header?.startsWith("sha256=")) return false;
const provided = header.slice("sha256=".length);
if (provided.length !== 64) return false;
const expected = createHmac("sha256", secret).update(rawBody).digest();
const providedBuf = Buffer.from(provided, "hex");
// non-hex input decodes to a shorter buffer
if (providedBuf.length !== expected.length) return false;
return timingSafeEqual(providedBuf, expected);
}
Signature first (401), JSON second (400), then the insert. A failed insert returns 500 on purpose, so Meta retries. The ticket's target for this path was 500 ms; it does a single insert.
2. The table
Each row holds the raw jsonb payload, a source (META or PADDLE), a status (PENDING, PROCESSING, DONE, FAILED), an attempt counter, the last error, processing_started_at and processed_at.
3. Claim with FOR UPDATE SKIP LOCKED
Every second (configurable), the worker claims up to 10 rows, oldest first:
UPDATE webhook_inbox
SET status = 'PROCESSING', processing_started_at = now()
WHERE id IN (
SELECT id FROM webhook_inbox
WHERE status = 'PENDING'
OR (status = 'PROCESSING'
AND processing_started_at < now() - INTERVAL '5 minutes')
ORDER BY received_at
LIMIT 10
FOR UPDATE SKIP LOCKED
)
RETURNING id;
SKIP LOCKED lets replicas run the same statement without blocking each other or claiming the same row. The claim is a single autocommit statement, so the lock only lives for the claim; afterwards, ownership is the PROCESSING status plus its timestamp, effectively a lease. It returns only IDs; the ORM then loads the rows, since raw RETURNING * yields snake_case rows, not entities.
4. Parse, route, record the outcome
A parser flattens Meta's two payload shapes (messaging[] and changes[]) into typed events and never throws; a malformed payload yields no events and the row is marked DONE with a warning. The dispatcher resolves the connected account from each entry's ID, skips unknown channels, and routes by field to five handlers: messages, reads, reactions, edits and story mentions.
Rows run one at a time, each in its own try/catch, so a poison row can't stall the batch. Success marks DONE. Expired-token or invalid-parameter errors from the Graph API go straight to FAILED, since retrying can't fix them. Anything else returns to PENDING, up to five attempts, then FAILED.
Idempotency and failure
The pipeline is at-least-once: Meta can deliver an event twice (two rows), a row can be retried after a partial failure, and an expired lease can be reclaimed while the original run is still alive. So the guarantee lives in the handlers, anchored on database constraints.
Messages. platformMessageId is unique, and the insert treats a violation as "already stored":
try {
await em.persist(message).flush();
return message; // first delivery
} catch (err) {
if (err instanceof UniqueConstraintViolationException) {
return null; // redelivery: already stored
}
throw err; // anything else: let the worker retry
}
Everything downstream (unread counter, reopening a closed lead, AI intent classification, the message SSE events) runs only when that insert returned a row.
Leads are unique on (channel, socialUserId); find-or-create catches the violation and re-reads the winner, so racing replicas converge. Read receipts use UPDATE … WHERE status <> 'READ' and publish only if rows changed. Reactions are unique on (message, actor) with insert-or-update. Edits overwrite the content by message ID.
Recovering abandoned rows
If a replica dies mid-batch, its rows stay in PROCESSING. My first version reset every such row at startup, which works with one process and breaks with several: a replica booting during a rolling deploy would steal rows another one was still processing. The TTL in the claim replaced it: five minutes in PROCESSING means abandoned. Dispatch is expected to take well under a second, and idempotent handlers make a rare double run harmless.
Graceful shutdown covers the common case: on SIGTERM the server stops taking requests, waits for the in-flight batch, then closes the connection pool. A daily pass purges DONE rows older than 30 days using a partial index; the DELETE is idempotent, so replicas need no coordination.
Reusing the pattern
- Paddle billing webhooks share the table with
source = 'PADDLE': raw body, signature check through Paddle's SDK, insert, 200. Since a reclaimed lease could run a handler twice, billing events are also deduplicated by Paddle's event ID. - Meta data-deletion requests have their own table but copy the claim:
FOR UPDATE SKIP LOCKEDand the same five-minute TTL. - Periodic jobs (token health checks, AI summaries, billing renewals, retention) are gated sub-passes of the same one-second tick, not separate timers; the ones that claim rows use
SKIP LOCKEDtoo.
Real-time updates with SSE
After a handler writes to the database, it publishes to an in-process broadcaster: a map from tenant ID to that tenant's open Server-Sent Events streams. GET /events requires the session cookie and subscribes to the caller's tenant. Handlers emit lead.created, message.created, lead.updated, reaction.changed, message.edited and message.read_by_customer. Each carries the inbox row ID as a correlation ID, and the pipeline's logs carry it too, so one grep traces a webhook to its lead, message and events.
Bun's default idle timeout (about 10 seconds) kept closing quiet streams, so the server raises it to 120 seconds and each stream sends a : ping comment every 25. And streams that throw on write are dropped, which cleans up disconnected clients.
How I tested it
- Signature verification: valid signature, tampered body, wrong secret, missing or malformed headers, 64-character non-hex signatures, UTF-8 payloads with emoji.
- The route, with the repository mocked: 401 for missing or bad signatures, 400 for invalid JSON, 500 when the insert fails, 200 with payload, source and signature persisted.
- Parser, dispatcher and handlers take dependencies as optional parameters, so tests inject stubs instead of global module mocks. The message suite covers duplicates explicitly: a redelivered message to a closed lead must neither reopen it nor emit events.
- SSE: tenant isolation and cleanup of dead streams.
- End to end: Meta's dashboard "Test" button sends nothing for the Instagram
messagesfield, so a script signs fixture payloads with the app secret and posts them to the endpoint. With it I verifiedPENDING→ claimed →DONE, lead and message created, and a replayed message ID with no duplicate.
Running the real worker also exposed a bug the unit tests had missed: my raw SQL used $1-style placeholders, which the query layer doesn't bind, so the status updates had been silently broken.
The honest gap: SKIP LOCKED and the TTL reclaim have no automated multi-connection test. A similar concurrency test elsewhere is skipped, because ORM forks inside one process can share a connection and hide the lock. That invariant rests on Postgres semantics, not on my suite.
Limitations and next steps
- SSE is in memory, per instance. An event published on replica A never reaches a browser connected to replica B. The worker is multi-instance safe and the pool is sized for up to five replicas; the broadcaster isn't. Dev caps the service at two instances, and the documented production plan is one until there's a shared bus (Redis pub/sub or Postgres
LISTEN/NOTIFY). - No replay. An event published while no client is connected is dropped with a warning.
- The worker shares the web process, so production needs an always-on instance with CPU always allocated; with request-based CPU, a one-second loop would freeze.
- No backoff. Failed rows retry on the next tick; a
next_attempt_atcolumn is planned for when handlers start hitting rate limits. - Pub/Sub is deferred, not dropped. The infra plan has the topics and a dead-letter topic with five max deliveries, but none of it is provisioned yet.
Closing
A table, one UPDATE … SKIP LOCKED, a lease timestamp and a few unique constraints give this project what it needs from a queue today: fast acknowledgements, durability across restarts, retries and harmless duplicates. What caps horizontal scaling right now is the in-memory SSE fan-out, not the queue.