A chat system looks like a sending problem: deliver a message from one person to another quickly. Messenger's engineering history shows that the hard part is a syncing problem. Each person has several devices, mobile networks drop and change constantly, and every device must end up with the same conversations in the same order without downloading everything again each time it wakes up.

This article follows Messenger's architecture as Meta's engineers have described it: the 2014 move to snapshot-plus-delta sync over MQTT backed by a service called Iris, the 2018 storage migration from HBase to MyRocks, the 2020 LightSpeed client rebuilt around SQLite, and default end-to-end encryption from December 2023. For a generic design built from scratch, see designing a real-time chat system; here the subject is the decisions one very large messenger actually made, and why.

Advertisement

From pull to push: the 2014 redesign

Early Messenger apps fetched state over HTTPS: ask the server what changed, receive a response, render. On mobile that is wasteful. Each poll costs a round trip and radio time even when nothing changed, and the response often repeats data the device already has. Facebook's 2014 engineering post, Building Mobile-First Infrastructure for Messenger, describes the replacement: a client fetches an initial snapshot of its messages once over HTTPS, then subscribes to deltas that the server pushes over MQTT, a lightweight publish-subscribe protocol built for low power and low bandwidth.

The team also switched the wire format from JSON to Thrift, which they report cut payloads by about half. Combined, the new protocol reduced non-media data use by 40 percent and, by reducing congestion, cut message send errors by 20 percent. The lesson generalises: on mobile, the protocol shape, snapshot once then push only what changed, matters more than server speed.

Iris: an ordered queue with many pointers

Messenger sync: one ordered update queue, many pointers, three storage tiersSender appsend over MQTTIris: ordered update queueseq 101 | 102 | 103 | 104 | 105enqueuePhonepointer at 105Laptop, offlinepointer at 102push deltawaitsHot: memorynewest updatesWarm: MySQL on flashabout one weekCold: long-term storefull history, snapshotsstorage pointerReconnect within windowsend updates after the device's pointerPointer fell off the windowfresh snapshot over HTTPS, then deltasEach consumer advances its own pointer, so an offline device never blocks delivery to the others.
Iris keeps updates in order and lets each consumer, whether a device or the storage tier, advance its own pointer. Recent updates live in memory, about a week in MySQL on flash, older history in long-term storage.

The service behind the protocol is Iris, which Meta describes as a totally ordered queue of messaging updates: new messages, but also state changes such as a message being read. Separate pointers track how far each consumer has got, one for each app and one for the storage tier. When a device goes offline, its pointer stays put while new updates are enqueued and the other pointers advance. When it returns, the server sends what lies after its pointer.

Iris tiers its data by age. The newest updates sit in memory and go straight to online apps and to storage. About a week of updates is kept in MySQL on flash, which serves apps that were briefly offline and covers a storage outage. Older history, and the full inbox snapshots, live in the long-term store. Meta reports that enqueueing became an order of magnitude faster than writing to traditional disk, and that semi-synchronous MySQL replication let it handle a database hardware failure within about 30 seconds.

The design separates two jobs that naive chat systems merge. Delivery reads from the head of a short, fast, ordered log. Durable history is just another consumer of that log, so a slow or failed history store delays history, not delivery.

Advertisement

The sync protocol in code

The core is small: a server log with sequence numbers and a client that remembers the last sequence it applied. The rules that make it robust are in the edge cases: duplicates after reconnect are ignored, a hole triggers catch-up, and a client whose pointer has fallen out of the retained window takes a fresh snapshot.

class Mailbox:
    """Server side: an ordered log of updates with one pointer per consumer."""
    def __init__(self, window):
        self.log = []                 # [(seq, update)] newest last, trimmed to window
        self.next_seq = 1
        self.window = window          # how many updates we keep for catch-up

    def append(self, update):
        seq = self.next_seq
        self.next_seq += 1
        self.log.append((seq, update))
        if len(self.log) > self.window:
            self.log.pop(0)
        return seq

    def since(self, last_seen):
        oldest = self.log[0][0] if self.log else self.next_seq
        if last_seen + 1 < oldest:
            return None               # gap: caller must resnapshot
        return [(s, u) for s, u in self.log if s > last_seen]

class Client:
    def __init__(self, server, mailbox_id):
        self.server, self.mid = server, mailbox_id
        self.state, self.last_seen = {}, 0

    def connect(self):
        delta = self.server.mailbox(self.mid).since(self.last_seen)
        if delta is None:
            self.state, self.last_seen = self.server.snapshot(self.mid)
            delta = self.server.mailbox(self.mid).since(self.last_seen) or []
        for seq, update in delta:
            self.on_push(seq, update)

    def on_push(self, seq, update):
        if seq <= self.last_seen:
            return                    # duplicate after reconnect: ignore
        if seq != self.last_seen + 1:
            return self.connect()     # hole in the sequence: catch up first
        apply_update(self.state, update)   # must be deterministic
        self.last_seen = seq

Two properties carry the design. Updates must be applied deterministically, so two devices that apply the same sequence reach the same state; that is why read receipts and edits are updates in the log rather than side effects. And the retained window bounds both memory and catch-up cost: a device offline beyond it pays one snapshot instead of an unbounded replay.

Worked example: a laptop that was closed for three days

A user's phone is online; their laptop has been closed since Friday with its pointer at sequence 102. Over the weekend, friends send 40 messages, the user reads some on the phone, and one message is edited. Each of those is an update with a sequence number, 103 onwards, and the phone's pointer has advanced with them.

On Monday the laptop reconnects with last seen 102. Three days is within a one-week warm window, so the server returns updates 103 to 145, and the laptop applies them in order: messages appear, read state flips for those already read on the phone, and the edit replaces the original text. If a push for 146 arrives while 140 is still in flight, the client sees the hole and catches up rather than rendering out of order. Had the laptop been closed for a month, its pointer would be outside the window, and it would download a fresh snapshot and resume deltas from there. Nothing in either path depended on the phone.

Storage: from HBase to MyRocks

Messenger's message store ran on HBase, on HDFS, from 2010. Meta's 2018 post, Migrating Messenger storage to optimize performance, describes moving every account to MyRocks, its MySQL storage engine built on RocksDB, running on flash servers. The reported results were 90 percent less storage consumed and read latency 50 times lower than before.

The migration is a template for moving a live, write-heavy store without downtime. Each account moved through three states: not migrated, double-writing and done. While double-writing, the migrator validated both the data and the API by sending reads to both systems and comparing the answers. Normal migration moved 99.9 percent of accounts within two weeks. The remainder, very high-traffic accounts such as large businesses running bots, used a buffered migration: data up to a cutoff was copied to a buffer tier while new writes for the account queued in Iris, then the queue drained into the new store.

NOT_MIGRATED, DOUBLE_WRITING, DONE = "not_migrated", "double_writing", "done"

def migrate_account(acct, old, new, iris):
    if acct.state == NOT_MIGRATED:
        acct.state = DOUBLE_WRITING            # writes now go to both stores
        new.bulk_load(acct.id, old.export(acct.id))
    if acct.state == DOUBLE_WRITING:
        for req in sample_reads(acct.id):
            if old.read(req) != new.read(req):  # data and API validation
                return rollback(acct)           # stay on old store, investigate
        acct.state = DONE                       # reads switch to new store

def migrate_heavy_account(acct, old, new, iris, cutoff):
    iris.hold_consumer(acct.id, "storage")     # new writes queue in Iris
    buffer = old.export(acct.id, until=cutoff)
    new.bulk_load(acct.id, buffer)
    iris.resume_consumer(acct.id, "storage", into=new)   # drain the queue

Iris is what made the buffered path possible. Because storage is just a consumer with a pointer, pausing that consumer is free: delivery continues, writes accumulate in order, and the new store catches up from the pointer. The message queues article explains the same consumer-offset pattern in Kafka-style systems.

The client as a database: LightSpeed

In 2020 Meta described Project LightSpeed, a rebuild of the iOS app. Core Messenger code shrank from more than 1.7 million lines to 360,000, the app started twice as fast and was a quarter of the size. The architectural move was to make SQLite the universal system on the device: views read from the local database, UI configuration is stored in it, and business logic runs as stored procedures written in CG-SQL, a language the team built that compiles to C for SQLite. A server broker acts as a single gateway between the app and server features.

This completes the sync design. If the server sends a deterministic stream of updates, the natural client is a local database that applies them, and the UI is a view over that database. Offline reading, instant rendering at launch and consistent state across screens then come for free. The cost is that schema changes must be migrated on billions of devices, so the local schema is versioned and evolved as carefully as a server schema.

What default end-to-end encryption changed

In December 2023 Meta began rolling out end-to-end encryption by default for personal chats and calls on Messenger and Facebook, publishing a whitepaper on the messaging protocol and another on Labyrinth, its protocol for end-to-end encrypted storage of message history across the devices linked to an account.

The architectural consequences follow from one fact: the server can no longer read message content. Every feature that ran on server-side plaintext must move to the client or become something the client controls. Search over history runs on the device against the local database. Link previews, if offered, must be generated without the server reading the conversation. In the usual multi-device design each linked device has its own keys, so a message is encrypted for each device of each participant, and multi-device fan-out grows with device count. New devices cannot simply download history from the server in readable form; that is the problem Labyrinth addresses, keeping history recoverable on a new device while it remains encrypted to the server.

The ordered-log design survives intact, because ordering, delivery and pointers never needed to read content. What changes is what an update carries: an encrypted payload the server routes but cannot inspect. Compare the device-key approach in iMessage's architecture.

Failure modes

FailureSymptomMitigation in this design
Device offline beyond retained windowMissing messages if deltas were assumedDetect gap, fall back to snapshot
Duplicate push after reconnectMessage shown twiceIgnore sequence numbers at or below last seen
Out-of-order deliveryEdits applied before the message they editApply strictly in sequence; catch up on holes
Storage tier outageHistory writes failStorage pointer pauses; Iris retains updates
Migration divergenceOld and new stores disagreeDouble-write with read comparison; roll back per account
Hot accountOne mailbox overwhelms a partitionBuffered migration; per-account limits
Client schema driftCrashes after app updateVersioned local schema with migrations tested on real data

Trade-offs

A per-mailbox ordered log makes reasoning simple and multi-device sync correct, but it concentrates writes for busy mailboxes and needs a gap-and-snapshot path that is exercised rarely, so it must be tested deliberately. Push over a persistent connection saves battery and data but means keeping very large numbers of long-lived connections and handling reconnect storms when a network region recovers; reconnect with jitter and serve catch-up from the warm tier. End-to-end encryption removes server-side features and moves their cost to clients. For notification delivery when the app is not connected, the notification system design covers the platform push path that complements MQTT.

What to do next

  1. Model your chat or collaboration sync as an ordered update log per mailbox with sequence numbers, and make every state change, including reads and edits, an update.
  2. Implement the client rules above: ignore duplicates, catch up on holes, resnapshot when the pointer falls outside the retained window, and test each path.
  3. Make durable history a consumer of the log with its own pointer, so storage outages and migrations do not stop delivery.
  4. Plan any storage move as a per-account state machine with double writes and read comparison, plus a buffered path for the hottest accounts.
  5. Treat the client store as a database replica with a versioned schema, and decide early which features will survive end-to-end encryption.
Key takeaway: Messenger is a sync system built around one idea: an ordered log of updates per mailbox, with every device and the storage tier advancing its own pointer. Snapshot once, then push deltas over MQTT; keep recent updates fast and older history durable; migrate storage behind the log with double writes and read comparison; make the client a local database that applies the stream. Default end-to-end encryption kept that structure and moved content-dependent features onto devices.