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”.
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.
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.
- 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
EventSourcetreats a non-200 response as fatal. The client must catch theerrorevent withreadyState === 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.
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.