DEV Community

Tal Aizikov
Tal Aizikov

Posted on AI-assisted

A Postgres table as our push queue: retries, dead letters, and FCM in Next.js

A push notification looks like a one-liner: call Firebase Cloud Messaging with a device token, a title and a body. Then the real questions arrive. What happens when FCM is down for a minute? What if the user turned reminders off? What if the token belongs to a phone that was wiped last month? And how do you know whether anyone opened the thing?

We're building MotionLab, personalized stretching and mobility sessions. In an earlier post we covered how AI sessions pick exercises. This one is about a less glamorous part of the same MotionLab-Web repo: the notification backend. There's no queue service and no worker fleet. One Postgres table is the queue, the state machine and the audit trail, and a handful of Next.js route handlers drive it.

Stack: Next.js 16.2 (App Router route handlers on the Node runtime), Supabase Postgres, firebase-admin for FCM, and Vercel Cron. Snippets are trimmed from the real code.

1. The table is the queue

Every notification is a row, and its status column is the state machine: pending, sent, delivered, retrying or dead.

create table if not exists public.notifications (
    id              uuid primary key default gen_random_uuid(),
    user_id         uuid not null references auth.users(id) on delete cascade,
    type            text not null,               -- maps to a preference category
    title           text not null,
    body            text not null,
    data            jsonb not null default '{}'::jsonb,
    channel         text not null default 'push', -- push | email | sms (future)
    status          text not null default 'pending',
    attempt_count   int  not null default 0,
    next_attempt_at timestamptz not null default now(),
    failed_reason   text,
    created_at      timestamptz not null default now(),
    sent_at         timestamptz,
    delivered_at    timestamptz,
    opened_at       timestamptz,
    read_at         timestamptz
);

-- Cron dispatch query: due rows still in flight
create index if not exists notifications_due_idx
    on public.notifications(next_attempt_at) where status in ('pending', 'retrying');
Enter fullscreen mode Exit fullscreen mode

A few things this buys:

  • Retries are just data. attempt_count, next_attempt_at and failed_reason live on the row, so a retry is an UPDATE, not a message re-published somewhere else.
  • The partial index only covers work in flight. Sent, delivered and dead rows pile up over time, but the index the dispatcher reads only holds pending and retrying rows. Postgres partial indexes are made for exactly this "small hot subset of a big table" shape.
  • Lifecycle timestamps are separate columns. sent_at, delivered_at, opened_at and read_at are each set once, which makes the metrics in section 7 plain SQL.

Device tokens get their own table with a unique token and an active flag (plus another partial index on active tokens per user), and preferences are one row per user with a boolean per category.

2. Enqueue: check preferences before writing anything

All sends go through one enqueue() function. The first thing it does is read the user's preferences and decide whether the notification should exist at all:

// notification.type → notification_preferences column. Unknown types are not gated.
const CATEGORY_COLUMN: Record<string, string> = {
  marketing: "marketing",
  reminder: "reminders",
  reminders: "reminders",
  system: "system_updates",
  system_update: "system_updates",
  // ...
};

if (channel === "push" && prefs?.push_enabled === false) {
  return { skipped: true, reason: "push_disabled" };
}
const col = CATEGORY_COLUMN[params.type];
if (col && prefs && prefs[col] === false) {
  return { skipped: true, reason: `category_disabled:${col}` };
}
Enter fullscreen mode Exit fullscreen mode

Two small decisions worth calling out:

  • The map accepts aliases. reminder and reminders both land on the reminders column, so a caller can't skip the gate because of a plural.
  • Defaults live in the schema, not the code. In notification_preferences, marketing defaults to false and the other categories default to true. Marketing is opt-in without any special-casing in TypeScript.

A skipped notification is never inserted. Otherwise the row goes in as pending with next_attempt_at = now(), a created event is logged, and enqueue() immediately calls attemptSend() on it. The cron job is the retry path, not the normal path.

3. One attempt: send, then record what happened

attemptSend() is the heart of it. It loads the row, bails out early if the row is already finished, gathers the user's active tokens, and hands everything to a channel:

if (n.status === "sent" || n.status === "delivered" || n.status === "dead") {
  return { skipped: true };
}

const result = await channelFor(n.channel).send(
  { id: n.id, title: n.title, body: n.body, data: n.data },
  { tokens }
);

// Permanently invalid tokens → deactivate.
if (result.invalidTokens.length > 0) {
  await supabase
    .from("device_tokens")
    .update({ active: false, updated_at: new Date().toISOString() })
    .in("token", result.invalidTokens);
}
Enter fullscreen mode Exit fullscreen mode

Then exactly one of three things happens to the row: it's marked sent, it's moved to dead, or it gets a new next_attempt_at. The backoff schedule is a plain array:

// Retry schedule: Attempt 1 immediate, then 30s, 2m, 10m, 1h. After the 5th → dead-letter.
export const BACKOFF_MS = [0, 30_000, 120_000, 600_000, 3_600_000];
export const MAX_ATTEMPTS = BACKOFF_MS.length;

const attempt = n.attempt_count + 1;
// ...success → status "sent"; attempt >= MAX_ATTEMPTS → status "dead"
const delay = BACKOFF_MS[attempt] ?? BACKOFF_MS[BACKOFF_MS.length - 1];
const nextAttemptAt = new Date(Date.now() + delay).toISOString();
Enter fullscreen mode Exit fullscreen mode

The indexing is easy to misread, so here it is spelled out. The first attempt fails with attempt === 1, so the wait is BACKOFF_MS[1], 30 seconds. Index 0 is never used as a delay; it documents that attempt one is immediate. After the fifth failure, attempt >= MAX_ATTEMPTS and the row goes to dead, with a failed event whose meta records the reason and a console.error line for the logs.

One thing that took a moment to internalize: next_attempt_at is a "not before", not a timer. Nothing wakes up at that instant. The row simply becomes eligible, and the next run of the drain job (section 6) picks it up. So the real retry delay is the backoff or the time until the next drain, whichever is longer. If you copy this pattern, the cron schedule is as much a part of your retry policy as the array is.

4. Channels: the dispatcher doesn't know it's doing push

The dispatcher only talks to an interface:

export interface NotificationChannel {
  send(notification: OutboundNotification, recipient: DeliveryRecipient): Promise<DeliveryResult>;
}

export interface DeliveryResult {
  /** True if at least one endpoint accepted the message. */
  success: boolean;
  /** Tokens FCM rejected as permanently invalid (deactivate them). */
  invalidTokens: string[];
  error?: string;
}
Enter fullscreen mode Exit fullscreen mode

channelFor(n.channel) returns a PushChannel for push (and as the default). EmailChannel and SmsChannel exist as stubs that throw, so the channel column and the dispatcher are already shaped for them.

The push channel does one job that's easy to forget: FCM's data payload must be a flat string-to-string map, and our data column is JSONB. So the channel flattens it and always adds the notification's own id:

// FCM data must be string→string; include notification_id so the client can ack/deep-link.
const data: Record<string, string> = { notification_id: n.id };
for (const [k, v] of Object.entries(n.data ?? {})) {
  data[k] = typeof v === "string" ? v : JSON.stringify(v);
}
Enter fullscreen mode Exit fullscreen mode

That notification_id is what makes read receipts possible later. It's also why "success" means at least one device accepted the message. A user with a phone and a tablet shouldn't see a retry because one of them is offline.

5. FCM: per-token results and token hygiene

The FCM wrapper uses sendEachForMulticast from the Firebase Admin SDK, which returns one result per token instead of a single pass/fail:

const INVALID_TOKEN_CODES = new Set([
  "messaging/registration-token-not-registered",
  "messaging/invalid-registration-token",
  "messaging/invalid-argument",
]);

const res = await messaging.sendEachForMulticast({
  tokens,
  notification: { title: msg.title, body: msg.body },
  data: msg.data ?? {},
  android: { priority: "high" },
});
return res.responses.map((r, i) => ({
  token: tokens[i],
  success: r.success,
  error: r.error?.message,
  invalidToken: r.error ? INVALID_TOKEN_CODES.has(r.error.code) : false,
}));
Enter fullscreen mode Exit fullscreen mode

Per-token results are what feed the active = false update in section 3. Tokens go stale all the time (reinstalls, wiped phones, revoked permissions). Flagging them instead of deleting them keeps a history, and the partial index on active tokens means dead ones cost nothing on the read path.

Two setup details: the admin app is initialized lazily from a service-account JSON env var and reused while the function instance stays warm (admin.apps.length > 0), and every route that touches it declares export const runtime = "nodejs", because firebase-admin needs the Node runtime, not the Edge runtime.

6. The drain: a cron route that only does catch-up

Retries are picked up by a Vercel Cron job that calls a GET route. Vercel sends Authorization: Bearer $CRON_SECRET, and the route refuses anything else, including the case where the secret isn't configured at all:

const auth = request.headers.get("authorization");
if (!process.env.CRON_SECRET || auth !== `Bearer ${process.env.CRON_SECRET}`) {
  return NextResponse.json({ error: "Unauthorized" }, { status: 401 });
}
Enter fullscreen mode Exit fullscreen mode

That !process.env.CRON_SECRET || guard matters. Without it, a missing env var turns the check into auth !== "Bearer undefined", which anyone can satisfy.

The drain itself is short. It selects due rows oldest-first, up to a limit, and runs each through the same attemptSend() used by enqueue():

const { data: due } = await supabase
  .from("notifications")
  .select("id")
  .in("status", ["pending", "retrying"])
  .lte("next_attempt_at", new Date().toISOString())
  .order("next_attempt_at", { ascending: true })
  .limit(limit);
Enter fullscreen mode Exit fullscreen mode

That query matches the partial index from section 1 exactly: same status filter, same sort column. Because there's only one send path, backoff and dead-lettering behave the same whether a row is on its first attempt or its fourth.

7. Read receipts through RLS, and metrics in one view

The server writes notifications with Supabase's service-role client, which bypasses row-level security. Signed-in users get read-only access to their own rows:

create policy notifications_owner_select on public.notifications
    for select using (auth.uid() = user_id);
Enter fullscreen mode Exit fullscreen mode

But clients do need to report "delivered", "opened" and "read". Rather than open up UPDATE, those go through three SECURITY DEFINER functions that can touch exactly one timestamp on exactly one row the caller owns:

create or replace function public.mark_notification_opened(p_id uuid)
returns void language plpgsql security definer set search_path = public as $$
begin
    update public.notifications
       set opened_at = coalesce(opened_at, now())
     where id = p_id and user_id = auth.uid();
    if found then
        insert into public.notification_events(notification_id, user_id, event)
        values (p_id, auth.uid(), 'opened');
    end if;
end $$;

grant execute on function public.mark_notification_opened(uuid) to authenticated;
Enter fullscreen mode Exit fullscreen mode
  • where ... user_id = auth.uid() means a guessed id from another account matches nothing.
  • coalesce(opened_at, now()) keeps the first open time when a client acks twice.
  • if found only logs an event when a row really matched.
  • set search_path = public pins name resolution, which the Postgres docs recommend for SECURITY DEFINER functions, since they run with the owner's privileges.

The delivered variant also moves status to delivered, but only from pending, sent or retrying, so a late ack can't revive a dead row.

Every step also lands in an append-only notification_events table, and a view turns the timestamp columns into per-type rates with Postgres FILTER clauses:

round(count(*) filter (where opened_at is not null)::numeric
      / nullif(count(*) filter (where delivered_at is not null), 0), 4) as open_rate,
Enter fullscreen mode Exit fullscreen mode

The nullif(..., 0) returns NULL instead of dividing by zero for a type that hasn't delivered anything yet.

8. Broadcasts and segments

Sending to a group goes through broadcast(), which resolves a segment to user ids and calls enqueue() for each of them, so every user's preferences still apply. Segments are one const object that serves as the allowlist, the type and the human-readable label:

export const SEGMENTS = {
  all: "All users with an active device",
  inactive_7d: "No session in the last 7 days",
  never_sessioned: "Has the app but has never completed a session",
  mobility_test_due: "Last mobility test was more than 30 days ago",
  // ...
} as const;

export type Segment = keyof typeof SEGMENTS;

export function isValidSegment(value: unknown): value is Segment {
  return typeof value === "string" && value in SEGMENTS;
}
Enter fullscreen mode Exit fullscreen mode

The admin route validates the request body with isValidSegment(), returns the valid keys in its 400 error, and echoes SEGMENTS[segment] in the response, so whoever triggered the send can see who it targeted in plain words. Adding a segment means adding one key and one case. The switch over Segment lets TypeScript flag a key that has no query yet.

Every segment starts from users with at least one active device token, since there's no point resolving people we can't reach. Most are then simple set differences. inactive_7d is "has ever completed a session" minus "completed one in the last 7 days", so someone who has never done a session isn't counted as lapsed. They have their own never_sessioned segment.

broadcast() then enqueues in batches of 25 with Promise.all, and catches errors per user so one bad row doesn't sink the batch:

batch.map((userId) =>
  enqueue({ userId, ...params }).catch((e) => {
    console.error(`[broadcast] enqueue failed for ${userId}:`, e);
    return { error: "exception" };
  })
)
Enter fullscreen mode Exit fullscreen mode

The result is three counters, enqueued, skipped and errors, where "skipped" means the user's own preferences said no.

What we'd tell someone building the same thing

  • If you already have Postgres, try a table before adding a queue service. Status, attempts and next_attempt_at on the row, plus a partial index on in-flight work, covers retries and dead letters for a small team.
  • Gate on preferences before inserting. Put the defaults in the schema so opt-in categories stay opt-in.
  • Keep one send path. The first attempt and every retry run the same function.
  • Treat next_attempt_at as eligibility. Your drain schedule sets the real floor on retry latency.
  • Ask FCM for per-token results. Deactivate tokens that are permanently invalid, and count a send as a success if any device got it.
  • Let clients write through narrow SECURITY DEFINER functions, not broad UPDATE policies.

If you want to see what all this is in service of, MotionLab builds stretching sessions around how your body feels today, and the blog has the mobility side of things.

How do you handle notification retries? Postgres table, a hosted queue, or something like pg-boss? We'd like to hear what's worked for you in the comments.

This article was drafted with AI assistance and reviewed by the MotionLab team.

Top comments (0)