Chat systems are one of the harder distributed systems problems disguised as a simple product: send a message, the other person sees it. Underneath is a persistent-connection management problem at massive scale, a message-ordering problem across unreliable networks, and a delivery-guarantee problem that has to survive users going offline mid-conversation for days.

Problem

We’re designing a chat system supporting one-on-one and small group conversations (up to a few hundred members), with text messages, delivery/read receipts, typing indicators, and presence (online/offline/last-seen). Users connect from mobile and web clients that go on and off networks constantly — the system has to make a message sent while the recipient was offline show up correctly the moment they reconnect, in the right order, exactly once.

The core tension is that chat wants two things that are in some tension with each other: low-latency real-time delivery when both parties are online, and durable, ordered, exactly-once-feeling delivery when they’re not. The architecture is largely about serving both without one compromising the other.

Requirements

Functional

  • Send and receive messages in 1:1 and group conversations in near real time.
  • Persist message history, retrievable on new device login or app reinstall.
  • Delivery receipts (sent, delivered, read) per message per recipient.
  • Presence: online/offline status and last-seen timestamp.
  • Typing indicators (ephemeral, no persistence needed).
  • Multi-device: a user can be logged in on phone and web simultaneously and see the same conversation state on both.

Non-functional

  • Message delivery latency under 200ms p99 when both parties are online and connected.
  • At-least-once delivery, with client-side dedup to present exactly-once to the user.
  • Per-conversation ordering guarantee: messages in a single conversation appear in the same order to every participant, even if the underlying transport reorders them.
  • Durable: a message accepted by the server is never lost, even if the recipient is offline for weeks.
  • Support millions of concurrent persistent connections without the connection layer becoming the bottleneck.

Capacity Estimation

Assume 50M daily active users, each sending an average of 40 messages/day (heavily used group chats push this number up, but 40 is a reasonable blended average), and 40% of DAU concurrently connected at peak (evening hours, global user base smooths this somewhat).

  • Messages/day: 50M × 40 ≈ 2B outbound sends/day. Each message fans out to an average of 3 recipients (mix of 1:1 and small groups), giving roughly 6B delivery events/day — close to our stated 15B/day scale target once read receipts and typing/presence events are included as messages on the same pipe.
  • Write QPS (average): 2,000,000,000 / 86,400 ≈ 23,000/s for message sends alone; peaking around 3x during evening hours ≈ 70,000/s.
  • Concurrent connections at peak: 50M × 0.4 ≈ 20M… at a more conservative sustained-peak estimate of 2M concurrent (matching our stated scale, accounting for a global, time-zone-distributed user base rather than everyone online at once).
  • Per-connection server cost: with a well-tuned event-loop server (Netty, or Go’s netpoller), each idle WebSocket connection costs roughly 20-30 KB of memory (buffers, TLS session state, connection metadata). 2M connections × 25 KB ≈ 50 GB of connection memory, spread across a fleet — comfortably achievable with a few dozen mid-sized connection-gateway hosts, not thousands.
  • Storage per message: message id, conversation id, sender id, body (avg 100 bytes for text), timestamp, status flags ≈ 250 bytes.
  • Daily storage: 2B × 250 bytes ≈ 500 GB/day.
  • Annual storage: 500 GB × 365 ≈ 182 TB/year, before replication — this is the number that pushes message storage toward a horizontally-sharded wide-column store rather than a single relational database from day one.
Back-of-the-envelope
Messages sent / day
2B
50M DAU x 40 msgs/day
Delivery events / day
~15B
fan-out + receipts + presence
Send QPS (peak)
~70,000/s
3x evening peak
Concurrent connections
~2M
~25KB/conn ≈ 50GB fleet-wide
Message storage / year
~182 TB
250 bytes/message, pre-replication

API Design

Chat is inherently bidirectional, so the primary interface is a persistent WebSocket connection carrying typed frames, backed by a REST API for history and account operations.

POST /v1/connect HTTP/1.1
Upgrade: websocket
Authorization: Bearer <session-token>

HTTP/1.1 101 Switching Protocols
Upgrade: websocket
Connection: Upgrade
// Client -> Server frame (send message)
{
  "type": "message.send",
  "clientMsgId": "c_8a1f-local-uuid",
  "conversationId": "conv_7f21",
  "body": "running 10 late, sorry",
  "sentAt": "2026-08-18T19:04:11Z"
}
// Server -> Client frame (ack)
{
  "type": "message.ack",
  "clientMsgId": "c_8a1f-local-uuid",
  "serverMsgId": "msg_9e102c",
  "sequence": 4821,
  "status": "sent"
}
// Server -> Client frame (incoming message, fanned out to recipient)
{
  "type": "message.new",
  "serverMsgId": "msg_9e102c",
  "conversationId": "conv_7f21",
  "senderId": "usr_44a1",
  "sequence": 4821,
  "body": "running 10 late, sorry",
  "sentAt": "2026-08-18T19:04:11Z"
}
GET /v1/conversations/conv_7f21/messages?before=4821&limit=50 HTTP/1.1

HTTP/1.1 200 OK
Content-Type: application/json

{
  "messages": [
    { "serverMsgId": "msg_9e099a", "sequence": 4820, "senderId": "usr_12b0", "body": "on my way", "sentAt": "2026-08-18T19:02:50Z" }
  ],
  "hasMore": true
}

clientMsgId is generated client-side before the network round trip and is the basis for client-side idempotency — if the client doesn’t see an ack within a timeout, it resends with the same clientMsgId, and the server treats a duplicate as a no-op that just re-returns the original ack.

Data Model

CREATE TABLE conversations (
  id          BIGINT PRIMARY KEY,
  type        VARCHAR(8) NOT NULL,   -- 'direct' or 'group'
  created_at  TIMESTAMPTZ NOT NULL DEFAULT now()
);

CREATE TABLE conversation_members (
  conversation_id BIGINT NOT NULL REFERENCES conversations(id),
  user_id         BIGINT NOT NULL,
  joined_at       TIMESTAMPTZ NOT NULL DEFAULT now(),
  last_read_seq   BIGINT NOT NULL DEFAULT 0,
  PRIMARY KEY (conversation_id, user_id)
);

CREATE TABLE messages (
  conversation_id BIGINT      NOT NULL,
  sequence        BIGINT      NOT NULL,
  id              BIGINT      NOT NULL,
  sender_id       BIGINT      NOT NULL,
  client_msg_id   VARCHAR(64) NOT NULL,
  body            TEXT        NOT NULL,
  sent_at         TIMESTAMPTZ NOT NULL,
  PRIMARY KEY (conversation_id, sequence)
);

CREATE UNIQUE INDEX idx_messages_dedup ON messages (conversation_id, sender_id, client_msg_id);

CREATE TABLE delivery_receipts (
  message_id  BIGINT      NOT NULL,
  user_id     BIGINT      NOT NULL,
  status      VARCHAR(16) NOT NULL,  -- 'delivered' or 'read'
  updated_at  TIMESTAMPTZ NOT NULL,
  PRIMARY KEY (message_id, user_id)
);

The messages table is keyed by (conversation_id, sequence), not a globally unique id alone — this is the load-bearing design decision. Partitioning storage by conversation_id means every conversation’s messages live together and its sequence column gives a cheap, per-partition total order, so “give me the next 50 messages after sequence 4820” is a simple range scan on one shard, never a cross-shard merge.

Field Purpose
sequence Monotonic per-conversation counter — the ordering primitive the whole system relies on
client_msg_id unique index Client-side idempotency; retried sends collapse to one row
conversation_members.last_read_seq Cheap “unread count” computation without scanning delivery_receipts

High-Level Architecture

High-level architecture

Clients hold a persistent WebSocket to a stateless connection gateway. The gateway doesn’t know how to route a message on its own — it hands off to a sequencer (which assigns the per-conversation order and persists), and a fan-out service that looks up which gateway instances hold connections for the recipients (via the routing/presence service) and pushes the message frame to them. If a recipient has no active connection, fan-out hands off to the push notification bridge instead.

Deep Dive

Connection management at scale

2 million concurrent WebSocket connections cannot live on one host. We run a fleet of stateless connection gateway instances behind a layer-4 load balancer with sticky routing. Each gateway instance holds an in-memory map of userId -> local connection, and registers userId -> gatewayInstanceId in a shared, low-latency store (Redis) so any other service can answer “which gateway, if any, holds this user’s connection right now” in a single lookup.

This registry is the crux of routing: sending a message doesn’t require knowing where the recipient is connected at send time — the fan-out service queries the registry per recipient, per message.

The sequencer and per-conversation ordering

Every conversation needs a single point that assigns monotonically increasing sequence numbers to its messages, because two clients racing to send at the same instant need a deterministic, agreed-upon order that every participant’s client will render identically. We shard this by conversation_id — each conversation is deterministically owned by one sequencer shard (consistent hashing), so within a shard, assigning the next sequence number is a simple, uncontended increment (SELECT nextval equivalent, or an atomic Redis INCR backed by periodic snapshot to the durable store).

Decision

Shard the sequencer by conversation_id, not globally

Consistent hashing of conversation_id to sequencer shard, each shard owning an independent counter space

A single global sequence counter would need every message in the entire system to serialize through one component — an obvious bottleneck at 70,000 sends/sec peak. Sharding by conversation means unrelated conversations never contend with each other, and the sequencer’s total throughput scales horizontally with shard count. The cost is that sequence numbers are only meaningful within a conversation, never comparable across conversations — which is fine, because nothing in the product needs a global ordering (see Requirements).

Delivery guarantees and client-side dedup

The server guarantees at-least-once delivery to a connected client — a message might be delivered twice if, say, the client’s ack gets lost and the gateway retries the push. Clients dedup on serverMsgId (or clientMsgId for their own sent messages) before rendering, which is why every frame carries a stable id rather than relying on arrival order or a lack of duplicates.

For offline delivery: the message is sequenced and persisted regardless of recipient connectivity (this durability write is the point we consider it “sent”) → fan-out checks the routing registry → if the recipient is offline, nothing more happens in real time, but the message is durably in messages. On reconnect, the client sends its last known sequence number per conversation, and the server replies with everything after that, in order. This makes “catch up after being offline” and “receive a live message” the same code path — a range read from a known point — rather than two mechanisms to keep in sync.

Read receipts and typing indicators

Read receipts are just another message type flowing through the same pipe (receipt.read frames), but they don’t need the same durability treatment as message content — a lost read receipt is a UX inconsistency (someone doesn’t see a checkmark update promptly), not a lost message. We write them to delivery_receipts asynchronously and batch/debounce on the client (don’t fire a read receipt per message if the user just scrolled through fifty of them in one motion; batch into one receipt for “read up to sequence N”).

Typing indicators skip persistence entirely — they’re pure pub/sub through the routing layer with no database write at all, TTL’d at a few seconds on the receiving client so a dropped “stopped typing” event self-heals instead of leaving a stuck “is typing…” indicator.

Multi-device fan-out

A user logged in on both phone and web is really two connections under one userId. The routing registry maps userId -> [gatewayInstance, connectionId] as a set, not a single value, and fan-out pushes to every active connection for that user. Read state has to reconcile across devices too — reading a conversation on web should clear the unread badge on mobile — which is handled by broadcasting receipt updates to all of a user’s own connections, not just the conversation’s other participants.

Failure Handling

Failure Detection Response
Connection gateway instance crashes Health check fails, LB deregisters All connections on that instance drop; clients auto-reconnect with backoff, resume from last known sequence per conversation
Sequencer shard unavailable Write timeout on sequence assignment Sender’s client sees a send timeout, retries with same clientMsgId; no message is lost because none was ever assigned a sequence, so there’s nothing to reconcile
Routing registry (Redis) unavailable Timeout on lookup Fan-out falls back to broadcasting to all gateway instances (each checks its local connection map) — more expensive, but correct; registry unavailability degrades performance, not correctness
Push notification bridge down Delivery failure from provider Message stays durably queued; delivered on next successful reconnect regardless — push is a convenience wake-up mechanism, not the delivery guarantee itself
Network partition between client and gateway Missed WebSocket ping/pong heartbeat Gateway closes the connection after a grace period, deregisters from routing; client reconnects and catches up via sequence-based range read
Duplicate message delivery Client detects repeated serverMsgId Silently deduped client-side; no server action needed, this is expected at-least-once behavior

Scaling

Connection gateway scaling

Stateless with respect to routing decisions (the registry, not local memory, is authoritative for cross-instance lookups), so horizontal scaling is adding gateway instances behind the load balancer. The practical ceiling per instance is bound by file descriptor limits and memory, not CPU, at these message rates — a well-tuned instance handles 50-100K idle connections comfortably.

Sequencer and message store scaling

Both shard on conversation_id and scale by adding shards via consistent hashing — the same story as the connection registry. A very high-traffic conversation can become a hot shard; the mitigation is capping practical group size rather than sub-sharding a single conversation’s sequence space, which would break the ordering guarantee that makes this design simple.

Fan-out amplification

A message to a 300-person group is 300 individual routing lookups and pushes from one write. At scale this is the dominant cost, not the write itself. We batch routing-registry lookups (one multi-get for all 300 recipients rather than 300 round trips) and let the fan-out service push in parallel across gateway instances, since each recipient’s push is independent once we know their location.

Message store growth

At 182 TB/year pre-replication, message storage moves to a wide-column store (Cassandra, ScyllaDB, or a similarly sharded system) rather than a single relational cluster, partitioned by conversation_id exactly as the logical schema suggests. Old, inactive conversations’ messages can move to cheaper cold storage after a threshold of inactivity, fetched on-demand if a user ever scrolls back far enough to need them.

Trade-off

WebSocket persistent connections

  • True bidirectional push, lowest latency for server-initiated messages
  • Higher per-connection resource cost, needs sticky routing infrastructure
  • Requires a connection-aware fan-out/routing layer

HTTP long-polling

  • Works everywhere without special infrastructure, simpler load balancing
  • Higher latency (poll interval) and higher request overhead per message
  • Doesn't scale well past moderate concurrency due to repeated connection setup

Recommendation — WebSockets for the primary client connection given our sub-200ms latency requirement and scale; long-polling remains a useful fallback for clients or networks that block WebSocket upgrades.

Observability

  • Connection count and churn rate per gateway instance — a spike in reconnects across the fleet (not just one instance) usually means a client-side release regression or a network provider issue, not a server problem, and this distinction saves an on-call engineer from chasing the wrong root cause.
  • End-to-end delivery latency (sequenced timestamp to recipient-ack timestamp), sampled and broken out by “both online” versus “recipient reconnect catch-up,” since these are fundamentally different latency budgets and averaging them together hides regressions in either one.
  • Sequencer shard lag — the gap between “message accepted” and “sequence assigned and persisted” per shard, which surfaces hot-shard problems (an unusually large or active group chat) before they show up as user-visible delivery delay.
  • Routing registry hit rate — how often fan-out finds a live connection versus falls back to push notification, which is a direct proxy for how much load is landing on the push provider and whether that capacity is provisioned correctly.

Every message frame carries a trace id threaded from client send through sequencing, persistence, and fan-out, so a single slow or lost message can be traced end-to-end across every hop rather than reconstructed from separate service logs.

Trade-offs

The central trade-off in this design is per-conversation sequencing versus global ordering. Sharding the sequencer by conversation_id is what makes the system horizontally scalable at all — a global sequence counter would cap total system throughput at whatever one serialized counter can sustain, nowhere near 70,000 sends/sec. The cost is that cross-conversation ordering (e.g., “which of these two messages, in two different chats, happened first”) is not something the system can answer precisely, only approximately via wall-clock timestamps. This is the right trade because no chat product actually needs that guarantee — users reason about order within a conversation, never across them.

The second trade-off is at-least-once delivery with client-side dedup versus attempting exactly-once at the protocol level. True exactly-once across an unreliable network and a client that can crash mid-processing is not achievable without unbounded coordination cost. At-least-once plus idempotent client rendering gets a user-perceived exactly-once experience at a fraction of the engineering cost, and this pattern — push the dedup responsibility to the edge, where an id is already available and cheap to check — recurs everywhere in this system, from sequencer to receipts to reconnect catch-up.

Final Architecture

Final architecture with sharded sequencing and multi-device fan-out

The finished system’s core insight is matching guarantee strength to what each piece of data actually needs: durable, per-conversation-ordered storage for message content; best-effort, ephemeral pub/sub for typing indicators; asynchronous, debounced writes for receipts. Sharding both the sequencer and the message store by conversation_id keeps the system horizontally scalable without ever requiring cross-shard coordination, and pushing idempotency to the client edge via stable message ids turns unreliable at-least-once delivery into a reliable-feeling product experience.