Notification & Activity Feeds Permalink to this section

Part of Real-Time Application Patterns.

A notification feed is the feature where users notice a single missing item. A dashboard that is a second stale is invisible; a mention that never showed up, or a “your export is ready” that appeared twice, is a support ticket. Notifications are append-only data with per-user visibility, which puts them at the opposite end from dashboards: every event matters, order matters, each user sees a different stream, and a client that was disconnected must receive what it missed. This guide covers the full path — durable storage and per-user channels on the server, a resumable SSE stream using event ids as cursors, unread counts that stay consistent across tabs and devices, and the client store that de-duplicates what the network inevitably delivers twice. It is written for teams adding a bell icon, an activity sidebar or an inbox to an existing product, and the code works with any backend that has a database and a pub/sub bus.

How It Works Permalink to this section

The core rule is that the database is the source of truth and the stream is a notification of the database. Every notification is written to durable storage first, with a per-user monotonic id, and only then published to the user’s channel. The SSE stream is a tail of that table, and reconnecting is simply “tail from where I was”.

Write first, publish second, stream from the cursor Flow from a domain event through a durable notification write and a publish to the user's channel, to the SSE node and the browser, which returns its cursor on reconnect. Write first, publish second, stream from the cursor Domain event comment, mention, job done 1 INSERT row id per user, durable 2 PUBLISH user:42 channel 3 SSE node tail + live frames Browser Last-Event-ID The live publish is an optimisation; the table is the guarantee.
Because the row exists before the publish, a client that misses the publish still finds the notification when it reconnects and replays from its cursor.

On the wire each notification is one named event whose id is the per-user sequence:

retry: 5000

event: notification
id: 8812
data: {"id":8812,"kind":"mention","actor":"Priya","target":"PR #412","ts":"2026-09-18T09:12:03Z","read":false}

event: unread
data: {"count":3}

event: notification
id: 8813
data: {"id":8813,"kind":"export_ready","target":"Q3 report","ts":"2026-09-18T09:12:41Z","read":false}

Two event types share the stream. notification carries items and has an id, so it advances the browser’s cursor. unread carries a derived count and deliberately has no id field, because it is state rather than a list item — a reconnect should replay items, not old counts. The spec’s rule that an event without an id field leaves the last event id unchanged makes this separation free, as explained in event ID and retry mechanism design.

Server-Side Implementation Permalink to this section

The server has three responsibilities: persist, publish, and stream from a cursor with no gap between the replay and the live tail.

-- Per-user sequence: ids are dense and ordered within each user's feed.
CREATE TABLE notifications (
  user_id     bigint      NOT NULL,
  seq         bigint      NOT NULL,
  kind        text        NOT NULL,
  payload     jsonb       NOT NULL,
  created_at  timestamptz NOT NULL DEFAULT now(),
  read_at     timestamptz,
  PRIMARY KEY (user_id, seq)
);

-- Allocate the next seq atomically for one user.
CREATE TABLE notification_counters (user_id bigint PRIMARY KEY, last_seq bigint NOT NULL);
// notify.js — called by any part of the product that wants to notify a user.
export async function notify(userId, kind, payload) {
  const row = await db.tx(async (t) => {
    const { last_seq } = await t.one(
      `INSERT INTO notification_counters (user_id, last_seq) VALUES ($1, 1)
       ON CONFLICT (user_id) DO UPDATE SET last_seq = notification_counters.last_seq + 1
       RETURNING last_seq`, [userId]);
    return t.one(
      `INSERT INTO notifications (user_id, seq, kind, payload) VALUES ($1, $2, $3, $4)
       RETURNING seq, kind, payload, created_at`, [userId, last_seq, kind, payload]);
  });
  // Publish only after commit. If this publish is lost, replay still finds the row.
  await redis.publish(`user:${userId}:notifications`, JSON.stringify(row));
  return row;
}

The stream handler subscribes first, then replays, then drains anything that arrived during the replay — the same ordering that prevents gaps in every resumable stream:

app.get('/api/notifications/stream', requireSession, async (req, res) => {
  const userId = req.user.id;                         // never from the query string
  openStream(res);                                    // headers, retry, heartbeat

  let cursor = Number(req.get('Last-Event-ID') ?? req.query.after ?? 0);
  const early = [];
  let replaying = true;

  const sub = await bus.subscribe(`user:${userId}:notifications`, (msg) => {
    const n = JSON.parse(msg);
    if (replaying) early.push(n); else send(n);
  });

  function send(n) {
    if (n.seq <= cursor) return;                      // already delivered — drop the duplicate
    cursor = n.seq;
    res.write(`event: notification\nid: ${n.seq}\ndata: ${JSON.stringify(n)}\n\n`);
  }

  const missed = await db.any(
    `SELECT seq, kind, payload, created_at FROM notifications
      WHERE user_id = $1 AND seq > $2 ORDER BY seq LIMIT 500`, [userId, cursor]);
  missed.forEach(send);
  replaying = false;
  early.forEach(send);                                // anything published during the query

  res.write(`event: unread\ndata: ${JSON.stringify({ count: await unreadCount(userId) })}\n\n`);
  req.on('close', () => sub.unsubscribe());
});

The same algorithm in Python, using FastAPI with asyncpg and redis.asyncio, shows that nothing about it is Node-specific. The early-buffer is an asyncio.Queue that the subscription fills while the replay query runs:

# notifications_stream.py — FastAPI, asyncpg, redis.asyncio
import asyncio, json
from fastapi import Depends, Request
from fastapi.responses import StreamingResponse

@app.get("/api/notifications/stream")
async def notifications_stream(request: Request, user=Depends(current_user)):
    cursor = int(request.headers.get("last-event-id") or request.query_params.get("after") or 0)
    pubsub = redis.pubsub()
    await pubsub.subscribe(f"user:{user.id}:notifications")      # 1. subscribe first

    async def gen():
        nonlocal cursor
        try:
            yield "retry: 5000\n\n"
            rows = await pg.fetch(                                 # 2. replay after the cursor
                "SELECT seq, kind, payload, created_at FROM notifications "
                "WHERE user_id = $1 AND seq > $2 ORDER BY seq LIMIT 500", user.id, cursor)
            for r in rows:
                cursor = r["seq"]
                yield f"event: notification\nid: {r['seq']}\ndata: {json.dumps(dict(r), default=str)}\n\n"
            yield f"event: unread\ndata: {json.dumps({'count': await unread_count(user.id)})}\n\n"
            while not await request.is_disconnected():             # 3. live tail, cursor-filtered
                msg = await pubsub.get_message(ignore_subscribe_messages=True, timeout=20)
                if msg is None:
                    yield ": hb\n\n"
                    continue
                n = json.loads(msg["data"])
                if n["seq"] <= cursor:
                    continue                                       # already sent by the replay
                cursor = n["seq"]
                yield f"event: notification\nid: {n['seq']}\ndata: {msg['data'].decode()}\n\n"
        finally:
            await pubsub.unsubscribe()
            await pubsub.aclose()

    return StreamingResponse(gen(), media_type="text/event-stream",
                             headers={"Cache-Control": "no-cache", "X-Accel-Buffering": "no"})

Redis holds published messages for a subscribed client in its output buffer until they are read, so messages published during the replay query wait there and are drained by the live loop; the seq <= cursor check removes any overlap with the replayed rows. The FastAPI SSE implementation guide covers the worker and disconnect details this sketch leaves out.

The LIMIT 500 bounds replay after a very long absence. If a user comes back with thousands of missed items, the client does not need them streamed one frame at a time; the first page of the inbox is fetched over plain HTTP and the stream resumes from the newest id. Per-user channels for notification delivery covers the channel design and its failure modes at scale.

Subscribe, replay, then drain — no gap and no duplicate Sequence diagram showing the SSE node subscribing to the user channel, querying the table for rows after the cursor, buffering a publish that arrives during the query, then sending replayed rows followed by the buffered one. Subscribe, replay, then drain — no gap and no duplicate Browser SSE node Bus Database GET Last-Event-ID: 8810 SUBSCRIBE user:42 seq > 8810 publish 8813 (buffered) rows 8811, 8812, 8813 8811, 8812, 8813 buffered 8813 is dropped: seq ≤ cursor
A notification published while the replay query runs is held in the early buffer and sent after the replay; the seq check discards it if the query already returned it.

Client-Side Consumption Permalink to this section

The client store is keyed by notification id. That one decision makes duplicates harmless: an item delivered twice is written to the same key twice.

// notifications-store.js
export function createNotificationStore() {
  const items = new Map();            // seq → notification
  let unread = 0;
  const listeners = new Set();
  const emit = () => listeners.forEach((fn) => fn({ items: [...items.values()].sort((a, b) => b.seq - a.seq), unread }));

  const es = new EventSource('/api/notifications/stream', { withCredentials: true });
  es.addEventListener('notification', (e) => {
    const n = JSON.parse(e.data);
    const isNew = !items.has(n.seq);
    items.set(n.seq, n);
    if (isNew && document.visibilityState === 'visible') showToast(n);
    emit();
  });
  es.addEventListener('unread', (e) => { unread = JSON.parse(e.data).count; emit(); });

  return {
    subscribe(fn) { listeners.add(fn); fn({ items: [], unread }); return () => listeners.delete(fn); },
    close() { es.close(); },
  };
}

Toasts are shown only for items that are new to this store and only when the tab is visible, so a reconnect replay never produces a burst of popups. Building a notification bell with SSE turns this store into a complete component, and unread counts and read receipts covers keeping the count correct when several tabs and devices mark items read.

With several tabs open, each tab would otherwise hold its own stream. On HTTP/1.1 that eats into the six-connection limit; on any protocol it multiplies server load. Sharing one SSE connection across tabs moves the stream into a single leader tab or a SharedWorker.

Edge Cases & Network Interference Permalink to this section

Notification streams are mostly idle. That makes them the stream type most likely to be closed by an intermediary that thinks the connection is dead.

Why quiet streams fail differently Two panels contrasting the failure modes of a quiet notification stream with those of a busy data stream. Why quiet streams fail differently Quiet stream (notifications) idle timeouts close it dead sockets go unnoticed missed item found only on reload fix: heartbeat + watchdog Busy stream (dashboards) buffering delays frames slow clients back up render storms in the tab fix: coalesce + frame loop
A notification stream may go hours without an item. Every idle timeout between the browser and the server gets a chance to close it.
  • Idle timeouts. Load balancers, corporate proxies and mobile carriers close connections idle for 30 to 350 seconds. A comment line every 15–25 seconds keeps every hop active.
  • Half-open connections. A laptop that sleeps on one network and wakes on another may keep a socket that looks open but delivers nothing. The browser does not notice; a watchdog that reconnects after two missed heartbeats does.
  • Authentication expiry. A session cookie that expires while the stream is open does not affect the open stream, but the reconnect will fail with a 401, and EventSource treats a non-200 response as fatal. The client must catch the error event with readyState === CLOSED, refresh the session, and reopen.
  • Cross-tenant leakage. A channel name derived from request input rather than the authenticated identity lets one user subscribe to another’s notifications. Derive channel names on the server only.
  • Muted and snoozed items. Filtering must happen before the frame is written, in the projector, so preferences apply identically to replay and live delivery.

Mitigation checklist:

Performance & Scale Considerations Permalink to this section

Notification streams have a distinctive capacity profile: very many connections, very few messages. A product with 200,000 daily users may hold 60,000 concurrent streams that together receive a few hundred notifications per minute. The binding constraints are memory per idle connection and the pub/sub subscription count, not CPU or bandwidth.

Memory held by 60,000 idle notification streams Bar chart comparing total memory for sixty thousand idle streams across four server runtimes. Memory held by 60,000 idle notification streams Go (goroutine per stream) ~0.9 GB Node.js (event loop) ~1.6 GB Python asyncio ~2.4 GB Thread per connection ~60 GB (1 MB stacks) approximate resident memory for connections alone
Idle connections are cheap in event-loop and goroutine runtimes and expensive in thread-per-connection ones. Measure your own, but plan for the idle case, not the busy one.

Per-user channels in Redis cost one subscription per connected user per node. At tens of thousands of users per node, prefer a single pattern subscription per node or a sharded bus, and route messages to local connections in memory; the trade-offs are covered in per-user channels for notification delivery and sharded pub/sub in Redis Cluster.

Volume spikes are the third. A product event that notifies many users at once — a comment on a thread with 5,000 followers, an announcement to a whole workspace — produces thousands of inserts and publishes in a burst. Write those notifications through a queue with a bounded worker pool rather than inline in the request that caused them, and publish in batches:

// Fan-out worker: one job per (event, recipient batch), not per recipient inline.
queue.process('notify-followers', 8, async (job) => {
  const { event, recipientIds } = job.data;           // up to 500 recipients per job
  const rows = await insertNotifications(recipientIds, event);   // one multi-row INSERT
  const pipe = redis.pipeline();
  for (const r of rows) pipe.publish(`user:${r.user_id}:notifications`, JSON.stringify(r));
  await pipe.exec();                                  // one round trip for the batch
});

The originating request returns immediately, the database sees a few large inserts instead of thousands of small transactions, and recipients receive the notification within a second or two of each other.

Replay queries are the other hot spot. The (user_id, seq) primary key makes “rows after cursor” an index range scan, which stays fast at any table size. A deploy that reconnects every client at once turns into 60,000 of those queries in a few seconds, so spread reconnects with a jittered retry: value and consider skipping the query when the client’s cursor equals the counter’s last_seq.

Validation & Debugging Permalink to this section

Exercise the three paths that matter: live delivery, replay after a gap, and duplicate suppression.

# Live: open the stream, then trigger a notification from another terminal.
curl -sN -b session.txt https://app.example.com/api/notifications/stream

# Replay: reconnect claiming an older cursor and confirm the missed items arrive in order.
curl -sN -b session.txt -H 'Last-Event-ID: 8800' https://app.example.com/api/notifications/stream \
  | grep '^id:' | head

# Idle survival: leave a stream open for 10 minutes and confirm heartbeats keep arriving.
curl -sN -b session.txt https://app.example.com/api/notifications/stream | grep --line-buffered '^:' \
  | while read l; do echo "$(date +%T) $l"; done

In DevTools, the EventStream tab lists every event with its id. After toggling offline and back, the replayed ids should continue the sequence without a hole. Server logs should record one line per connection with the cursor it resumed from and how many rows it replayed:

{"evt":"notif_stream_open","user":42,"cursor":8800,"replayed":13,"replay_ms":4}

A spike in replayed after a deploy is normal. A steady stream of opens with replayed > 0 during normal operation means connections are dropping more often than they should — look for an intermediary with a short idle timeout.

Production Checklist Permalink to this section

Frequently Asked Questions Permalink to this section

Why not use a global auto-increment id instead of a per-user sequence?

A global id works as a cursor too, and it is simpler to allocate. A per-user sequence makes gaps detectable — the client can tell that 8812 follows 8811 — and keeps ids small. Either is correct as long as ids are monotonic within the user's feed.

Should notifications be delivered by Web Push instead of SSE?

They solve different problems. SSE delivers to an open page in real time; Web Push reaches a user whose page is closed. Most products use both: SSE for the in-app bell, and push or email for items that matter while the user is away.

What happens if Redis loses a publish?

Nothing visible, provided the notification was committed to the database first. The connected client misses the live frame, but the next reconnect replays from its cursor and finds the row. For faster repair, send a periodic cursor check that makes the client reconnect when the server's latest id is ahead of its own.

How long should notifications be kept for replay?

Keep the rows as long as the inbox shows them — often 30 to 90 days — because replay is just a query against the same table. The stream only replays a bounded page; older items are reached through the inbox's normal pagination.

Can activity feeds for a whole team use the same design?

Yes, with the channel and sequence keyed by team or project instead of user. Authorisation moves to connect time — check membership before subscribing — and a revoke message on the user's control channel closes streams for members who lose access.

How do I mark notifications as read across devices?

Write read state to the database and publish a read event on the same user channel. Every connected tab and device receives it and updates its local store and unread count. The unread-counts guide covers the race between marking read and new items arriving.

Deep Dives