System design practice

Chat

medium websocket pub-sub durability presence

Store before the ack, then deliver over pub/sub or push.

Solve it in your browser Read the lesson first

Design one-to-one messaging for a chat app. Phones keep a WebSocket open to the service while the app is in the foreground; when it is not, the only way to reach the user is a push notification.

Functional requirements

Use these use case and scenario names exactly: the traffic, requirements and tests in problem.proschi refer to them.

Scale

Constraints

What is given

problem.proschi declares the sender and the recipient and the external push provider (APNs and FCM; about 200 ms per notification), and holds the traffic, requirements and tests. Add the components, the connections and the two use cases.

How your design is checked

You write the design as text in Proschi. Tests run in your browser: a simulation of the traffic above checks latency, availability, cost and what happens when a machine fails. How the simulation works.

Start designing

Lesson · 11 min read

Learn it: Chat #

Chat: store before the ack, deliver after it #

One-to-one chat is a favourite interview problem because it mixes three ideas that each look simple alone: a durable write, a long-lived connection, and routing a message to whichever server the recipient happens to be connected to. Get the order wrong and you either lose messages or make every sender wait for the slowest phone on the network. This lesson builds the design in the right order.

What you'll learn #

The problem, explained #

Who uses it. 50 million daily active users on phones. While the app is in the foreground, it keeps a WebSocket open to the service. When it is in the background, the only way to reach the user is a push notification through APNs or FCM.

Functional requirements.

Non-functional requirements. p99 (the latency that 99% of requests beat) under 100 ms from send to ack, and under 200 ms for history. Sending available 99.95%. A message is never lost once the sender saw the ack. The ack never waits for delivery: neither for the push provider nor for the recipient's connection. After the ack, the service looks up presence first and then delivers the message. Any single machine can fail. At most $4,000 a month.

What is given, and why. given.proschi declares the sender, the recipient and the external push provider, with a capacity of 10k notifications a second and roughly 200 ms per call. Everything between them is yours to design.

What the tests check.

Back-of-the-envelope #

QuantityArithmeticResult
Sendsgiven10k rps
Online deliveries via pub/sub70% × 10k7k rps
Push notifications30% × 10k3k rps (provider limit 10k)
Presence lookupsone per send10k rps
Message writesone per send10k rps
Gateway work10k sends + 7k deliveries17k rps
History readsgiven2k rps
Messages per day, upper bound10k/s × 86,400 s864M
Storage per day (assume ~200 bytes each)864M × 200 Babout 170 GB/day, about 63 TB/year

The storage line is an upper bound (the peak rate held all day) with an assumed message size, but it shows the shape: data grows forever, writes are constant, and reads are for recent messages in one conversation.

Gateways do two jobs. A gateway receives the sender's message and delivers messages to the users connected to it. Every online delivery passes through a gateway a second time. In Proschi, a service replica handles 2k requests a second, so 17k rps needs 8.5 replicas at 100%. Divide by a 70% target and check that one replica fewer still stays under 100%.

The store must take 10k writes a second. A relational database in the model is single-primary: 5k writes a second per shard, whatever the replica count. 10k writes would saturate it unless you shard it. A partitioned store like Cassandra accepts writes on every replica, 20k a second each, so two replicas are lightly loaded. Real chat systems (Discord is the famous write-up) chose wide-column stores for exactly this append-heavy, partition-by-conversation pattern.

Latency. The ack path is load balancer → gateway → message store → ack: a 2 ms hop, a 10 ms gateway, a 5 ms write, plus queueing and tails. That fits 100 ms easily. Adding the push provider (about 200 ms) before the ack would not.

Connections, not just requests. The model counts requests per second. A real gateway is also sized by open connections and memory per connection. With tens of millions of users online at peak, you need enough gateway servers to hold them all; mention it, even though the simulation does not count it.

Cost. Service replicas $100, a load balancer $50, a Redis replica $150, a queue or broker replica $200, a Cassandra replica $500. With $4,000, every tier needs to be sized, not padded.

Concepts #

The ack is a promise: store first #

An ack (acknowledgement) tells the sender "your message is safe". The only honest way to say that is after the message is in durable, replicated storage. If you ack on receipt and store later, a gateway crash between the two loses messages the sender believes were sent, and nothing can recover them.

Once the message is stored, everything else can be retried: delivery to the recipient, the push notification, syncing to the sender's other devices. The store becomes the source of truth; connections and notifications are just ways of telling people to look.

The trade-off is that the ack waits for one database write. With a store designed for fast appends, that is a few milliseconds. When not to do this: ephemeral signals such as "typing…" indicators, which are fine to lose and should never touch the database.

title "Durable comments"

reader "Reader"        [Actor]
api    "Comments API"  [REST API]  x2
store  "Comment Store" [Cassandra] x2
fanout "Comment Bus"   [Kafka]     x2

reader -> api    : HTTPS
api    -> store  : write
api    -> fanout : publish

usecase "Post comment" {
  reader -> api    : POST /threads/5/comments
  api    -> store  : INSERT comment c_9
  store --> api    : ok
  api   --> reader : 201 {"id": "c_9"}
  api   ->> fanout : CommentPosted c_9
}

WebSockets and stateful gateways #

HTTP is request/response: the client asks, the server answers. For chat, the server needs to talk first. A WebSocket (RFC 6455) starts as an HTTP request, upgrades to a long-lived, two-way connection over TCP, and lets either side send messages at any time.

That changes the servers. A gateway that holds WebSockets is stateful: user Bob is connected to gateway 7, not to "the service". Losing gateway 7 drops Bob's connection (his app reconnects to another gateway), and anyone who wants to reach Bob must know he is on gateway 7. Load balancers must support long-lived connections, and deploys must drain connections gradually.

Alternatives: long polling (the client keeps a request open until the server has something, then re-opens it) works everywhere but costs a request per message; server-sent events push one way only. When not to use WebSockets: for occasional updates where a push notification or a periodic poll is enough.

Presence and pub/sub routing #

The sender is on gateway 3, the recipient on gateway 7. Two pieces get the message across:

If presence says nobody is connected, send a push notification instead. Presence can be slightly stale (a phone that just lost signal still looks online for a few seconds), which is fine because the message is already stored: the recipient fetches it on the next app open, and clients deduplicate by message id.

Trade-offs: presence is one more store to keep available, and very large group chats need a different design (fan-out to many gateways). Slack's real-time messaging architecture uses the same split: gateway servers hold the client connections, and separate servers track presence.

title "Live scores"

source   "Match Data"  [Third Party API]
feed     "Score Feed"  [Worker]          x2
registry "Who Watches" [Redis]           x2
bus      "Score Bus"   [NATS]            x2
edge     "Socket Edge" [WebSocket]       x3
fan      "Fan"         [Actor]

source -> feed     : webhook
feed   -> registry : lookup
feed   -> bus      : publish
bus    -> edge     : deliver
edge   -> fan      : WebSocket

usecase "Goal scored" {
  source    -> feed     : GoalScored match 12
  feed      -> registry : GET watchers:match12
  registry --> feed     : edge-2
  feed     ->> bus      : PUBLISH edge.2 goal
  bus      ->> edge     : goal (on edge-2)
  edge     ->> fan      : SCORE 1-0
}

Designing it step by step #

1. Scope. Ask: one-to-one only, or groups (one-to-one here)? Delivery guarantees (never lose after ack; at-least-once delivery, where a message may arrive twice but never zero times, with deduplication on the client is fine)? Ordering (per conversation, by time)? Multi-device (history must work on any device)? Read receipts, typing indicators, media (out of scope, but name them)? Confirm the numbers: 10k sends a second, 70% of recipients online, 2k history loads a second.

2. High-level design. Draw three paths.

Explain why history goes through a separate stateless API: it is ordinary request/response and should not occupy gateway capacity reserved for live traffic.

3. Deep dive.

4. Wrap-up. Walk the requirements: the ack follows the write; delivery is after the ack in both scenarios; no single replica anywhere; the budget fits. Then mention extensions: message ids generated with a time-ordered id scheme, delivery and read receipts as messages in the reverse direction, group chats with a per-group fan-out, end-to-end encryption, and syncing a user's other devices by having every device of a user subscribe through presence.

Common mistakes #

Delivering before the ack (wrong/deliver-before-ack). The gateway stores the message, then publishes it, waits for the recipient's connection to confirm, and only then acks the sender. It feels "more correct", because the ack now means "delivered". But the sender's latency is now tied to the recipient's phone, which might be on a slow mobile network, and a stuck connection stalls the sender. In the model the recipient is an actor with no latency, so p99 only rises from about 51 ms to about 64 ms and stays under the limit: the numbers alone do not catch it. The flow test does: it fails Delivery happens after the ack. A good reminder that the model is optimistic about clients you do not control.

Other classic mistakes.

In the interview #

Open with the ordering, because it is the crux: "Store, ack, then deliver. The ack is a durability promise, not a delivery promise." Then draw the gateway, presence and broker, and trace one message through "Online" and one through "Offline".

Likely follow-ups, with short answers:

Further reading #

Now design it

More system design problems