WebSocket Message Only Received by Users on the Same Server: The Pub/Sub Backplane

Realtime · Intermediate · 7 min read · published

This article was written by Claude (Anthropic) and published automatically.

What this solves: Your chat or presence feature worked on one instance and broke the moment you scaled to three. Here's the backplane architecture that fixes cross-node delivery without melting Redis.

The Forces at Play

You ship a chat feature, it works perfectly in dev, and then you scale to three pods. Now a WebSocket message is only received by users on the same server — the customer sees the agent's reply about a third of the time, and nobody can reproduce it locally. Nothing is broken in the usual sense; the code is doing exactly what you wrote.

The tension is this: a WebSocket connection is stateful and pinned. Once the TCP socket is established, that user lives on exactly one process for the life of the connection. But the events they need are produced anywhere — by another user on a different pod, by a background job, by a webhook hitting an HTTP handler that has no sockets at all. So you have a routing problem: producers know a room ID, but delivery requires knowing which process holds the socket.

The naive answers each fail in a specific way:

The pattern that survives is a pub/sub backplane with per-room channels plus a durable, replayable log.

The Shape

flowchart TB
    C1["Client A<br/>room r88"] --> LB{{"Load Balancer<br/>(no stickiness needed)"}}
    C2["Client B<br/>room r88"] --> LB
    C3["Client C<br/>room r12"] --> LB

    LB --> N1["WS Node 1<br/>subs: room.r88"]
    LB --> N2["WS Node 2<br/>subs: room.r88, room.r12"]

    API["HTTP API / Worker<br/>(holds zero sockets)"] --> PUB["publish(roomId, event)"]
    N1 -->|"client sends msg"| PUB
    N2 -->|"client sends msg"| PUB

    PUB --> LOG[("Durable Event Log<br/>per-room, seq-numbered")]
    LOG -->|"assign seq"| BUS["Broker: topic room.{id}"]

    BUS -.->|"only rooms it holds"| N1
    BUS -.->|"only rooms it holds"| N2

    N1 --> C1
    N2 --> C2
    N2 --> C3

    N1 <-->|"register / TTL heartbeat"| PRES[("Presence Store<br/>user -> node, expiring")]
    N2 <-->|"register / TTL heartbeat"| PRES

    N1 -->|"resume from last_seq"| LOG

The key structural move: nodes never address each other. They address rooms. Subscription membership at the broker becomes the routing table, and it's maintained as a side effect of clients joining and leaving.

How Data Flows Through It

Agent on Node 1 sends a message in room r88; the customer's socket lives on Node 2.

  1. Ingress. Node 1 receives the frame, authorizes that this user may write to r88. Authorization happens here and nowhere downstream.
  2. Append. Node 1 appends the event to the durable log for r88 and gets back a monotonic seq (e.g. a Redis Stream ID, or a Postgres sequence per room). This is the ordering authority — not the broker.
  3. Publish. Node 1 publishes the enriched event to topic room.r88.
  4. Fan-out. The broker delivers to exactly the nodes subscribed to room.r88 — Node 1 and Node 2. Node 3, holding no r88 sockets, receives nothing and burns no CPU.
  5. Local delivery. Each node looks up its in-process Map<roomId, Set<Socket>> and writes the frame. Node 1 also echoes to the sender, which now sees the server-assigned seq and can reconcile its optimistic local copy.
  6. Reconnect. The customer's phone sleeps; the socket dies. On reconnect they land on Node 3 (no stickiness) and send last_seq: 4471. Node 3 reads 4472..HEAD from the log, flushes the gap, then attaches the live subscription — and dedupes by seq to cover the overlap window.
// The one contract every producer and consumer agrees on.
interface RoomEvent {
  roomId: string;
  seq: string;        // assigned by the log, not the sender — the ordering authority
  type: 'message' | 'typing' | 'presence';
  senderId: string;
  payload: unknown;
  publishedAt: string;
}

// Client resume handshake
interface Subscribe { roomId: string; lastSeq?: string }

Note that typing events skip the log entirely — they're publish-only, because a typing indicator that arrives 30 seconds late is worse than one that never arrives.

What Each Piece Owns

Load balancer — owns TCP distribution of the upgrade request. Deliberately does not own session affinity. If your design needs stickiness for correctness, the backplane isn't doing its job.

WebSocket node — owns socket lifecycle, authorization on ingress, the local roomId -> sockets map, and backpressure on slow clients. It deliberately owns no durable state: killing any node must lose nothing but connections.

Durable event log — owns ordering and replay. It does not own delivery; it never pushes to a client.

Broker (pub/sub) — owns live fan-out and, via subscription membership, the routing table. It deliberately does not own durability. Treat every published message as best-effort — the log is your safety net.

Presence store — owns "who is online, on which node," with short TTLs refreshed by heartbeat so a SIGKILLed pod expires itself. It does not own message routing; presence is a read model for UI and targeted single-user pushes, not a prerequisite for room delivery.

Where It Breaks Down

Subscription churn is the first bottleneck, not connection count. A user scrolling a room list can cause thousands of SUBSCRIBE/UNSUBSCRIBE commands per second on a single-threaded Redis. Fix: keep room subscriptions alive on a lazy timer (unsubscribe 30s after the last local socket leaves) and batch subscribe calls.

Redis Cluster broadcasts pub/sub to every node. Classic PUBLISH in a cluster propagates across all shards, so sharding gives you zero fan-out relief. Use sharded pub/sub (SPUBLISH/SSUBSCRIBE, Redis 7+) so room.{id} hashes to one shard, or move to NATS/Kafka.

Hot rooms break the per-room model. A 500,000-viewer live event means one channel, one shard, one hotspot. At that point you split the room into room.r88.part{0..15} with clients hashed into partitions, and accept that per-partition ordering is all you get.

Slow consumers cause memory blowups, not drops. A mobile client on a bad network stops reading; the node's send buffer grows. Without a bound you OOM the pod and disconnect thousands of healthy users. Cap the per-socket queue and force a disconnect — the client will reconnect with last_seq and self-heal. This is the payoff of having a log.

Partial failure: if the broker is up but the log is down, you must fail the write. Publishing without a seq produces messages that can never be replayed or ordered — a silent, permanent data gap that surfaces weeks later as "the history is missing a message."

The operational burden is the replay path. Live delivery is exercised constantly; the resume-from-cursor branch only runs on reconnect, so bugs there hide for months. Add a synthetic client in staging that reconnects every 10 seconds and asserts zero gaps and zero duplicates.

When This Is Overkill

One process holding an in-memory Map<roomId, Set<Socket>> is genuinely correct — and dramatically simpler — up to tens of thousands of concurrent sockets on a single well-tuned node. A lot of production real-time features never need more.

Cheaper designs to reach for first:

The signals you've outgrown the single node, in order of urgency:

  1. You need two instances for availability, not throughput — that alone kills the in-memory map.
  2. A non-socket producer (webhook, cron, another service) needs to push to a user.
  3. You've added sticky sessions to make a bug go away.

Signal 3 is the loudest. Stickiness as a correctness mechanism means every deploy silently reshuffles your routing table.

Key takeaway: A WebSocket tier is only a router — the moment you run more than one node, message fan-out must live in a shared backplane with per-room channels and a replayable log, not in process memory.

Real-world challenge

You scaled a support-chat service from one instance to three behind an ALB. Now roughly two thirds of messages sent by an agent never appear for the customer, and vice versa — but sometimes everything works fine. Separately, users report that messages sent while their phone was locked are permanently missing after the socket reconnects. Logs show no errors on any node. How do you diagnose and fix this?

Diagnosis step 1 — is it cross-node? Correlate a failed delivery with which pod each socket landed on. Add the pod name to the connection log and to a x-served-by header on the upgrade response:

conn open room=r_88 user=agent_3 pod=ws-7f4c
conn open room=r_88 user=cust_1 pod=ws-91ab

If successful deliveries always share a pod and failures never do, the fan-out is in-process (Map<roomId, Set<Socket>>) with no backplane. Success rate ≈ 1/N is the fingerprint.

Fix 1 — add a backplane. Publishing goes to Redis/NATS on room.{id}; each node subscribes only for rooms it holds. Nodes no longer talk to each other directly.

async function publish(msg: RoomEvent) {
  const seq = await appendToLog(msg);        // durable, ordered
  await broker.publish(`room.${msg.roomId}`, { ...msg, seq });
}

Diagnosis step 2 — the reconnect gap. Pub/sub is fire-and-forget: anything published while the socket was down is gone. Confirm by checking whether the missing messages exist in your message table but were never re-sent.

Fix 2 — resume from a cursor. Client sends last_seq on reconnect; the node reads the gap from the durable log before attaching the live subscription, and dedupes by seq.

Sticky sessions are not the fix — they only hide the bug until a pod restarts and rebalances everyone.