tialuk · · 7 min read

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

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

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

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

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.

← Back to the portfolio · hello@abrahamkazerian.dev