Skip to content

The finished design, decision by decision

How to design a Collaborative Editor (Google Docs)

Not the one correct diagram, but a design you can defend under these constraints: the finished architecture, then every stage's question with the reasoning that answers it, the tradeoffs it accepts, and where another engineer could land differently.

A real-time collaborative editor

The short answer

9 parts, each with one job. The map below shows how requests and data move between them; the stages after it explain why each part is there.

Editor client
Applies local edits immediately, keeps unacknowledged ops in a pending buffer, and integrates remote ops by sequence number.
Document router
Routes every connection for a document to that document's current owner.
Document owner
The single sequencer for a document: validates and integrates ops, appends them durably, then acks and broadcasts. Holds presence in memory.
Postgres op log
Durable, ordered history of every operation, and the ownership epoch for each document.
Pub/sub
Carries a hot document's sequenced ops from its owner to the edge servers holding its viewers.
Fan-out edge
Holds read-mostly connections, batches frames per viewer, and drops viewers that fall too far behind into catch-up mode.
Viewers
Read-mostly participants in a large document.
Snapshot storage
Holds document snapshots, each tagged with the exact sequence number it includes.
Snapshotter
Folds ops into periodic snapshots and compacts old history into coarser versions.
123456789CLIENTEditor clientEDGEDocument routerSERVICEDocument ownerDATABASEPostgres op logLOG / STREAMPub/subSERVICEFan-out edgeCLIENTViewersOBJECT STORESnapshot storageWORKERSnapshotter

Select a component to see what it is responsible for and which state it owns.

  1. 1Editor client → Document router: WebSocket: ops, acks, remote ops, presence
  2. 2Document router → Document owner: Route by document id to the current owner
  3. 3Document owner → Postgres op log: Append ops at next seq (batched, epoch-fenced)
  4. 4Document owner → Snapshot storage: Load latest snapshot on open
  5. 5Document owner → Pub/sub: Publish sequenced ops
  6. 6Pub/sub → Fan-out edge: Per-document subscription
  7. 7Fan-out edge → Viewers: Batched frames, bounded buffers
  8. 8Snapshotter → Postgres op log: Read ops since the last snapshot
  9. 9Snapshotter → Snapshot storage: Write snapshot tagged with its seq
  • Server push
  • Request / response
  • Bulk data
  • Asynchronous

Why does this design work?

Each document has one owner that does the only inherently serial work: assigning sequence numbers. Everything else is designed to happen anywhere and more than once.

  • Clients apply edits locally first and keep custody of each op until the owner acknowledges it, and the owner acknowledges only after the op is durable. Reconnection is a replay from the durable log by position, plus a resend of pending ops whose ids make duplicates harmless.
  • Concurrent edits converge through a merge algorithm rather than by choosing a winner, so concurrency and offline work are normal operating conditions rather than errors.
  • Ownership is a fenced lease, so deploys and crashes move documents between servers without ever producing two histories.
  • Presence is soft state, cheap and self-healing, kept out of the durable path entirely.
  • Fan-out is separated from sequencing, so a document with hundreds of viewers scales delivery without breaking the single-writer order.

Invariants, and where they are enforced

  • Each document has exactly one history: every sequence number is assigned once.

    One owner per document, holding a lease with an epoch; appends are conditioned on the current epoch, and a unique (doc_id, seq) constraint rejects any second writer. Enforced by Document owner, Postgres op log.

  • An acknowledged op is never lost.

    The owner acks only after the op's batch commits; clients keep every op until it is acked and resend on reconnect, deduplicated by op id. Enforced by Document owner, Postgres op log, Editor client.

  • Clients that have applied the same set of ops show the same document.

    Ops are integrated through a convergent algorithm (server-ordered OT or a CRDT) and applied idempotently by op id. Enforced by Document owner, Editor client.

  • Nobody appears present for longer than a few seconds after they leave.

    Presence is soft state refreshed by heartbeats and expired on a TTL; it is never persisted. Enforced by Document owner.

What does it rely on?

  • Clients persist pending ops locally and reuse op ids when resending.
  • The router sends all of a document's connections to its current owner, and ownership changes are fenced.
  • The merge algorithm is correct; periodic checksum comparisons catch divergence if it is not.
  • Postgres can absorb batched appends from all active documents; past that, the op log is partitioned by document.
  • Documents are small enough to hold in an owner's memory.

What tradeoffs does it make?

ChoiceGainsCosts
Single owner per documentFree ordering, local fan-out and in-memory state.Document-aware routing, ownership handoffs, and hot-document limits.
Ack after durable commit (group commit)No acknowledged edit can be lost.A few milliseconds of added latency per batch.
OT or CRDT mergingConcurrent and offline edits preserved.Algorithmic complexity and per-character metadata (CRDT).
Soft-state presenceNo durable writes for ephemeral data; ghosts expire on their own.Brief inaccuracy after restarts.

What are the reasonable alternatives?

Hosted sync service or CRDT backend (e.g. a managed realtime database)
Better when the team wants collaboration as a feature rather than a core competency, and the service's data model, permissions and pricing fit.
Peer-to-peer CRDT sync with a relay server
Better when local-first operation and privacy matter most, and the server should store opaque data.
Paragraph-level locking
Better when collaboration is occasional, documents are structured, and predictability beats fluidity.

When does it stop working?

  • A single document's sequencing (not delivery) outgrows one process, for example thousands of simultaneous typists.
  • Documents become too large to hold in memory, which calls for per-section ownership and a different model.
  • Clients lose local storage while holding unsynchronized offline edits.
  • Regulatory requirements demand that deletions purge every snapshot and log entry, which conflicts with immutable history.

Every stage, decided and explained

Spoilers, for the whole investigation: each stage's question and its answer, the reasoning behind it, and the tradeoffs it accepts. If you have not worked through the stages yet, you may want to do that first.

Work through the stages

Stage 1 of 12 · Model

Size the problem

Real-time systems fail in ways that depend on numbers: message rates, fan-out factors, connection counts. Work some out before you choose anything. Useful figures: 20,000 connections at peak, about a tenth of users typing at any moment, 5-10 operations per second per typist, and one document with 30 editors and 200 viewers.

What you need to know first

"Last write wins" on the whole document means: whoever saves second replaces whatever the first person wrote. With two people typing in the same second, someone's sentence disappears.

Collaborative editing has to merge operations (insert "x" at this point, delete these characters), not replace whole documents.

20,000 people connected, a tenth of them typing at about 7 operations a second. About how many operations a second arrive?

About 14,000 ops per second.

20,000 × 0.1 × 7 = 14,000 ops a second. As 14,000 separate committed transactions that's heavy; batched per document it's modest.

The all-hands doc: 30 editors at about 7 ops a second each, every op delivered to about 230 participants. About how many outbound messages a second?

About 46,000 messages per second.

30 × 7 ≈ 210 ops a second; × 230 recipients ≈ 48,000 messages a second (about 46,000, excluding each sender) for one document. Delivery, not ingestion, is the hot path.

Do 20,000 idle WebSocket connections need a large server fleet?

No: an idle connection costs tens of kilobytes on an event-driven server; what costs is message rate and buffering.

20,000 × ~30 KB is well under a gigabyte. A few servers hold the connections; traffic decides the rest.

What the stage asks

Which statements hold?

  1. Fails

    Sending the whole document on each change works if the change is debounced to once a second.

    Bandwidth is the smaller problem. Whole-document last-write-wins means any two people typing within the same second overwrite each other, which is exactly the bug in testing. The unit of change has to be an operation, merged rather than replaced; see Conflict resolution and convergence.

  2. Holds

    Polling once a second cannot meet the ~200 ms latency goal.

    Average delay would be half a second plus request time, and polling faster multiplies requests for every idle client.

  3. Holds

    The all-hands document alone can require tens of thousands of outbound messages per second.

    30 editors × ~7 ops/s ≈ 200 ops/s, each delivered to ~230 participants, is about 46,000 messages a second for one document, unless you batch. Fan-out, not ingestion, is the hot path; see Backpressure and capacity.

  4. Depends

    20,000 concurrent WebSocket connections require a large server fleet.

    Idle connections are cheap on an event-driven server: tens of kilobytes each, so 20,000 fit in well under a gigabyte. What costs is message rate and per-connection buffering. A handful of servers holds the connections; the fleet size is set by ops and fan-out. See Persistent connections.

  5. Depends

    Writing every operation durably to Postgres is infeasible at this scale.

    About 2,000 typists × ~7 ops/s ≈ 14,000 ops a second. As 14,000 separate committed transactions, that is a heavy load for one primary. Batched per document, committing every 10-20 ms with the ack waiting for the batch, it becomes a few hundred transactions a second carrying small rows, which is very manageable.

The reasoning

  1. Concurrent edits are normal: merge operations instead of replacing documents.
  2. About 14,000 ops a second inbound means writes must be batched.
  3. One hot document can need tens of thousands of outbound messages a second; connections themselves are cheap.

Three numbers frame the design:

  • ~14,000 ops/s inbound, so writes must be batched (without weakening acknowledgements).
  • ~46,000 msgs/s for one hot document, so fan-out needs batching and its own capacity.
  • 20,000 connections, which is less than it sounds. Connections are not the bottleneck; traffic is.

And one non-number: concurrent edits are normal, not an edge case. The design has to merge rather than replace.

Stage 2 of 12 · Decide

Choose the transport

Each active editor streams 5-10 small operations a second to the server and receives everyone else's from it, ideally within 200 ms. Order matters: an editor's operations depend on the ones before them.

What you need to know first

Transports differ in direction, overhead and ordering:

TransportDirectionOrdering of one client's messages
Pollingclient askseach request separate
Server-Sent Eventsserver → clientordered downstream only
HTTP POST per opclient → serverseparate requests can arrive out of order
WebSocketboth ways, one connectionordered both ways

See Server push: polling, long polling, SSE, WebSockets and Persistent connections.

Each keystroke is a separate POST. Op 2 is sent before op 3, but op 3's request arrives first. What must the server do?

Reorder by a client sequence number, holding op 3 until op 2 arrives

Op 3 was computed assuming op 2 had happened. Separate requests lose ordering, so the server has to rebuild it. One WebSocket preserves it for free.

A persistent socket comes with commitments: heartbeats to detect dead clients, a protocol to resume after reconnecting, and a plan for deploys that close every socket.

What the stage asks

How should editor clients talk to the server?

  1. Sound

    One WebSocket per open document: ops up, acks and remote ops down

    High-rate, small, ordered messages in both directions are what WebSockets are for. One connection gives in-order delivery of a client's ops, minimal per-message overhead, and a natural place for heartbeats. You inherit reconnection and resumption, which the design needs anyway because clients go offline.

  2. Defensible

    Server-Sent Events for remote ops; an HTTP POST per local operation

    Downstream is fine. Upstream, ten POSTs a second per typist adds HTTP overhead, and separate requests can arrive out of order (or retry out of order), so the server must reorder each client's ops by a client sequence number. It is workable, and sometimes chosen for proxy-hostile environments, but you rebuild what the socket gives you.

  3. Defensible

    Long polling for remote ops, POSTs for local ones

    It delivers quickly and works everywhere, but every delivered batch costs a new request, messages between polls must be buffered per client, and you still have the POST ordering problem. It is a historical fallback rather than a design choice today.

  4. Flawed

    Peer-to-peer WebRTC data channels between editors

    Thirty editors in a full mesh means 435 connections, NAT traversal failures, and no server that holds the durable history or enforces permissions. The requirement that no acknowledged edit is lost needs a durable party in the loop.

What a strong answer covers

  • Traffic is high-rate in both directions, so a persistent full-duplex channel fits.
  • A single connection preserves the order of a client's own ops; separate requests do not.
  • Names what the choice costs: reconnection, heartbeats, deploy handling, routing.Supporting

The reasoning

  1. High-rate, ordered traffic in both directions fits one persistent WebSocket per document.
  2. Separate requests lose a client's operation order.
  3. Sockets bring heartbeats, reconnect protocols and deploy handling as obligations.

Compare this with the status page in the video pipeline, where polling was the right call. The difference is the traffic: a few updates over twenty minutes versus dozens of messages per second in both directions. Same question, different constraints, different answer; see Server push: polling, long polling, SSE, WebSockets.

The costs of the socket are now commitments: heartbeats to detect dead clients, a resume protocol for reconnects, and a plan for deploys. Several later stages are those commitments coming due.

Stage 3 of 12 · Decide

Where does a document's live state live?

Five editors of one document may be connected to five different servers. Each op must be ordered relative to every other op for that document, persisted, and delivered to the other four within 200 ms.

What you need to know first

Every document needs one authority that decides the order of its operations. Clocks can't do it: client clocks disagree by seconds, so ordering by timestamp gives different orders on different machines. See Ordering.

Give each document an owner server: it holds the document in memory, assigns each op the next sequence number, persists it and broadcasts it. Ordering becomes a counter in memory.

That turns ordering into routing: every connection for document 42 must reach its current owner. A router can map document IDs to servers with consistent hashing, so that when servers come and go, most documents stay where they are. See Consistent hashing.

Think of the keys as documents and the nodes as collaboration servers. Every key that moves is a document whose owner changes, and whose editors must reconnect.

2,000 keys spread over cache nodes. Add a node and count the keys that now belong somewhere else: each one is a cache miss, or data to copy.
Placement
4 nodes
4 nodeskeys by hash mod N

Keys per node. The line is a perfectly even share.

  • Node 1500
  • Node 2508
  • Node 3512
  • Node 4480

The busiest node holds 1.0× an even share.

Alternative: any server accepts any client, and each op takes the next sequence number from a row in Postgres. What's the cost?

Every keystroke waits for a round trip to a contended database row before it can be ordered.

It works, but the document's counter row becomes a lock every op queues on, and every server must hold a copy of the document.

What the stage asks

How should servers coordinate on a document?

  1. Flawed

    Any server accepts any client; servers write ops to Postgres and poll it for others' ops

    Polling adds its interval to every edit's latency and multiplies database reads by the number of servers times the number of open documents. Ordering also depends on whichever insert commits first, and every server has to cope with it.

  2. Defensible

    Any server accepts any client; each op gets the next sequence number from Postgres, then is broadcast to other servers via pub/sub

    It works. The database row lock on the document's counter is the sequencer, and pub/sub delivers. But every keystroke now waits for a round trip to a contended row before it can be ordered, every server holds a copy of every open document, and transforming ops needs the latest state on whichever server received them.

  3. Sound

    Route every connection for a document to one owner server, which sequences, persists and broadcasts its ops

    One process holds the document in memory and assigns sequence numbers locally, so ordering costs nothing. Broadcast is local to that process. The cost moves to routing (a document-aware router or directory), to hot documents that outgrow one server, and to keeping ownership unique during deploys.

  4. Flawed

    Clients order ops by timestamp; servers just relay them

    Client clocks disagree by seconds, so ops would be ordered wrongly and differently on different clients. Order needs a single authority, not a consensus of clocks; see Ordering.

What a strong answer covers

  • Each document needs one authority assigning order; clocks cannot provide it.
  • Partitioning by document gives each its own sequencer while documents scale out independently.
  • Names what this creates: routing, hot documents, and keeping ownership unique during moves.Supporting

The reasoning

  1. Each document needs one ordering authority; client clocks can't provide it.
  2. A single owner per document makes sequencing a local counter.
  3. The hard parts move to routing connections to the owner and keeping ownership unique.

Partitioning by document turns a distributed ordering problem into a routing problem. Within one owner, ordering is a counter in memory. The hard parts are now:

  • Routing: every connection for document 42 must reach its current owner. Consistent hashing at the router keeps most documents in place as servers come and go; a small directory (doc → owner, epoch) handles explicit moves.
  • Uniqueness: during a deploy or crash, two servers must never both act as owner. That is a lease with fencing, enforced at the op log; see Leases and fencing tokens.
  • Hot documents: one owner handles one document's sequencing. The all-hands doc will test that.

Tradeoffs

ChoiceGainsCosts
Single owner per documentFree ordering, local broadcast, in-memory state.Document-aware routing and an ownership handoff protocol.
Any server plus a database sequencerStateless routing; any server can serve any client.A database round trip on a contended row for every op.

Stage 4 of 12 · Decide

Merging concurrent edits

Alice and Bob both see "The cat sat." Alice inserts "black " before "cat" at position 4. At the same moment, Bob deletes "sat" at positions 8-10. Each sends an op computed against the text they saw. Applied naively in the server's order, Bob's delete removes the wrong characters on Alice's machine. And offline users may send hours of such ops at once.

What you need to know first

Alice and Bob both edit "The cat sat." Alice inserts "black " at position 4. Bob deletes positions 8–10 ("sat"). Each op was computed against the text they saw.

Once Alice's insert has happened, everything after position 4 has shifted by 6 characters. Applied as-is, Bob's "delete 8–10" now removes the wrong characters.

After Alice's insert the text is "The black cat sat." Applied naively, Bob's delete of positions 8–10 (counting from 0) removes which characters?

"k c": positions 8–10 of "The black cat sat." are "k", " " and "c". The text becomes "The blacat sat.", and "sat" survives. Bob's delete needed to be shifted right by 6, to positions 14–16.

Two families of algorithms make concurrent edits converge. See Conflict resolution and convergence:

  • Operational transformation (OT): a central server transforms each incoming op against the ops it hadn't seen (shift Bob's delete by Alice's insert), then applies it in one agreed order.
  • CRDTs: every character gets a stable identity, so ops say "delete character #a17" instead of "delete position 8". Ops then commute: they give the same result in any order.

A user edits offline for three hours, then reconnects. Which approach handles that more naturally?

CRDTs: ops merge in any order, so a long divergence is just more ops to exchange.

OT can handle it, but transforming hours of ops against hours of others' ops is expensive and complex. CRDTs pay instead with per-character metadata.

What the stage asks

How should concurrent edits be merged?

  1. Flawed

    Last write wins on each paragraph

    Any two people editing the same paragraph at once lose one person's work. A finer grain only shrinks the window; for offline edits it is enormous.

  2. Defensible

    Lock a paragraph while someone is editing it

    It is simple and always convergent, and some products do this. But locks held by a client that went offline need leases and timeouts, collaborators keep hitting "Bob is editing", and offline editing is impossible by definition.

  3. Sound

    Operational transformation: the owner transforms each incoming op against the ops it had not seen, then sequences it

    With a single owner per document already in place, OT fits naturally: each op carries the sequence number it was based on, and the owner transforms it past everything since. Long offline sessions are the weak spot, because hours of divergent ops must be transformed against hours of others, which is expensive and historically bug-prone.

  4. Sound

    A sequence CRDT: every character has a stable identity, so ops commute and merge in any order

    Convergence no longer depends on ordering, so offline edits merge by simply exchanging ops. The owner still sequences ops for the log and catch-up, but correctness does not hinge on it. The cost is per-character metadata and tombstones that need compaction.

What a strong answer covers

  • Concurrent edits to the same region are normal, so the merge must preserve both rather than choose one.
  • Explains the mechanism: transformation against concurrent ops (OT) or commutative ops via stable identities (CRDT).
  • Addresses offline: long divergence is expensive for OT and natural for CRDTs.

The reasoning

  1. Concurrent edits computed against different versions must be merged, not chosen between.
  2. OT transforms ops against unseen ones in one agreed order; CRDTs make ops commute with stable identities.
  3. Long offline sessions favour CRDTs; both still benefit from server sequence numbers.

Both OT and CRDTs guarantee convergence. They differ in where the correctness lives: OT relies on the central sequencer to transform ops in one agreed order, while CRDTs build it into the data structure so that order does not matter; see Conflict resolution and convergence.

The rest of this investigation works with either, because both still benefit from a server that assigns each op a sequence number: it gives the durable log an order, gives reconnecting clients a position to resume from, and gives history a timeline.

What neither gives you is intent: two people rewriting the same sentence will converge on a mixture. That is a product problem (comments, suggestions, presence that shows who is where) rather than an algorithm problem.

Where another engineer could land differently

For structured data like form fields or task status, per-field last-write-wins is often fine, because a field is an atomic value rather than a sequence people type into together.

Stage 5 of 12 · Model

Trace one keystroke

Alice types a character. Follow the operation from her keyboard to Bob's screen. Getting the order right is what makes "no acknowledged edit is lost" true.

What you need to know first

Custody: at every moment, an edit must be held by someone who won't forget it: either the client's pending buffer (saved locally for offline use) or the durable log on the server.

The acknowledgement hands custody over. So the server may only ack once the log has the edit, and the client may only drop the edit from its buffer once it has the ack. See Durability.

Why does Alice's client apply her keystroke locally before the server has seen it?

So typing feels instant; the op stays in her pending buffer until the server acknowledges it.

Waiting for a round trip per keystroke would make typing feel sluggish. The pending buffer keeps the op safe meanwhile.

Alice never sees the ack for op a:812 and resends it. How does the owner avoid applying it twice?

Each op carries an id (client id + counter). The owner recognises a:812, finds the sequence number it already assigned, and acks with that instead of applying the op again.

What the stage asks

Order the life of a single operation.

In this order

  1. 1Alice's client applies the op to her local document immediately
  2. 2Client gives the op an id (client id + counter) and adds it to the pending buffer
  3. 3Client sends the op with the last sequence number it has seen
  4. 4Owner checks permissions and integrates the op into the in-memory document
  5. 5Owner appends the op to the log with the next sequence number; the batch commits
  6. 6Owner acks to Alice with the sequence number
  7. 7Alice's client removes the op from its pending buffer
  8. 8Owner broadcasts the sequenced op to the document's other clients
  9. 9Bob's client integrates the op and advances its last-seen sequence

The orderings that carry guarantees:

  • Apply locally first. Alice's keystroke appears instantly; the network never sits between her and her own typing. That is why ops must merge (OT/CRDT) rather than wait.
  • Durable before ack, ack before forgetting. The client keeps the op until the ack; the owner acks only after the commit. At every moment, at least one party durably holds the op.
  • Durable before broadcast. If Bob saw an op that was then lost in a crash, his document would contain text that exists nowhere else, and reconnecting would make the two diverge. Broadcasting after the commit costs a few milliseconds and rules that out.

The reasoning

  1. An edit is always in someone's custody: the client's pending buffer or the durable log.
  2. Ack only once the log has the op; the client drops it only after the ack.
  3. Op ids make resends idempotent.

The design rule is an invariant about custody: an op is always held by someone who will not forget it, either the client's pending buffer (persisted to IndexedDB for offline use) or the durable log. Custody is handed over by the ack, and the ack is only sent once the log has it; see Durability.

Op IDs make every step retryable: if Alice resends an op because she never saw its ack, the owner recognizes the ID and returns the existing sequence number instead of applying it twice; see Idempotency.

Stage 6 of 12 · Decide

Presence and cursors

Avatars show who is in the document; coloured cursors show where they are. Cursors move with every keystroke and every click. When someone closes their laptop, their avatar should disappear within seconds.

What you need to know first

Presence (who's here) and cursors (where they are) describe now. They're useless a minute later, and they rebuild themselves within seconds when clients reconnect and re-announce. That makes them soft state: keep them in memory, refresh them with heartbeats, expire them on a short TTL. See Soft state.

A laptop lid closes without a goodbye message. What removes that user's avatar?

Their heartbeats stop, and the presence entry expires on its short TTL.

A client that vanishes can't announce leaving. Expiry handles it without anyone needing to notice.

Cursors move with every keystroke, but only the latest position matters. So cursor updates can be coalesced (send only the newest every 50–100 ms), throttled, and sent at most once: a lost cursor update is replaced by the next one.

What the stage asks

Where should presence and cursor positions live?

  1. Flawed

    A Postgres row per user per document, updated on every cursor move

    Cursor moves are as frequent as keystrokes, so this doubles database writes for data nobody needs a minute later. And a crashed client never deletes its row, so ghosts need a cleanup job anyway, which is just expiry done the hard way.

  2. Sound

    In the owner's memory: refreshed by heartbeats, expired on a short TTL, broadcast throttled, never persisted

    Presence describes now. If the owner restarts, clients reconnect and re-announce within seconds, so durability buys nothing. Expiry handles clients that vanish without a goodbye. Cursor updates are coalesced (only the latest position matters) and sent at-most-once.

  3. Defensible

    Redis keys with a TTL, plus pub/sub to broadcast changes

    Correctly soft, and the right choice if there were no single owner per document. With an owner already holding every participant's connection, Redis adds a network hop and a dependency for information the owner already has.

  4. Flawed

    Record cursor moves in the op log so history shows where everyone was

    It bloats the durable log, the snapshots and every reconnect catch-up with data that is worthless seconds later. History should record what changed in the document, not where people looked.

What a strong answer covers

  • Presence describes the present; losing it on a crash costs nothing because clients re-announce.
  • It must expire automatically (heartbeat plus TTL), or disconnected users linger as ghosts.
  • Cursor updates can be throttled, coalesced and sent at-most-once, since only the latest matters.Supporting

The reasoning

  1. Presence and cursors are soft state: in memory, refreshed by heartbeats, expired by TTL.
  2. Only the latest cursor matters, so coalesce, throttle and send at most once.
  3. Keep edits durable and ordered; keep presence cheap and disposable.

The system now has two kinds of state with opposite needs:

EditsPresence and cursors
If lostData loss, never acceptableRebuilt within seconds by re-announcement
OrderOne sequence per documentOnly the latest value matters
StorageDurable log, group commitOwner's memory, TTL expiry
DeliveryAt-least-once, deduplicated by op idAt-most-once, coalesced

Treating them the same, in either direction, is a common and expensive mistake: durable presence wastes writes, and ephemeral edits lose work. See Soft state.

Stage 7 of 12 · Break it

Acknowledged, then lost

To reduce database load, an engineer changed the owner to buffer ops in memory and flush them to Postgres every two seconds. Find the design decisions that turned a crash into data loss.

What you need to know first

Group commit batches many writes into one transaction without weakening any of them: ops for a document accumulate for 10–20 ms, one transaction inserts them all, and then each op is acknowledged.

Throughput is the same as buffering longer; latency rises by a few milliseconds; and an ack still means "durable".

The owner acks each op immediately and flushes to Postgres every 2 seconds. The process is killed 1.5 s after the last flush. What's lost?

Up to 1.5 s of acknowledged ops, which clients already dropped from their buffers

The ack handed custody to a server that only had the ops in memory. Clients won't resend what they think is saved.

After the crash, Bob reconnects claiming last_seq = 5131, but the log only reaches 5120. What should the new owner do?

Not trust it. Bob has seen ops the log doesn't contain, so his document is ahead of the truth. Force a resync: send him the snapshot and log as they really are. Then he and everyone else converge on the same history.

What the stage asks

Select the lines where the design is at fault.

Logdoc 4711: owner and client logs
  1. 110:04:00.000 owner-3 flushed ops 5101-5120 to postgres
  2. 210:04:00.180 client-alice send op a:812 (base seq 5120)
  3. 310:04:00.182 owner-3 op a:812 → seq 5121 (buffered) ack → alice

    The ack is sent while the op exists only in memory. An acknowledgement must mean durable, so either flush before acking or ack when the batch commits (group commit).

  4. 410:04:00.183 client-alice ack a:812; removed from pending

    That is correct client behaviour, but it is why the early ack is fatal: custody passed to a server that had not secured the op.

  5. 510:04:00.183 owner-3 broadcast seq 5121 → bob, carol
  6. 610:04:01.402 owner-3 … ops 5122-5131 buffered and acked
  7. 710:04:01.950 owner-3 killed (OOM)
  8. 810:04:03.100 owner-5 takes ownership; loads log up to seq 5120
  9. 910:04:03.300 client-bob hello last_seq=5131 → owner-5 resumes Bob from 5131

    Bob claims a sequence number the log does not contain. The owner must detect last_seq > log head and force a resync, not trust it, or Bob's document permanently diverges.

  10. 1010:04:03.320 client-alice hello last_seq=5131, pending=[]

What the fix has to do

  • Ack only after the op is durably persisted; batching is fine if acks wait for the batch.
  • Clients keep ops until acked and resend on reconnect; the op id makes the resend idempotent.
  • On reconnect, the server compares the client's last seen position with the log and forces a resync if the client is ahead.
  • Broadcasting before durability let other clients see ops that were then lost.Supporting

The reasoning

  1. Ack only after the op is durable; group commit keeps batching without weakening the ack.
  2. Broadcast only after durability, so nobody sees ops that are later lost.
  3. Treat a client's last-seen position as a claim and resync if it's ahead of the log.

The batching was not the mistake. The early acknowledgement was. Group commit keeps both properties: ops for a document accumulate for ~10-20 ms, one transaction inserts them all, and then every op in the batch is acked and broadcast. Throughput is the same, the guarantee is intact, and latency rises by a few milliseconds; see Durability.

The second lesson is about trusting client positions. A client's last_seq is a claim. Before resuming, the owner checks it against the log's head. If the client is ahead (it saw ops that were lost), the owner sends a full resync: a snapshot plus the client's own unacknowledged ops to re-apply on top.

Stage 8 of 12 · Break it

Write the reconnect protocol

Implement the client side of reconnection. The client knows its last-seen sequence number and holds its pending ops, each with a unique id. The server can return every op after a given sequence number.

What you need to know first

Reconnection uses three things the design already has:

  • A position: the client's last-seen sequence number.
  • A replay: the log can return every op after that position, so the server doesn't need to remember the client. See Append-only logs.
  • Idempotent resends: pending ops carry their original ids, so any the server already sequenced are skipped.

Alice has 40 pending ops; the log has 25 ops from Bob she hasn't seen. How do her edits avoid overwriting Bob's?

Bob's ops are integrated and Alice's pending ops are rebased on top (transformed, or merged as CRDT ops).

Alice's ops were computed against an older document. The merge algorithm adjusts them so both sets of edits survive.

Two more rules: apply remote ops in sequence order and ignore any with seq ≤ lastSeq (duplicates). And reconnect with exponential backoff plus jitter, so thousands of clients don't reconnect in the same instant. See Retries, backoff and jitter.

What the stage asks

Implement onReconnect and onServerMessage so that nothing is lost, duplicated or applied out of order.

Reference implementation

class DocClient {
  lastSeq = 0;
  pending: Op[] = [];
  attempt = 0;

  onReconnect(socket: Socket) {
    this.attempt = 0;
    // The log is the buffer: ask for everything after our position,
    // and offer every op we still hold. Same ids as before.
    socket.send({ type: "hello", lastSeq: this.lastSeq, pending: this.pending });
  }

  onDisconnect() {
    const delay = Math.random() * Math.min(30_000, 500 * 2 ** this.attempt++);
    setTimeout(() => this.connect(), delay); // full jitter
  }

  onServerMessage(msg) {
    switch (msg.type) {
      case "op": {
        if (msg.op.seq <= this.lastSeq) return;            // duplicate delivery
        if (msg.op.seq !== this.lastSeq + 1) return this.requestFrom(this.lastSeq);
        const mine = this.pending.findIndex((p) => p.id === msg.op.id);
        if (mine >= 0) {
          // Our own op, sequenced (its ack may have been lost).
          this.pending.splice(mine, 1);
        } else {
          // Rebase: integrate the remote op, transforming pending ops past it
          // (OT), or merging by identity (CRDT, where order does not matter).
          this.doc.integrateRemote(msg.op, this.pending);
        }
        this.lastSeq = msg.op.seq;
        return;
      }
      case "ack": {
        this.pending = this.pending.filter((p) => p.id !== msg.opId);
        // lastSeq advances when the sequenced op itself arrives in order.
        return;
      }
      case "resync": {
        // The server's history differs from what we saw. Its log wins;
        // our unacknowledged work is re-applied on top.
        this.doc = Doc.from(msg.snapshot);
        this.lastSeq = msg.seq;
        for (const op of this.pending) this.doc.applyLocal(op);
        return;
      }
    }
  }
}
  • The server keeps no per-client queue. Catch-up reads the durable log after lastSeq, so it works across servers, restarts and arbitrarily long gaps (with a snapshot for very long ones).
  • Op ids make resends safe. The owner keeps a recent op-id index per document. If it already sequenced one of Alice's pending ops before the tunnel, it returns the existing sequence number instead of applying it again.
  • The client's own op returns as a sequenced op in the catch-up stream. The client recognizes it by id and drops it from pending, which covers the case where the ack was lost but the op was committed.
  • Gaps trigger a re-request, not a guess. Order is maintained by sequence number, never by arrival.

What a strong answer covers

  • On reconnect the client sends its last seen sequence and receives missed ops from the log, rather than relying on server memory.
  • Pending ops are resent with their original ids, so the server deduplicates any it had already sequenced.
  • Missed remote ops are integrated before or alongside the pending ops (transform or CRDT merge), so local edits are rebased rather than overwritten.
  • Remote ops are applied in sequence order; duplicates (seq <= lastSeq) are ignored.
  • Handles a resync by replacing state with the snapshot and re-applying pending ops on top.Supporting
  • Reconnects use exponential backoff with jitter.Supporting

The reasoning

  1. Resume from the client's last-seen sequence by replaying the log, not server memory.
  2. Resend pending ops with original ids and rebase them over missed remote ops.
  3. Apply remote ops in order, ignore duplicates, and back off with jitter.

Reconnection is where the earlier decisions pay off. Sequence numbers give the client a position; the durable log turns that position into a replay; op ids make every resend idempotent; the merge algorithm turns 40 local ops and 25 remote ones into one document. None of it needs the server to remember anything about Alice personally; see Append-only logs.

Stage 9 of 12 · Break it

Deploy day

Deploys are the most common "failure" this system will ever see, and they happen every day. Evaluate each statement about getting through one.

What you need to know first

Load-balancer draining waits for in-flight requests to finish before stopping a server. A WebSocket is one request that never finishes, so draining alone just waits out the timeout and cuts it.

The server has to hand off actively: stop accepting ops, flush pending batches, release ownership, and tell clients to reconnect.

20,000 clients reconnect with a random delay spread evenly over 10 seconds. About how many reconnections a second do the new servers face?

About 2,000 per second.

20,000 ÷ 10 = 2,000 a second, instead of 20,000 handshakes, permission checks and catch-up reads in the same instant. Jitter turns a spike into a ramp.

During handoff, two servers can briefly both believe they own a document: the old one mid-flush, the new one starting. Ownership is a lease with an epoch number that increases with each new owner.

Appends to the log are conditioned on the owner's epoch, and (doc_id, seq) is unique, so a stale owner's write fails and it learns it lost.

Should presence be saved before a deploy so avatars survive it?

No: clients reconnect and re-announce within seconds; presence rebuilds itself.

That's the benefit of treating presence as soft state.

What the stage asks

Which statements hold?

  1. Fails

    Load-balancer connection draining is enough: WebSocket connections will finish on their own within the 30-second drain window.

    Draining waits for requests to complete, and a WebSocket never completes. The server must actively hand off: stop accepting, flush pending batches, release ownership, and tell clients to reconnect.

  2. Holds

    If every client reconnects the instant its socket closes, the new servers can be overwhelmed.

    20,000 simultaneous handshakes, permission checks and catch-up reads arrive at once. Clients should reconnect after a random delay; servers can also stagger closing their connections. See Retries, backoff and jitter.

  3. Holds

    During a handoff, the old and new servers can briefly both believe they own a document.

    The old owner may be mid-flush when the new owner's lease begins, or paused. Ownership is a lease, and leases can overlap in belief even when they do not overlap in time.

  4. Holds

    A unique constraint on (doc_id, seq) in the op log, plus appends conditioned on the owner's epoch, prevents a stale owner from forking history.

    The stale owner's insert either collides with a sequence number the new owner already used, or fails the epoch check. Either way it learns it lost, and its unacknowledged ops will be resent by clients to the real owner.

  5. Fails

    Presence must be persisted before the deploy so avatars survive the restart.

    Clients reconnect and re-announce within seconds; presence rebuilds itself. That is the benefit of treating it as Soft state.

The reasoning

  1. Draining doesn't end WebSockets: hand off actively and tell clients to reconnect.
  2. Reconnect with jitter so new servers see a ramp, not a spike.
  3. Fence ownership with epochs and a unique (doc, seq) so a stale owner can't fork history.

A graceful handoff for each document:

  1. The old owner stops accepting ops for the document and commits any buffered batch (acking those ops).
  2. It releases its lease; the new owner acquires it with a higher epoch.
  3. Clients receive "reconnect" and back off with jitter; the router sends them to the new owner.
  4. The new owner loads the latest snapshot and log tail. Clients resume from their lastSeq and resend pending ops.

And for the ungraceful version (a crash), the same steps happen with step 1 skipped: the lease expires instead of being released, and the epoch fence makes sure a zombie owner cannot write. Deploys exercise the crash-recovery path daily, which is the best way to keep it working.

Stage 10 of 12 · Change it

The all-hands document

The single-owner design made ordering free. Now one document's traffic exceeds what one process can deliver. Sequencing is still cheap (200 ops/s); delivery is not (about 46,000 messages/s).

What you need to know first

Find the part of the work that must be serialized and keep only that part serialized. Here:

  • Sequencing (assigning order) must be single-writer, and it's cheap: about 200 ops a second.
  • Delivery (sending to 230 sockets) is expensive, and parallel by nature.

Instead of sending each op separately, each recipient gets one frame every 50 ms containing all new ops. With 230 recipients, how many frames a second?

About 4,600 frames per second.

20 frames a second × 230 recipients = 4,600 frames a second, down from about 46,000 individual messages. Same ops, a tenth of the overhead.

A viewer on a bad connection reads slower than ops arrive. What should the server do with their send buffer?

Bound it; when it fills, drop the buffered ops and let the client catch up from the log by sequence number

An unbounded buffer turns one slow viewer into growing memory on the owner. The log makes dropping safe: the client asks for what it missed.

What the stage asks

How do you keep this document healthy?

  1. Defensible

    Move hot documents to a larger server

    It buys headroom quickly and is a reasonable stopgap. It does nothing about slow viewers' buffers growing, and the next, bigger document hits the same wall.

  2. Flawed

    Split the document's editors across several servers, each sequencing its own clients' ops

    Two sequencers mean two histories. Sequence numbers stop defining a single order, catch-up positions become ambiguous, and the log's uniqueness invariant has to be abandoned. Even with a CRDT, where content would converge, you lose the single timeline that history and resumption depend on.

  3. Sound

    Batch ops into frames every ~50 ms per recipient, bound every connection's send buffer, and move viewers who exceed it to catch-up-from-log mode

    Batching cuts per-message overhead by an order of magnitude, from ~46,000 sends a second to ~4,600 frames. Bounded buffers turn a slow viewer from a memory leak into a client that briefly falls behind and catches up from the log by sequence number, the same path as a reconnect.

  4. Sound

    Keep one owner for sequencing; publish sequenced ops to pub/sub, and serve viewers from separate edge servers that subscribe

    Sequencing and delivery are different jobs with different scaling needs. The owner does the cheap, single-writer part; any number of edge servers do the expensive, embarrassingly parallel part. Viewers lose nothing: ops carry sequence numbers, so edges can detect gaps and catch up from the log.

What a strong answer covers

  • Sequencing must stay single-writer, but fan-out delivery can be distributed.
  • Per-connection buffers must be bounded; slow clients fall back to catching up from the log.
  • Batching or coalescing reduces per-message overhead.Supporting

The reasoning

  1. Keep only the necessary part serialized: sequencing stays single-writer, delivery scales out.
  2. Batch ops into frames per recipient to cut per-message overhead.
  3. Bound send buffers; slow clients fall back to catching up from the log.

The general move is to find the part of the work that must be serialized, keep only that part serialized, and scale everything else out. Here the serialized part (assigning sequence numbers) is tiny; the expensive part (delivering bytes to 230 sockets) is parallel by nature.

Both sound options compose. Batching and bounded buffers make each server efficient and safe; edge fan-out spreads delivery across servers. Pub/sub can be ephemeral here, because a dropped message shows up as a sequence gap that the edge repairs from the durable log; see Publish/subscribe and Backpressure and capacity.

Stage 11 of 12 · Change it

History and fast loads

The op log is the source of truth, and replaying it from the start is how the owner rebuilds a document. That cost grows forever.

What you need to know first

When the log is the source of truth, loading a document means replaying it, and that cost grows with the document's whole life. A snapshot stores the document as it was at a given sequence number. Loading becomes: latest snapshot, then replay only the ops after it. See Append-only logs.

Snapshots are taken every 2,000 ops. At most how many ops does a load replay after the latest snapshot?

About 2,000 ops.

At most 2,000, compared with 2.3 million from the beginning. Load time is now bounded, whatever the document's age.

Why must each snapshot record exactly which sequence number it includes?

So replay starts at the next op, with no gap and no op applied twice

A snapshot without its position can't be combined with the log safely.

Snapshots never change once written, so they suit Object storage: keyed by {doc}/{seq}, written once, cached forever.

What the stage asks

How should documents be loaded and history served?

  1. Flawed

    Keep replaying the full log, but make apply faster

    Load time stays proportional to the document's entire life. Optimizing the constant factor does not change the slope.

  2. Sound

    Periodically write snapshots tagged with the sequence number they include; load = latest snapshot + ops after it

    Load cost becomes bounded: one object fetch plus at most a few thousand ops. The sequence number on the snapshot is essential: it says exactly where replay resumes. Snapshots at coarser intervals for older periods double as version history.

  3. Flawed

    Store only the latest document state and delete the log

    Loading is fast, but reconnecting clients can no longer catch up by position, version history is gone, and a bug in the in-memory state is now permanent with nothing to rebuild from.

  4. Defensible

    Write a new snapshot after every op

    Loads are instant, but every keystroke now writes the whole document: 50 KB × thousands of ops a second. Snapshotting every few hundred ops or few minutes gets nearly the same load time at a fraction of the cost.

What a strong answer covers

  • A snapshot must record exactly which sequence number it reflects, so replay starts at the next op with no gap or duplicate.
  • The log after the latest snapshot remains essential for catch-up and correctness.
  • Older history can be compacted into coarser snapshots for the version sidebar.Supporting

The reasoning

  1. A snapshot is a cached fold of the log, keyed by the sequence number it includes.
  2. Load = latest snapshot + ops after it, so load time stays bounded.
  3. Keep the log after the latest snapshot for catch-up and correctness.

Snapshots are a cache of the log's fold, with the cache key being the sequence number. Like any cache they can be rebuilt from the source of truth, which is why the log after them must be kept, and why a snapshot without its position is useless; see Append-only logs.

Snapshots are immutable objects keyed by {doc}/{seq}, which makes them a natural fit for Object storage: write once, cache forever, never update in place.

Stage 12 of 12 · Defend it

Defend the guarantees

In the design review, someone asks: "Walk me through why every client ends up with the same document, and why we never lose an edit someone saw as saved. What failure would break each guarantee?"

What you need to know first

Three ideas carry this design, and the same three appear in jobs and payments:

  • A single authority per unit of state: one owner per document, fenced by an epoch.
  • Custody transferred only once durable: ack after commit; clients hold ops until acked.
  • Idempotent application of anything repeatable: op ids, sequence numbers.

Defending the guarantees means naming the failure that would break each one.

Which failure would break the 'no acknowledged edit is lost' guarantee?

Acking an op before it's committed to the log

The client drops it on the ack. A crash before the commit then loses it with nobody holding a copy.

What the stage asks

Explain the convergence and durability guarantees, the mechanism behind each, and what would have to fail to break them.

Reference answer

Convergence. Every op has a unique id and is integrated through a convergent merge (OT transformed in the owner's order, or a CRDT whose ops commute). Duplicate deliveries are ignored by id. So any two clients that have applied the same set of ops show the same text, regardless of when or how often they received them.

One history. Each document has exactly one owner at a time, holding a lease with an epoch. Appends to the op log are conditioned on that epoch, and (doc_id, seq) is unique, so even a zombie owner cannot fork the history. Sequence numbers therefore mean the same thing to everyone.

No acknowledged edit is lost. Custody is explicit: the client keeps every op (persisted locally) until it is acknowledged; the owner acknowledges only after the batch containing it commits. A crash at any point leaves the op with the client, who resends it with the same id, or in the log, from which it is replayed.

What would break them: acking before commit (the incident in this investigation); a second sequencer without fencing; a client clearing its pending buffer before the ack, or losing its local storage while offline; a bug in the merge algorithm, which is why production systems periodically compare document checksums between clients and server.

What is not promised: that the merged text matches what both authors meant when they edit the same words, that edits appear within 200 ms during a handoff, or that presence is exact.

What a strong answer covers

  • Convergence: ops are merged by a convergent algorithm (server-ordered OT or CRDT) and applied idempotently by op id.
  • One history per document: a single owner with an epoch-fenced lease and a unique (doc, seq) constraint.
  • Durability: acks are sent only after commit, and clients hold ops until acked and resend them on reconnect.
  • Names what would break each, e.g. acking before commit, a second sequencer without fencing, a client losing its pending buffer before reconnecting, a merge bug.
  • Distinguishes what is not guaranteed: intent preservation, latency during handoffs, presence accuracy.Supporting

The reasoning

  1. Convergence comes from a convergent merge plus idempotent application by op id.
  2. Durability comes from acking only after commit and clients holding ops until acked.
  3. Name what breaks each guarantee: early acks, an unfenced second sequencer, a lost pending buffer.

The same three ideas from the other investigations carried this system: a single authority per unit of state (the document owner, like the job lease and the payment attempt), custody that is only transferred once durable, and idempotent application of anything that can be repeated. The technologies were different; the reasoning was not.

How you did

Now try it as an interview question

  • “Design Google Docs.”
  • “Design a multiplayer whiteboard or Figma-like editor.”
  • “How would you add live presence and cursors to an existing app?”
  • “Design a chat system that works offline and syncs when reconnected.”
  • “Your WebSocket servers need to be redeployed daily. How do you do it without disrupting users?”

The interview mode mixes stages from this and other investigations with concept recall and questions about your own projects.

Back to the last stage