Designing a Collaborative Editing Platform (Google Docs / Notion)
Difficulty: Advanced Topics: Real-Time Collaboration, OT/CRDT, WebSocket, Conflict Resolution, Presence Awareness Asked at: Google, Notion, Figma, Microsoft, Amazon Prerequisites:WebSockets and Consistency Models
1. Understanding the Problem
A collaborative editing platform lets multiple users simultaneously edit the same document in real-time - seeing each otherβs cursors, changes appearing character-by-character as they type, without any user overwriting anotherβs work. The hard part? When two users type at the same position in the same millisecond, you need a deterministic way to merge both changes without data loss, all while maintaining sub-100ms latency so typing feels instant.
2. Naive First Cut
flowchart LR
User1["User 1 Browser"]:::client
User2["User 2 Browser"]:::client
API["API Server"]:::service
DB["Postgres DB"]:::data
User1 --> API
User2 --> API
API --> DB
classDef client fill:#4c3a5e,stroke:#818cf8,color:#e2e8f0
classDef service fill:#1a3a2a,stroke:#4ade80,color:#e2e8f0
classDef data fill:#3b3520,stroke:#fbbf24,color:#e2e8f0
| Color | Meaning |
|---|---|
| π Purple-Orange | Client apps |
| π΅ Blue | Edge / Gateway |
| π’ Green | Backend services |
| π‘ Yellow | Data stores |
| π£ Purple | Async (Kafka / Pub/Sub) |
| π΄ Pink | External services |
How this breaks:
- Last-write-wins in Postgres means User 1βs edits silently disappear when User 2 saves - data loss
- No way to push changes to other users in real-time - polling every second is too slow and wasteful
- Storing the full document on every keystroke creates massive write amplification (100 chars/min Γ millions of docs)
- No conflict resolution - concurrent edits at the same position produce garbled text
- Single API server canβt maintain persistent connections for millions of active editors
- No presence awareness - users have no idea others are editing, leading to conflicting sections
The rest of the doc evolves this into a production-grade real-time collaborative editing system using operational transforms and persistent WebSocket connections.
3. Prior Art Weβre Drawing From
- Google Wave OT - Pioneered Operational Transformation for real-time collaborative editing. Jupiter protocol with a central server that transforms concurrent operations to maintain consistency. (Google Research)
- Figma CRDT - Uses a custom CRDT (Conflict-free Replicated Data Type) for multiplayer design editing with no central coordination for conflict resolution. Demonstrated CRDTs can work at scale with structured documents. (Figma Engineering Blog)
- Yjs - Open-source CRDT framework used by Notion, JupyterLab, and others. YATA algorithm for text sequences with O(1) amortized insertion and efficient encoding. (Yjs GitHub)
- Automerge - Research CRDT library that treats documents as mergeable JSON structures. Demonstrates how CRDTs handle offline editing and eventual convergence. (Ink and Switch)
- Google Docs Jupiter Protocol - Server-mediated OT where each client sends operations to a central server that transforms and broadcasts them. Single point of serialization ensures total ordering. (Operational Transformation FAQ)
4. Functional Requirements
Core (Top 3)
- Real-time collaborative editing - multiple users type simultaneously in the same document with changes appearing within 100ms for all participants
- Document storage and retrieval - create, open, and save documents; persist all content durably
- Version history - view previous versions of a document, restore any earlier state, see who made what changes
Below the Line
- Rich text formatting (bold, italic, headings, lists)
- Comments and suggestions (track changes)
- Offline editing with sync on reconnect
- Access control and sharing permissions
- Real-time presence indicators (cursors, selections)
- Document templates and search
5. Non-Functional Requirements
Core
| NFR | Target |
|---|---|
| Edit propagation latency | < 100ms from one userβs keystroke to appearing on another userβs screen |
| Consistency | Eventual consistency - all users converge to the same document state regardless of operation order |
| Availability | 99.99% - document editing must never be βdownβ during work hours |
| Scale | Support 100M+ documents with up to 50 concurrent editors per document |
Below the Line
- Version history retained for 30+ days
- Document size up to 1M characters
- Support 10M+ daily active users across all documents
- Sub-second document open time (cold start)
6. Scale Estimation (Back-of-Envelope)
- Users: 100M DAU, 5M concurrent editing sessions at peak, spread over ~2M open documents
- Write QPS: a typist emits ~5 ops/sec, and roughly 20% of open sessions are actively typing at any instant, so ~1M ops/sec sustained across all documents
- Hot document: ~50 concurrent editors is the realistic ceiling, giving ~250-500 ops/sec on a single document
- Read QPS: 50K document opens/sec (initial state load + presence sync)
- Storage: ~10TB of current document state (500M documents Γ avg 20KB), before version history
- Bandwidth: ~10 Gbps (1M ops/sec Γ ~100 bytes, fanned out to ~2.5 collaborators each, plus framing and presence)
The numbers that shape the design: 1M ops/sec to serialize, 5M live connections to hold, and a 100ms budget for an edit to reach the other editors.
7. Core Entities
- Document - id, title, owner, permissions, current version number, created/updated timestamps
- Operation - documentId, version number, userId, type (insert/delete), position, content, timestamp
- Session - documentId, userId, WebSocket connection reference, cursor position, selection range
- Version Snapshot - documentId, version number, full document content, timestamp (periodic checkpoint)
- User - id, name, email, avatar, color assignment for presence
8. API / System Interface
POST /api/v1/documents
Body: { title, content?: "" }
Response: { documentId, version: 0, createdAt }
Auth: JWT Bearer token
Note: Creates an empty document
GET /api/v1/documents/{docId}
Response: { documentId, title, content, version, collaborators[] }
Auth: JWT Bearer token
Note: Returns latest snapshot + any pending ops after snapshot
WebSocket /ws/v1/documents/{docId}/edit
Client sends: { type: "op", ops: [{ type: "insert", pos: 12, content: "hello" }], version: 42 }
Server broadcasts: { type: "op", ops: [...transformed...], userId, version: 43, serverTimestamp }
Client sends: { type: "cursor", pos: 17, selectionEnd?: 25 }
Server broadcasts: { type: "presence", userId, name, color, pos, selectionEnd }
Auth: JWT ticket in connection params
GET /api/v1/documents/{docId}/history
Query: ?fromVersion=1&toVersion=100
Response: { versions: [{ version, userId, timestamp, summary }] }
Auth: JWT Bearer token
POST /api/v1/documents/{docId}/restore
Body: { targetVersion: 42 }
Response: { documentId, newVersion: 101, content }
Auth: JWT Bearer token (owner/editor only)
9. High-Level Design
FR1: Real-Time Collaborative Editing
The requirement is that two people can type in the same document and each sees the otherβs changes as they happen. Read plainly, that needs somewhere to hold the document, a live channel to each editor, and a rule for what to do when edits arrive.
Build the simplest version of that. One service, holding the document in memory, applying edits in the order they land and echoing each one to the other editors. Everything else on this page β the transform algorithm, the operation log, the pub/sub fabric, the sharded worker pool β answers a non-functional requirement, so each belongs in a deep dive.
New components we need:
- API Gateway - Entry point for HTTP requests (document CRUD). Handles auth and routing.
- Collaboration Service - One process. It holds the current text of every open document in memory, holds the WebSocket for every editor of those documents, applies each incoming edit to its copy, and echoes the edit to the other editors of that document.
π‘ WebSocket = a persistent connection that stays open so the server can push an edit the instant it happens rather than waiting to be asked. Learn more β
flowchart LR
U1["User 1 Browser"]:::client
U2["User 2 Browser"]:::client
GW["API Gateway"]:::edge
CS["Collaboration Service"]:::service
MEM["In-memory doc text<br/>and editor sockets"]:::data
U1 -->|"1. Send edit over WebSocket"| CS
U2 -->|"2. Send edit over WebSocket"| CS
GW -->|"3. Auth on connect"| CS
CS -->|"4. Apply edit in arrival order"| MEM
CS -->|"5. Echo edit to user 2"| U2
CS -->|"6. Echo edit to user 1"| U1
classDef client fill:#4c3a5e,stroke:#818cf8,color:#e2e8f0
classDef edge fill:#1e3a5f,stroke:#60a5fa,color:#e2e8f0
classDef service fill:#1a3a2a,stroke:#4ade80,color:#e2e8f0
classDef data fill:#3b3520,stroke:#fbbf24,color:#e2e8f0
Step-by-step flow:
- User 1 opens the document. The Gateway authenticates them, and the Collaboration Service loads the document text into memory and keeps their WebSocket
- User 1 types βhelloβ at position 12 β the client sends
{type: "insert", pos: 12, content: "hello"} - The Collaboration Service applies it to its in-memory copy at position 12
- It writes the same edit down the socket of every other editor of that document
- Each of those clients applies the edit to its own local copy at position 12
- User 1βs client already applied the edit locally when they typed it, so the round trip is confirmation, not news
Why send edits rather than the whole document? A keystroke is a handful of bytes; a document is tens of kilobytes. At 1M ops/sec, broadcasting whole documents would be ~20 GB/sec of fan-out to communicate a few megabytes of actual change. Sending the edit is not an optimization we are deferring, it is the only version that works at all.
What we have deliberately left broken. For two people taking polite turns, this is correct and complete. It has three holes, and the first one is a correctness bug, not a performance one:
- Step 3 applies edits in arrival order, and step 5 replays them at face value. When two edits are genuinely concurrent, the positions in them refer to different versions of the text, so the clients do not just lose an edit β they end up with different documents and never reconverge. That is Deep Dive 1, and it is the heart of this problem.
- One process holds every connection and every open document. 5M live WebSockets do not fit on one box, and the moment there are two boxes, step 4 cannot reach an editor connected to the other one. Routing across a fleet is Deep Dive 2; deciding where each documentβs edit state lives is Deep Dive 4.
- The document exists only in memory. A crash or a deploy loses everything anyone typed. That is what FR2 is for.
FR2: Document Storage and Retrieval
FR1 keeps the document in memory, which means a crash loses it. Documents need to survive the browser closing, the process restarting, and a year of nobody opening them.
In simple terms: You close the tab and come back tomorrow. Your document has to still be there, exactly as you left it.
New components we need:
- Document Service - Handles document CRUD: create, open, list, delete, and permissions.
- Document DB (Postgres) - One row per document: id, title, owner, permissions, and the documentβs current content.
flowchart LR
User["User Browser"]:::client
GW["API Gateway"]:::edge
DS["Document Service"]:::service
DB[("Postgres documents")]:::data
CS["Collaboration Service"]:::service
User -->|"1. Open document"| GW
GW -->|"2. Forward to doc svc"| DS
DS -->|"3. Read row with content"| DB
DS -->|"4. Return document"| User
CS -->|"5. Save content after each edit"| DB
classDef client fill:#4c3a5e,stroke:#818cf8,color:#e2e8f0
classDef edge fill:#1e3a5f,stroke:#60a5fa,color:#e2e8f0
classDef service fill:#1a3a2a,stroke:#4ade80,color:#e2e8f0
classDef data fill:#3b3520,stroke:#fbbf24,color:#e2e8f0
Step-by-step flow:
- User opens a document β
GET /documents/{docId}hits the Document Service - Document Service checks the callerβs permission on that document, then reads the row
- It returns the title and the full current content in one read. There is nothing to assemble
- The client renders it and opens a WebSocket to the Collaboration Service for live editing
- As edits come in, the Collaboration Service writes the updated content back to the same row, so the persisted copy tracks the in-memory one
Why store the content on the document row? Because opening a document is then a single primary-key read, which is the cheapest thing a database does, and it is genuinely the right answer while documents are small and edits are infrequent. The read path here is hard to beat β it is the write path that will not hold.
What we have deliberately left broken. Two things, and they are related:
- Step 5 rewrites the whole document on every edit. A 20KB document rewritten per keystroke, at 1M ops/sec across the fleet, is roughly 20 GB/sec of database writes to express maybe a few MB/sec of actual change. Postgres will not do that, and no amount of hardware makes rewriting 20KB to change one character a reasonable idea.
- Overwriting means the past is gone. Each save destroys the previous content, so there is no answer to βwhat did this say yesterdayβ β which is exactly what FR3 has to deliver next.
Both point the same way, and Deep Dive 3 is where the operation log and periodic snapshots earn their place.
FR3: Version History
Users need to see what the document looked like at a point in the past, who changed it, and be able to go back to it. FR2 overwrites the content on every save, so right now none of that is possible.
In simple terms: Your boss asks βwhat did this document say last Tuesday?β Something has to have kept a copy.
New components we need:
- History Service - Lists the saved versions of a document, returns any one of them for preview, and restores one as the current content.
The storage side needs one new table rather than a new datastore: document_versions,
holding the full content at each save point along with who saved it and when. FR2βs
documents row stays as the current state; this table is the trail behind it.
flowchart LR
User["User Browser"]:::client
GW["API Gateway"]:::edge
HS["History Service"]:::service
CS["Collaboration Service"]:::service
VER[("document_versions")]:::data
DOC[("documents")]:::data
User -->|"1. View version history"| GW
GW -->|"2. Forward to history svc"| HS
HS -->|"3. List versions for doc"| VER
HS -->|"4. Read one version content"| VER
HS -->|"5. Restore writes current row"| DOC
CS -->|"6. Append a version on save"| VER
classDef client fill:#4c3a5e,stroke:#818cf8,color:#e2e8f0
classDef edge fill:#1e3a5f,stroke:#60a5fa,color:#e2e8f0
classDef service fill:#1a3a2a,stroke:#4ade80,color:#e2e8f0
classDef data fill:#3b3520,stroke:#fbbf24,color:#e2e8f0
Step-by-step flow:
- Instead of only overwriting the
documentsrow, the Collaboration Service also appends a row todocument_versionson a periodic save β say every few minutes of activity, or when the last editor closes the document - User clicks βVersion Historyβ β
GET /documents/{docId}/historyreturns the list: version number, author, timestamp - User picks version 450 β the History Service reads that row and returns its content directly. No reconstruction is needed, because we stored the whole thing
- The UI renders it as a read-only preview
- If the user clicks βRestore,β the History Service writes that content back as the current
documentscontent and appends it as a new version, so restoring is itself an undoable event - Active editors are told the document was replaced and reload it
Why is restore a new version rather than a rollback? Because rolling back by deleting versions makes the restore itself unrecoverable. Appending keeps the history a strictly growing record, which means a mistaken restore is one more restore away from being fixed.
What we have deliberately left broken. This satisfies the requirement and it is expensive and coarse:
- Every version is a full copy. A 1MB document saved 1,000 times over six months is ~1GB for one document, almost all of it identical bytes repeated.
- You can only travel to a save point. βWhat did it say at 10:32β has no answer unless a save happened to land at 10:32. Everything typed between saves is invisible to history, and fine-grained undo across a collaborative session is not expressible at all.
- The timeline is per-save, not per-change. We can say βAlice saved at 14:05,β not βAlice rewrote the pricing section,β because a full-content row does not record what moved.
All three come from storing states instead of changes, and Deep Dive 3 is where that inverts.
10. Technology Choices
| Tier | Purpose | Stores | Access Pattern | Primary | Alternatives |
|---|---|---|---|---|---|
| Document Store | Persistent document content | Full document snapshots + metadata | Read/write by docId | Postgres (or Spanner) | CockroachDB, TiDB |
| Operation Log | Ordered stream of edit operations | Insert/delete ops with position and version | Append-only, read by docId + version range | Cassandra (or DynamoDB) | ScyllaDB, FoundationDB |
| Real-time Relay | Push operations to connected editors | WebSocket messages | Fan-out per document session | WebSocket Gateway + Redis Pub/Sub | SSE, gRPC streaming |
| Presence Cache | Active cursors and selections | userId, cursor position, color | High-QPS reads/writes, TTL-based | Redis Cluster | Memcached |
| Event Bus | Async events (save, snapshot, version) | Document lifecycle events | Pub/sub per document | Kafka or Redpanda | Kinesis, Pub/Sub |
| Object Store | Version snapshots and exports | Document snapshots, PDF exports | Batch writes, occasional reads | S3 | GCS, MinIO |
| Search Index | Full-text document search | Document titles and content | Text search by user | Elasticsearch | Typesense, Meilisearch |
Why a central OT server, not pure CRDT? Pure CRDTs (like Yjs or Automerge) work great for peer-to-peer scenarios and offline editing. But for a Google Docs-style product where we need a canonical server-side version, fine-grained access control, and version history, server-mediated OT gives us a single serialization point. This makes snapshotting, permissions, and undo/redo simpler. The tradeoff: the server is on the critical path for every operation. We mitigate this with per-document sharding.
Why Cassandra for the operation log? Operations are append-only and partitioned by documentId. We need fast sequential reads (replay ops from version X to Y) and high write throughput. Cassandraβs partition-key based access and LSM storage handle this naturally. We never update or delete individual operations.
11. Data Modeling
Postgres / Spanner (Document Store β persistent content):
CREATE TABLE documents (
doc_id UUID PRIMARY KEY,
owner_id UUID NOT NULL,
title VARCHAR(500),
current_version BIGINT NOT NULL DEFAULT 0,
content TEXT, -- latest full snapshot (periodically rebuilt from ops)
created_at TIMESTAMP NOT NULL,
updated_at TIMESTAMP NOT NULL
);
CREATE TABLE doc_permissions (
doc_id UUID NOT NULL,
user_id UUID NOT NULL,
role VARCHAR(10) NOT NULL, -- owner, editor, viewer
PRIMARY KEY (doc_id, user_id)
);
Cassandra / DynamoDB (Operation Log β ordered stream of edits):
Table: operations
PK: doc_id
SK: version_number (ascending, monotonic)
Columns: user_id, op_type (insert/delete/retain), position, content, timestamp, client_id
Redis (Presence Cache β active cursors):
Key: "doc:presence:{docId}" β Hash { userId β JSON { cursor_pos, selection, color, name } }
TTL: 30s per field (refreshed on every cursor move)
Key: "doc:version:{docId}" β Integer (latest version number for OT/CRDT conflict detection)
Access Patterns:
| Query | Data Source | How |
|---|---|---|
| Open document | Postgres | SELECT content, current_version FROM documents WHERE doc_id = ? |
| Apply edit (real-time) | Cassandra + Redis | Transform op against doc:version:{docId}, append to operations log, broadcast via Pub/Sub |
| Load recent edits (sync on reconnect) | Cassandra | SELECT * FROM operations WHERE doc_id = ? AND version_number > lastSeen |
| Show active collaborators | Redis | HGETALL doc:presence:{docId} |
| Snapshot for recovery | S3 | Periodic: rebuild document from ops, store full snapshot |
How Concurrent Edits Are Resolved (Operational Transform):
- User A types βhelloβ at position 5 β sends
{op: INSERT, pos: 5, content: "hello", baseVersion: 42}via WebSocket - Server checks: is baseVersion == current doc version? If yes: apply directly, increment version to 43, broadcast to all clients
- If no (another edit arrived first): transform Aβs operation against the conflicting op. E.g., if User B inserted 3 chars at position 2 (version 42β43), then Aβs position shifts from 5 to 8
- Transformed op is applied, stored in
operationslog with version 44, broadcast to all clients - Every 100 operations: a background job rebuilds the full document text from ops and writes a snapshot to Postgres + S3 (compaction)
12. Deep Dives
1) Two people type at the same position in the same instant. Whose character survives?
Problem: FR1 applies edits in arrival order and echoes them to the other editors verbatim. A position in an edit is only meaningful against the version of the text the author was looking at, and by the time the edit is applied elsewhere that text has moved.
Bad: apply-in-arrival-order and echo the raw edit, exactly as FR1 does. It is worth walking one example all the way through, because the failure is worse than βan edit gets lost.β
The document is Hello World. Alice inserts , at position 5, intending Hello, World. At the same moment Bob deletes position 3 (the second l), intending Helo World. Neither has seen the otherβs edit.
Server, Bob's delete arrives first:
"Hello World" -> delete pos 3 -> "Helo World"
"Helo World" -> insert , at 5 -> "Helo ,World"
Alice's client, her own edit applied locally first:
"Hello World" -> insert , at 5 -> "Hello, World"
"Hello, World" -> delete pos 3 -> "Helo, World"
Bob's client, his own edit applied locally first:
"Hello World" -> delete pos 3 -> "Helo World"
"Helo World" -> insert , at 5 -> "Helo ,World"
Alice is now looking at Helo, World. Bob is looking at Helo ,World. The server holds a third opinion that happens to match Bobβs. Nothing is retried and nothing reconciles, so those three copies stay different for the rest of the session, and every subsequent edit is applied against a different base on each machine β so they diverge further, not less. Two people editing a shared document now cannot see the same document, which is the one thing the product exists to do.
This is not a load problem. It happens with two users on a fast network, and at our scale of 1M ops/sec the concurrent-edit window is being entered constantly rather than occasionally.
Good: Lock-based editing. Lock a paragraph or section while someone is typing in it, so concurrent edits to the same region cannot happen. This does fix correctness β but it fixes it by removing the feature. Users see βsection locked by Bobβ instead of fluid typing, and the whole point was simultaneous editing.
Great: Server-mediated Operational Transformation using the Jupiter/dOPT protocol. (Borrowing from Google Wave and Google Docs.)
In simple terms: When two users type at the same time, their edits might conflict (User A inserts at position 5, User B deletes at position 3 β now position 5 is wrong). The server βtransformsβ each operation based on what happened before it, so both users end up with the same correct document.
How OT works:
The core OT algorithm maintains a server version counter and transforms each incoming operation against all operations that happened between the clientβs base version and the serverβs current version.
Transform rules for text operations:
- insert(pos, char) vs insert(pos2, char2):
if pos < pos2: insert stays at pos (the other insert is after us)
if pos > pos2: insert shifts to pos+1 (the other insert pushed us right)
if pos == pos2: break tie by userId (deterministic)
- insert(pos, char) vs delete(pos2):
if pos <= pos2: insert stays at pos
if pos > pos2: insert shifts to pos-1 (the deleted char pulled us left)
- delete(pos) vs insert(pos2, char):
if pos < pos2: delete stays at pos
if pos >= pos2: delete shifts to pos+1
- delete(pos) vs delete(pos2):
if pos < pos2: delete stays at pos
if pos > pos2: delete shifts to pos-1
if pos == pos2: delete becomes no-op (already deleted)
flowchart LR
C1["Client 1"]:::client
C2["Client 2"]:::client
OT["OT Engine"]:::service
VER["Version Manager"]:::service
OL["Op Log"]:::data
C1 -->|"1. Op + baseVersion"| OT
C2 -->|"2. Op + baseVersion"| OT
OT -->|"3. Assign version"| VER
VER -->|"4. Append to log"| OL
OT -->|"5. Transformed op"| C1
OT -->|"6. Transformed op"| C2
classDef client fill:#4c3a5e,stroke:#818cf8,color:#e2e8f0
classDef service fill:#1a3a2a,stroke:#4ade80,color:#e2e8f0
classDef data fill:#3b3520,stroke:#fbbf24,color:#e2e8f0
Why OT over CRDT for this use case?
| Factor | OT (Server-mediated) | CRDT (Decentralized) |
|---|---|---|
| Server involvement | Required (single serialization point) | Optional (peers can sync directly) |
| Version history | Trivial (linear version chain) | Complex (DAG of causal versions) |
| Undo/redo | Simple (reverse the operation) | Hard (tombstones, causal dependencies) |
| Memory overhead | Low (operations are tiny) | Higher (tombstones never deleted) |
| Offline support | Harder (need server to transform) | Native (merge on reconnect) |
| Correctness proof | Well-understood (25+ years) | Newer (still evolving for rich text) |
For a Google Docs-like product with always-online users, centralized permissions, and strong version history needs, OT wins on simplicity. If we needed peer-to-peer or offline-first (like a note-taking app), CRDT would be the choice.
2) With 10M people editing, how does one keystroke reach only the others in that document?
Problem: FR1 put every editorβs WebSocket on one Collaboration Service, which is what made step 4βs βecho to the other editorsβ a local operation. Scaling out breaks exactly that.
Bad: the single server from FR1. A WebSocket instance tops out around 100K connections β file descriptors, plus a few tens of KB of read and write buffer each β so 5M concurrent editors need at least 50 instances before anything else is considered. One instance is also a total outage surface: every deploy disconnects every editor at once.
The interesting part is what breaks second. Add the 50th instance and FR1βs broadcast quietly stops being correct: Aliceβs server holds her socket and has no way to reach a co-editor whose socket lives on a different host. Users in the same document silently stop seeing each other, which is a correctness failure introduced by the act of scaling, not a slowdown.
Good: Multiple WebSocket Gateway instances behind a load balancer with sticky sessions. Capacity is solved and the routing question is now explicit rather than hidden: when the Collaboration Service produces an operation, how does it find the instance holding a given editorβs connection? Sticky sessions do not answer that, they only keep one client pinned to one instance.
Great: WebSocket Gateway fleet + Redis Pub/Sub for last-mile routing + connection registry.
flowchart LR
CS["Collaboration Service"]:::service
RPS["Redis Pub/Sub"]:::async
WSG1["WS Gateway 1"]:::edge
WSG2["WS Gateway 2"]:::edge
WSG3["WS Gateway 3"]:::edge
U1["User 1"]:::client
U2["User 2"]:::client
U3["User 3"]:::client
CS -->|"1. Broadcast op via Pub/Sub"| RPS
RPS -->|"2. Push to gateway 1"| WSG1
RPS -->|"3. Push to gateway 2"| WSG2
RPS -->|"4. Push to gateway 3"| WSG3
WSG1 -->|"5. Deliver to user 1"| U1
WSG2 -->|"6. Deliver to user 2"| U2
WSG3 -->|"7. Deliver to user 3"| U3
classDef client fill:#4c3a5e,stroke:#818cf8,color:#e2e8f0
classDef edge fill:#1e3a5f,stroke:#60a5fa,color:#e2e8f0
classDef service fill:#1a3a2a,stroke:#4ade80,color:#e2e8f0
classDef async fill:#3b1f5e,stroke:#c084fc,color:#e2e8f0
Mechanism:
- When a user connects via WebSocket, the Gateway registers
(docId, userId) β gatewayInstanceIdin Redis - Gateway subscribes to Redis Pub/Sub channel
doc:{docId}:ops - When Collaboration Service produces a transformed op, it publishes to
doc:{docId}:opschannel - Only gateways with users in that document are subscribed β they receive the op and push to connected users
- Each gateway instance handles ~100K connections. 100 instances = 10M concurrent users
Connection handling:
- Heartbeat: Client sends ping every 30 seconds. If server receives no ping for 60 seconds, connection is considered dead - clean up session.
- Reconnection: Client stores last received version. On reconnect, sends
{resumeFrom: lastVersion}. Server replays any missed operations from the op log. - Graceful shutdown: When a gateway instance is being drained (deploy), it sends a
REDIRECTmessage to all connected clients with a new gateway URL. Clients reconnect within 5 seconds - zero downtime deploys.
Scaling math: 5M concurrent editors across ~2M open documents averages ~2.5 editors per document, so each document channel has a handful of subscribers. Total fan-out is 1M ops/sec Γ ~2.5 recipients = ~2.5M Pub/Sub deliveries/sec. A Redis node handles roughly 1M messages/sec, so this needs about 4 shards partitioned by docId, and you would run 8 for headroom and failure domains rather than 4. On the connection side, ~100K connections per gateway instance means ~50 instances for 5M editors, run at ~70 for headroom.
Note the asymmetry worth remembering: connections are the expensive resource here, not messages. The message fabric is a handful of Redis nodes; the connection fabric is dozens of machines whose only job is holding sockets open.
3) A document has 100K edits. How do we show last Tuesday without replaying all of them?
Problem: FR3 stores a full copy of the document at each save point, and FR2 rewrites the whole document row on every edit. Both store states where the thing that actually happened was a change, and that choice is what makes history both expensive and coarse.
Bad: the full copy per version that FR3 built. Three costs, and the storage one is only the most obvious.
Storage. A 1MB document saved 1,000 times over six months is ~1GB for a single document, and the overwhelming majority of those bytes are identical across consecutive versions. Multiply by 500M documents and the version history dwarfs the 10TB of current state by orders of magnitude.
Resolution. History exists only where a save landed. βWhat did this say at 10:32β is unanswerable, and per-keystroke undo across a collaborative session cannot be expressed at all, because the intermediate states were never recorded.
Write cost. FR2βs per-edit rewrite is the same mistake in the live path: ~20 GB/sec of writes across the fleet at 1M ops/sec Γ 20KB, to express a few MB/sec of real change.
Good: Invert it β store every operation, never the document, and replay from the beginning to reconstruct any version. Storage collapses, because an op is tens of bytes rather than a megabyte, and resolution becomes perfect: every keystroke is a point you can travel to. The new cost is read latency. Opening a 100K-operation document means applying 100K operations before the first character renders, and that is seconds of CPU on the critical path of every document open β 50K of which happen per second.
Great: Periodic snapshots + operation segments. (Borrowing from event sourcing best practices.)
Mechanism:
- Every 100 operations (or every 5 minutes of activity), create a snapshot: serialize the full document state, store in S3 with version tag
- To reconstruct any version V: find the nearest snapshot before V, replay only the operations between that snapshot and V
- Worst case: replay 99 operations (not 100K)
- Snapshots are immutable - old ones are never deleted (they serve as checkpoints in version history)
Undo implementation:
- Local undo (before ACK): Client reverses its own pending operation locally. Trivial.
- Collaborative undo (after ACK): Cannot simply reverse the operation because other usersβ ops may have been applied since. Instead: generate an inverse operation and send it as a new operation through OT. Example: undo of
insert("hello", pos 5)=delete(5 chars at pos 5). This inverse op goes through the full OT pipeline, getting transformed against any concurrent operations. - Undo stack per user: Each user has their own undo stack. Undoing User Aβs last edit doesnβt affect User Bβs edits.
Storage efficiency:
- Operations are tiny: average 20-50 bytes each (type + position + 1-5 chars)
- 100K operations for a 6-month document β 5MB in Cassandra
- Snapshots: one every 100 ops Γ 100K ops = 1000 snapshots Γ 1MB each = 1GB in S3
- Total per document: ~1GB for a heavily edited doc over 6 months. With S3 tiering (Glacier after 30 days), cost is negligible.
4) Two million documents are open at once. Where does each oneβs edit state live?
Problem: Deep Dive 1 gave the Collaboration Service real per-document state β a version counter and a buffer of recent operations to transform against. Deep Dive 2 spread the connections across a fleet but said nothing about where that state lives, and OT only works if every operation for a document is serialized against one authoritative copy of it.
Bad: the single Collaboration Service from FR1, still holding all of it in process memory. Put numbers on the memory: 2M open documents at ~20KB of text each is ~40GB of document state, plus a transform buffer per document, plus the connection state Deep Dive 2 already showed does not fit. That is before any headroom, on one machine, in one process.
Capacity is again not the worst of it. This process is the single serialization point for every document in the system, so it is also a single failure domain for every document in the system: it crashes, and 5M editors lose not just their connection but the authoritative version counter their pending edits were based on. And because the state is in-process, no other machine can take over β the state died with it.
Good: Shard by docId using consistent hashing across N Collaboration Service instances. Each instance owns a subset of documents. Works until an instance crashes - those documents become unavailable.
Great: Consistent hash ring with virtual nodes + stateless OT workers backed by Redis for document state.
flowchart LR
WSG["WebSocket Gateway"]:::edge
LB["Load Balancer"]:::edge
CS1["Collab Worker 1"]:::service
CS2["Collab Worker 2"]:::service
CS3["Collab Worker 3"]:::service
RED["Redis Doc State"]:::data
OL["Operation Log"]:::data
WSG -->|"1. Receive ops"| LB
LB -->|"2. DocId hash"| CS1
LB -->|"3. DocId hash"| CS2
LB -->|"4. DocId hash"| CS3
CS1 -->|"5. Load doc state"| RED
CS2 -->|"6. Load doc state"| RED
CS3 -->|"7. Load doc state"| RED
CS1 -->|"8. Persist op"| OL
classDef edge fill:#1e3a5f,stroke:#60a5fa,color:#e2e8f0
classDef service fill:#1a3a2a,stroke:#4ade80,color:#e2e8f0
classDef data fill:#3b3520,stroke:#fbbf24,color:#e2e8f0
Mechanism:
- Each document is assigned to a Collaboration Worker via consistent hashing on docId
- The workerβs OT state for each document (current version, last N ops for transform buffer) lives in Redis, not in-process memory
- When a worker receives an operation: reads doc state from Redis, transforms, writes new version + op atomically (Redis transaction / Lua script), publishes to Pub/Sub
- If a worker crashes: any other worker can pick up the document because state is in Redis, not local memory
- Rebalancing: add a new worker node to the hash ring β it takes ownership of some docIds β reads their state from Redis β starts processing immediately
Why Redis for OT state and not just Cassandra?
The OT transform operation needs atomic read-modify-write: βread current version, transform against recent ops, increment version, write new op.β This must be serialized per document. Redis single-threaded execution + Lua scripts give us this atomicity. Cassandra doesnβt support atomic read-then-write.
Handling hot documents (50+ concurrent editors):
- A viral document with 50 editors generates 50 ops/sec. Each op requires transform against the last ~10 operations = manageable.
- If a single document exceeds 200 ops/sec (unlikely in text editing), the worker can batch operations in 50ms windows and transform them together before broadcasting.
5) Fifty moving cursors is 50,000 messages a second. How do we avoid sending them all?
Problem: Seeing where your collaborators are is most of what makes a document feel shared. It is also the highest-frequency event in the system: a cursor moves on every keystroke, every arrow key and every mouse click, which is several times more often than the document actually changes.
Bad: Reuse FR1βs broadcast for cursors β every cursor event echoed to every other editor the moment it arrives. An actively working user generates on the order of 20 cursor events/sec between typing, arrow keys and clicking. On a 50-editor document that is 1,000 events/sec, each fanned out to 49 other people: ~50,000 messages/sec for one document.
Almost all of it is waste. Nobody can perceive a cursor position that is superseded 50ms later, so we are spending the entire budget transmitting states that are never rendered, on the same connections the actual edits need.
Good: Throttle each client to one cursor update every 200ms, so 5/sec instead of 20/sec. That is 50 Γ 5 Γ 49 β 12,500 messages/sec per hot document β a 4x cut for no perceptible loss, since 200ms is around the threshold where cursor motion still reads as smooth. Better, and still 12,500 individual messages to say something that could be said once.
Great: Throttled cursor updates + position transformation + server-side aggregation.
Mechanism:
- Client throttles cursor position sends to every 200ms (5 updates/sec per user max)
- Cursor positions are sent as document positions (character offset), not screen coordinates
- When other usersβ operations shift text, the server transforms all active cursor positions using the same OT rules (cursor at pos 10 shifts to pos 11 after an insert at pos 5)
- Server batches all cursor positions for a document every 200ms and sends one aggregated presence frame to each editor:
{cursors: [{userId: "A", pos: 12, color: "#ff6b6b"}, {userId: "B", pos: 45, selection: [45,60], color: "#4ecdc4"}]} - Presence data stored in Redis with 10-second TTL. If no update in 10 seconds, cursor disappears (user left or went idle)
What the aggregation buys. The frame count no longer depends on how many people are moving, only on how many people are watching: 5 frames/sec Γ 50 editors = 250 frames/sec for that document, against 50,000 in the Bad case and 12,500 with throttling alone. That is a 200x reduction from where we started, and the frame carries strictly more information than any individual message did, because it describes the whole room at one instant rather than one cursor at one instant.
Color assignment: Each user gets a deterministic color based on a hash of their userId from a predefined palette of 12 high-contrast colors. Consistent across sessions.
13. Design Self-Audit
| Question | Answer |
|---|---|
| Dedicated search index? | Yes - Elasticsearch for full-text document search by title and content. Indexed asynchronously via Kafka on document save events |
| Stale reads after writes? | After typing, your own changes appear instantly (optimistic local apply). Other users see changes within 100ms via WebSocket push. Acceptable. |
| Single points of failure? | Redis doc state has replicas with auto-failover. Collaboration workers are stateless (state in Redis). Cassandra op log uses RF=3. |
| Dead-letter / reconciliation? | Failed snapshot jobs go to DLQ and retry. If a clientβs op is rejected (version mismatch), client rebases and retries automatically. Orphaned sessions (no heartbeat > 60s) are cleaned by a sweeper. |
| Data freshness across caches? | Cursor presence TTL = 10s (stale cursors auto-disappear). Document metadata cache TTL = 30s. Op log is source of truth - no cache invalidation needed. |
| Cost at scale? | Redis (OT state + Pub/Sub): 20 shards Γ r6g.large β $4000/month. Cassandra op log (RF=3): 12 nodes β $6000/month. WebSocket Gateways (100 instances): $10K/month. S3 snapshots: $500/month. Total hot-path: ~$21K/month for 10M DAU. |
14. Core Flows
Flow 1: Concurrent Edit with OT Resolution
sequenceDiagram
participant U1 as User 1
participant U2 as User 2
participant WSG as WebSocket Gateway
participant CS as Collaboration Service
participant OL as Operation Log
Note over U1,U2: Both at document version 10
U1->>WSG: op: insert "A" at pos 5 (baseVersion: 10)
U2->>WSG: op: delete pos 3 (baseVersion: 10)
WSG->>CS: User1 op (baseVersion 10)
CS->>CS: No concurrent ops. Assign version 11.
CS->>OL: Append op (version 11)
CS->>WSG: Broadcast to User 2
WSG->>U2: Transformed op: insert "A" at pos 5 (version 11)
WSG->>CS: User2 op (baseVersion 10)
CS->>CS: Transform against version 11. Delete at pos 3 stays pos 3.
Note over CS: Insert at pos 5 does not affect delete at pos 3
CS->>OL: Append transformed op (version 12)
CS->>WSG: Broadcast to User 1
WSG->>U1: Transformed op: delete pos 3 (version 12)
Note over U1,U2: Both converge to same state
Walk-through:
- Both users start at version 10 of the document
- User 1 inserts βAβ at position 5 - arrives at server first, gets version 11 directly (no transformation needed since itβs based on the current version)
- User 2βs delete at position 3 arrives based on version 10, but server is now at version 11. Server transforms it against version 11βs operation (insert at pos 5). Since the delete is at pos 3 and the insert was at pos 5, the delete position is unaffected - stays at pos 3
- Both users converge: the document has βAβ inserted at pos 5 AND character at pos 3 deleted
Non-obvious failure path: What if the Collaboration Service crashes mid-transform? The operation was never assigned a version number, so the client will timeout and retry. Operations are idempotent (same content + same base version = same transform). The client retries with the same baseVersion, and the server re-processes it.
Flow 2: Document Open with Snapshot Reconstruction
sequenceDiagram
participant User
participant GW as API Gateway
participant DS as Document Service
participant PG as Postgres
participant S3 as Object Store
participant CAS as Cassandra Op Log
participant WSG as WebSocket Gateway
User->>GW: GET /documents/{docId}
GW->>DS: Fetch document
DS->>PG: Get metadata + latest snapshot version
PG-->>DS: snapshotVersion: 950 currentVersion: 987
DS->>S3: Fetch snapshot at version 950
S3-->>DS: Full document content at v950
DS->>CAS: Read operations 951 to 987
CAS-->>DS: 37 operations
DS->>DS: Apply 37 ops to snapshot
DS-->>User: Document content + version 987
User->>WSG: Open WebSocket for docId
WSG->>WSG: Register session
WSG-->>User: Connected. Listening for ops from v988+
Non-obvious failure path: What if a snapshot is corrupted in S3? The Document Service falls back to the previous snapshot (version 850) and replays 100 more operations (851-987). Slower, but guarantees correctness. If ALL snapshots are lost, the service can rebuild from operation 0 - the op log is the source of truth.
Document Session State Machine
stateDiagram-v2
[*] --> LOADING : User opens document
LOADING --> CONNECTED : Snapshot loaded + WebSocket open
CONNECTED --> EDITING : User types
EDITING --> CONNECTED : Idle timeout 30s
CONNECTED --> RECONNECTING : WebSocket drops
RECONNECTING --> CONNECTED : Reconnected within 30s
RECONNECTING --> OFFLINE : Disconnected over 30s
OFFLINE --> SYNCING : Connection restored
SYNCING --> CONNECTED : All buffered ops sent and ACKed
CONNECTED --> CLOSED : User closes document
CLOSED --> [*]
Each transition: EDITING buffers operations locally and sends over WebSocket. RECONNECTING queues new operations client-side. SYNCING replays all buffered operations to the server with proper base versions for OT resolution.
15. Final Architecture
flowchart TD
U1["User Browsers"]:::client
LB["Load Balancer"]:::edge
GW["API Gateway"]:::edge
WSG["WebSocket Gateway Fleet"]:::edge
DS["Document Service"]:::service
CS["Collaboration Service OT"]:::service
HS["History Service"]:::service
SS["Snapshot Service"]:::service
SUMM["Change Summarizer"]:::service
PRES["Presence Aggregator"]:::service
KF["Kafka"]:::async
RPS["Redis Pub/Sub"]:::async
PG["Postgres Metadata"]:::data
CAS["Cassandra Op Log"]:::data
RED["Redis Doc State + Presence"]:::data
S3["S3 Snapshots"]:::data
ES["Elasticsearch"]:::data
U1 -->|"Open document"| LB
LB -->|"Route HTTP"| GW
LB -->|"Route WebSocket"| WSG
GW -->|"Forward to doc svc"| DS
WSG -->|"Forward to collab svc"| CS
CS -->|"Load doc state"| RED
CS -->|"Append to op log"| CAS
CS -->|"Broadcast via Pub/Sub"| RPS
RPS -->|"Push to WS gateways"| WSG
DS -->|"Fetch doc metadata"| PG
DS -->|"Store doc snapshot"| S3
DS -->|"Read op log"| CAS
SS -->|"Compact op log"| CAS
SS -->|"Save compacted snapshot"| S3
KF -->|"Feed doc summarizer"| SUMM
KF -->|"Update search index"| ES
CS -->|"Publish doc change"| KF
HS -->|"Read version history"| CAS
HS -->|"Fetch old snapshots"| S3
PRES -->|"Track active cursors"| RED
PRES -->|"Broadcast presence"| RPS
classDef client fill:#4c3a5e,stroke:#818cf8,color:#e2e8f0
classDef edge fill:#1e3a5f,stroke:#60a5fa,color:#e2e8f0
classDef service fill:#1a3a2a,stroke:#4ade80,color:#e2e8f0
classDef async fill:#3b1f5e,stroke:#c084fc,color:#e2e8f0
classDef data fill:#3b3520,stroke:#fbbf24,color:#e2e8f0
classDef external fill:#4a1942,stroke:#f472b6,color:#e2e8f0
How it works end-to-end:
- User opens document β browser connects via Load Balancer to WebSocket Gateway for real-time sync
- User types an edit β operation sent over WebSocket to Collaboration Service (OT engine)
- OT transforms and applies β Collaboration Service checks Redis Doc State for current version, transforms against concurrent ops
- Operation persisted β written to Cassandra Op Log as an immutable event
- Broadcast to collaborators β transformed op published via Redis Pub/Sub to all WebSocket Gateway instances holding active editors
- Periodic snapshot saved β Snapshot Service compacts the op log into a full document snapshot stored in S3
- Async indexing and summarization β Kafka carries events to Change Summarizer (version history labels) and Elasticsearch (full-text search)
- Presence tracked β Presence Aggregator maintains cursor positions in Redis, broadcast via Pub/Sub for live cursor indicators
Want a deep dive on rich text OT (formatting operations), offline editing with CRDT fallback, or access control for shared documents? Drop a comment below π
Key Technologies
| Term | What it is |
|---|---|
| Operational Transformation (OT) | Algorithm that adjusts concurrent edit positions so multiple usersβ changes merge without conflicts or data loss. |
| CRDT | Conflict-free Replicated Data Type - a data structure that converges to the same state across replicas without coordination; used in peer-to-peer/offline-first editors. |
| WebSocket | Persistent bidirectional connection between client and server enabling sub-100ms operation push to all collaborators. |
| Kafka | Event bus used for async document lifecycle events (snapshots, version summaries, change notifications). |
| Redis Pub/Sub | In-memory publish/subscribe messaging used to route transformed operations to the correct WebSocket Gateway instance holding each userβs connection. |
| Postgres | Relational DB storing document metadata, permissions, and snapshot version pointers with ACID guarantees. |
| Version Vector | Data structure tracking the latest version each client has seen, enabling the server to identify which operations need transformation on arrival. |
Whatβs Expected at Each Level
This section helps you calibrate your depth. You donβt need to cover everything - just know whatβs expected for your level.
Mid-level
Understand the core problem - multiple users editing the same document simultaneously. Propose a server that broadcasts changes to all connected clients via WebSocket. With prompting, recognize that conflicts arise when edits overlap and that naive βlast writer winsβ destroys data.
Senior
Explain OT (Operational Transformation) or CRDT for conflict resolution - articulate how concurrent operations are transformed to maintain consistency. Propose WebSocket for real-time sync. Discuss cursor/presence indicators, document versioning, and how to handle offline edits that merge on reconnect using buffered operations.
Staff+
Compare OT vs CRDT trade-offs at scale (OT needs a central server for linear ordering, CRDT is peer-to-peer but has larger payloads and tombstone overhead). Discuss undo/redo in collaborative context (transforming undo against concurrent operations), document permissions model with real-time access revocation, and how Google scales this to millions of concurrent docs with strong consistency per document using per-document session processes.
π― Key Takeaways
- Operational Transform (OT) resolves concurrent edits by transforming positions
- Server-mediated OT gives linear version history - simpler than CRDT for online editing
- Snapshot + operation log avoids writing full document on every keystroke
- Redis Pub/Sub routes operations to the correct WebSocket gateway
Related Designs
- Chat System (WhatsApp) - similar WebSocket fan-out, presence tracking
- Notification System - multi-channel push delivery
- Stock Broker (Robinhood) - event sourcing, ordered operations
Related Concepts
Understand the building blocks used in this design:
- WebSockets vs SSE β β carry real-time collaborative edits and presence between clients
- Consistency Models β β how concurrent edits eventually converge via OT or CRDT
- Event Sourcing & CQRS β β the ordered stream of edit operations is the documentβs source of truth
Discussion
Newest first