IM Systems in Depth

Chapter 0

The whole picture

What parts make up a messaging system, and how does one message get from Ana’s phone to every one of Ben’s devices?

This chapter is a first draft. It will be revised once the first eight chapters are written.

Sending a message looks simple. Ana sends it to a server, and the server passes it to Ben. You can write a working version in an afternoon, and on the office Wi-Fi it really does work.

The trouble starts outside the office:

  • Ana presses Send in a lift and loses signal for a second. Did the message go out? Her app does not know.
  • So the app sends it again, and Ben gets the same message twice.
  • Ana and Ben write at almost the same moment, and each of them sees the two messages in a different order.
  • Ben has not opened the web app on his computer all day. When he does, it must find out which messages it missed, without asking each of his hundreds of conversations one by one.
  • In a group of 500, one message must reach about 750 devices; in a live room of 1,000,000 viewers, it is pushed a million times.
  • A million long-lived connectionslong-lived connection长连接A TCP or WebSocket connection a device keeps open to the server. With it, the server can push a message at any time instead of waiting to be asked.See the glossary send nothing most of the time, yet each one holds memory. Routers at the mobile carrier quietly forget connections that stay silent too long, and one restart for a deploy sends 100,000 phones back in the same second.

Meanwhile, people expect a lot from chat: nothing lost, nothing shown twice, everything in order, and all at once. A web page that is half a second slow goes unnoticed; a lost message gets a “did you get my message?” right away.

Messaging is hard not because one of these problems is very hard, but because they arrive together and pull on each other. Retries stop messages from getting lost, and retries create duplicates. Long-lived connections make chat instant, and they bring a million connections and reconnect storms. Fewer heartbeats save the battery, and then dead connections go unnoticed. Each fix brings the next problem.

This series, IM Systems in Depth, shows step by step how a messaging system grows from the afternoon version into one that handles all of this. Before we start, let’s look at the whole system once.

1. What a messaging system must do

Most backends work by question and answer: the client asks, the server answers. You refresh, and a page comes back. Messaging works the other way round. Ben does nothing, and the server must bring Ana’s message to his phone. That ability comes from the long-lived connection, so strictly it belongs to the connection layer, not to messaging itself: the same connection layer can carry live-stream comments, push notifications or device control (chapter 19). Messaging adds a whole set of guarantees on top of it.

In one sentence: every message Ana sends must appear on every one of Ben’s devices, exactly once, in the same order everyone else sees, even if one of those devices was offline all day.

Every part of that sentence has a cost:

  • “Every message”: the network loses packets, so you need acknowledgments (ACKs)ACKACK(确认)The receiver’s “got it” to the sender. In this book, the server’s ACK means the message is in the message log, not that anyone has received it.See the glossary and retries.
  • “Every one of Ben’s devices”: Ben has a phone and the web app on his computer, and they must agree.
  • “Exactly once”: the network can only deliver at least once, and retries create duplicates, so repeats are droppeddedup去重When a message ID arrives again, return the first result instead of storing another copy. The network can only deliver at least once; dedup makes each message show once.See the glossary by message ID and each message is shown once.
  • “In the same order”: no two clocks agree, so timestamps cannot decide the order.
  • “Offline all day”: when a device comes back, it must find out what it missed, and fetch it efficiently.

Replace “Ben” with “a group of 500”, or “a live room of 1,000,000”, and the same promise becomes a very different engineering problem.

2. The running example

The app: an ordinary social chat app, with 1:1 chats and groups of up to 500, on phones, desktops and the web. It has no name, so it never looks like a real product.

The people: Ana and Ben chat with each other, and both are in “the hiking group” (500 members). Ben has a phone and the web app on his computer.

The numbers: the whole series shares one set of assumptions, and every estimate comes from it. For example:

  • 10% of users are online at the peak;
  • each user sends 40 messages a day;
  • each person has 1.5 devices on average;
  • a message is about 200 bytes;
  • 10% of packets are lost. That is chosen high so failures are easy to see; good networks lose far less.

So one message in the hiking group is written into 500 members’ inboxesinbox信箱A per-person index of “new messages for you”, with its own per-user sequence. Write fan-out writes it, and devices pull from it by sync cursor.See the glossary and pushed to about 500 × 1.5 = 750 devices. By v3 the app has 10,000,000 daily users, with 1,000,000 online at the peak.

These are round numbers chosen for the example, not data from any real company.

3. The whole system, and one message’s journey

The diagram below is the system this series builds by the end of the trunk (version v3), drawn by layer:

  • At the top are two devices: Ana, who sends, and Ben, who receives.
  • Across the internet are the IM server’s four layers: connection (one long-lived connection per device), message (receiving), dispatch (delivering to every device) and storage, with shared infrastructure at the bottom.
  • The column on the left is the business layer: accounts, friends and groups, moderation. It sits beside the others because messages do not flow through it; the layers ask it when they need to. It is also the part that changes most often, and keeping it beside the path keeps the other layers stable.
  • The dashed boxes, the business systems (your app’s backend) and OS pushOS push系统推送A notification sent through Apple’s APNs, Google’s FCM or a phone maker’s channel to a phone with no live connection. A third party, best effort, with no guarantee.See the glossary (Apple, Google, phone makers), are not part of the IM.

A dotted line runs between the message and dispatch layers. Above it is receiving: the sender waits for the ACK. Below it is delivering, and from there on everything is asynchronous. The message logmessage log消息日志A durable, append-only queue: the hand-off between receiving and delivering. A message counts as received once it is in the log, and the server ACKs then. It usually runs on a message queue, such as a topic partitioned by conversation.See the glossary is the one hand-off between the two. Many teams say “access, logic, storage”: the connection layer is the access layer, and the message and dispatch layers together are the logic layer.

Look at the whole picture first; point at or tap any layer or card to see what it does and why it is hard. Then press “Follow this message”: Ana asks the hiking group “Still on Saturday?”, and the diagram lights its eight steps one at a time. The four steps of the overview on the home page are steps 1–2, 3–4, 5–7 and 8 here.

How one message crosses the system: first it is received safely (steps 1–4), then delivered to every device asynchronously (5–7); a device that was offline pulls it when it returns (8).

Sender: Ana’s phoneStill on Saturday?IM SDK: outbox, local database, syncReceiver: Bena sync cursor per devicePhoneonlineWebofflineInternet: TCP or WebSocket, TLSIM serverBusinesscalled, not passedAccounts, authRelations, rulesblocks, reviewchecked before ACKGroupsmember listsOpen APImsgs, callbacksConnectionholds connections,pushes to devicesgrows with devicesLocatorpicks a gatewayGatewaysone connection per online device: auth, heartbeats, trafficMessagereceives: no loss,no repeats, in ordergrows with messagesMessage serviceone owner per chat: dedup, seqDispatchdelivers toevery devicegrows with fan-outRecipientschats, groupsFan-outwrite or readRoutingfinds gatewaysOnline pushOffline pushSyncStoragemessages, stategrows with retentionMessagesper chatRead cursorsread up toOnline routesdevice → gatewayMetadata, cachechats, filesInboxesone per userInfrastructureMessage queueDiscovery, configRate limitsMetrics, tracing↑ sync↓ asyncMessage logdurable queueoutside the IMBusiness systemsapp backend, supportoutside the IMOS pushAPNs, FCM, vendorsaccounts, relations, messageswakes a phone in the background12345678
  1. Send: the message enters a gateway over Ana’s connection
  2. The gateway reads the conversation ID and forwards it to the message service that owns it
  3. Check whether the message ID was seen recently (a resend gets back its original seq and ACK; a much later resend is stopped by a unique key in storage), then check the sender may post, assign the chat’s next seq, append to the message log (a durable queue)
  4. Once the append succeeds, the message service ACKs: the sender waits only this long
  5. Dispatch reads the log, writes the message store, then delivers (asynchronous from here on)
  6. Work out the recipients, fan out (write fan-out copies into each inbox), look up online routes
  7. Push through the gateway Ben’s phone is on. A phone with no live connection gets an OS push instead, only a hint; the web app uses no OS push here and catches up when opened
  8. A device that was offline comes back with its sync cursor (its position in its own inbox), and the sync service returns what it missed from the inbox (write fan-out) or the group’s timeline (read fan-out); an online device that sees a gap in the chat’s seq pulls too, and overlaps between push and pull are dropped by seq

blue: the message’s pathblue dashed: the pull grey: calls, writes, notificationsdashed box: outside the IMdotted: the sync/async boundary

Many teams say “access, logic, storage”: the connection layer is the access layer, and the message and dispatch layers together are the logic layer. The message log usually runs on the message queue in the infrastructure row, for example as a topic partitioned by conversation.

Point at or tap a layer or card to see its job, why it is hard, and its chapters.

A few words first, since every chapter uses them (all terms are in the glossary):

  • The message log is the hand-off between receiving and delivering: a message counts as received once it is in the log. The message store keeps the history (one timeline per conversation). An inbox is a per-person index of “new messages for you”.
  • Write fan-out: when a message is sent, write one entry into each member’s inbox; reading means reading your own inbox. Read fan-out: keep one copy in the group’s timeline, and each reader fetches it from there. The 500-member hiking group uses write fan-out; bigger groups cannot afford it and switch to read fan-out (chapter 10).
  • Two kinds of number: the group’s seq (4,811, 4,812, …) sets the order in the group and is used for ordering and dedup; each inbox has its own per-user sequence, and each device’s sync cursor points into it.

This is only a walk, to show where each thing happens. From chapter 1 on, each chapter takes one of these mechanisms away, and you watch in the simulator what goes wrong.

4. Group chat and live rooms: where dispatch parts ways

In both a group and a live room, one message goes to many people, but the two take different paths, and the difference is all in the delivering part. Note that a live room of 1,000,000 viewers is already another product: it equals this app’s whole peak online. The figure assumes a busy room sends 100 messages a second, and its viewers are spread over 20 gateways, 50,000 connections each.

One message, delivered two ways. A group finds every device of every member and makes sure each can catch up; a live room finds only the gateways, and each gateway pushes to the connections it holds for that room. Losing some is allowed.
Group: 500 membersDispatchOnline routesmember → devices → gatewayone lookup per memberInboxesone each, 500 writesfan outbatched per gateway, one copy per devicegateway 3gateway 7gateway 12…500 lookups, one per member500 inbox writesabout 750 devices, one copy eachoffline devices catch up later from the inboxLive room: 1,000,000 viewersRoom routergateways register on joinNo per-user inbox, no catch-upa room may keep a seq, for dedup and order onlysent to these 20 gateways onlygateway 150,000 connsgateway 250,000 connsgateway 2050,000 conns…20 sends to gatewaysno per-person lookups or inboxes: each gateway knows its room connectionseach gateway pushes to its own 50,000 connections

100 messages a second × 1,000,000 viewers = 100,000,000 a second. So each gateway packs a second of messages into one batch (about 2 KB) and pushes each connection once a second: 50,000 pushes a second per gateway. Beyond that, sample and prioritise (gifts first).

In a group, every member matters: every message must arrive and every device must be able to catch up, so dispatch finds each member’s devices and writes each member’s inbox. In a live room, a million people watch at once, the screen shows a dozen messages, and a late message is worth nothing. So dispatch finds only the gatewaysgateway网关A connection-layer server. Each holds many long-lived connections and handles auth, heartbeats and traffic without reading messages. One of the few stateful parts: a restart drops all its connections.See the glossary, not the people, and losing some is fine. Chapter 40 builds this design step by step from the numbers.

Live rooms are one example. The last chapters (38–44) change the product in other ways too, and the same decisions flip.

5. How the system grows

The diagram above is the end point. The series starts at the beginning: one chat app grows from 1,000 users to 1,000,000 people online. All four layers are there from the start, but at first they are crowded into one program. Each step of growth forces one of them out and makes it deeper. Switch versions below:

1,000 users“We want chat in our app.” Make it work: one program and one table; the client polls, then switches to a long-lived connection.

Business (beside)a few permission checksClientsends requests, shows repliesConnectionpolling → a long-lived connectionMessagestore what arrivesDispatchpush to whoever is onlineStorageone messages tableone programone database
new in this version The smallest chat … Push instead of poll (chapters 1–2)

6. How to read a chapter

From chapter 1 on, each chapter makes one design decision, always in the five steps listed on the home page: the problem, watch it fail, estimate, fix it, and the cost and other answers. At the end of each chapter the map gains one piece, with a decision card. The simulator is a sandbox. It talks to no real server, so break it on purpose: raise the packet loss, switch a mechanism off, and see what happens.

7. Where to start

  • To follow the story: chapter 1, “The smallest chat” (planned, the next one).
  • For one layer only: tap it in the picture above and see which chapters cover it, or open the map.
  • For every chapter at once: the chapter map on the home page.
  • For a word you don’t know: the glossary. Words with a dotted underline show their meaning when you point at or tap them.