Collaborative Editor (Google Docs)
Worked solution

Collaborative Editor (Google Docs) — a worked solution

OT vs CRDT. Causal ordering. Real conflict-freedom. Server-authoritative single-writer doc-actor with WAL-before-ack — per-doc serialization, persisted before broadcast, pinned to a home region.

Try it yourself first.

You will remember almost none of this if you read it cold. The workspace walks the same 10 stages and runs the architecture you draw through a simulator, so you find out where your version breaks before you see ours.

Open the Collaborative Editor (Google Docs) workspace

The problem

Build the core of a real-time collaborative editor along the lines of Google Docs / Office Online — the real-time editing surface, deliberately without the peripheral product features (comments, suggestions, @mention notifications) that sit on top of it. Many users open the same doc simultaneously, type in shared paragraphs, see each other's cursors live, and pick up exactly where they left off after closing the laptop for a week. The hard problems are not "how do I store text" — they are causal ordering of concurrent ops, real conflict-freedom under partial failure, sub-200ms keystroke→broadcast at planet scale, and durability of edits the user already saw on their other devices. The load-bearing claim of this design: server-authoritative single-writer doc-actor, with WAL-before-ack to a per-doc Kafka partition fenced by a Raft-issued lease epoch. Everything else is variation on that spine.

The reference architecture

Reference architecture for Collaborative Editor (Google Docs): 14 components — Author Client (Browser / Desktop), Peer Client (concurrent editor), Snapshot CDN, Anycast L4 LB + WAF, Collab WS Gateway, Control API Gateway, Identity + ACL, Doc-Actor Lease Coordinator, Doc-Actor Pool (single-writer per doc), Op Log (canonical, source of truth), Snapshot + Version Store, Doc Metadata DB, Presence + Idempotency Cache, DLP / Audit Worker — connected by 27 flows.Author Client (Browse…Tiptap/ProseMirror + Y-…Peer Client (concurre…Tiptap/ProseMirror + Y-…Snapshot CDNCloudflare / Fastly (si…Anycast L4 LB + WAFKatran (BPF) + Envoy WA…Collab WS GatewayGo + nhooyr.io/websocke…Control API GatewayEnvoy + WAF + per-token…Identity + ACLCustom OIDC/SAML + per-…Doc-Actor Lease Coord…etcd v3 cluster (Raft) …Doc-Actor Pool (singl…Rust + Tokio, embedded …Op Log (canonical, so…Kafka 3.7 KRaft + KIP-4…Snapshot + Version St…S3 versioned, KMS-encry…Doc Metadata DBPostgres 16 + Patroni +…Presence + Idempotenc…Redis 7 Cluster (Stream…DLP / Audit WorkerFlink + per-tenant poli…
14 components, 27 flows. A dashed line is an asynchronous hop. This is the reference design, not the only one that works.

What each component is for

Author Client (Browser / Desktop)Tiptap/ProseMirror + Y-CRDT mirror, IndexedDB TransactionQueue, app-layer ping 10s

Holds the editing surface (Tiptap/ProseMirror), a local Y-CRDT mirror of the doc state, and an IndexedDB-backed transaction queue. Applies every keystroke optimistically to the local mirror, then ships the op over a single long-lived WebSocket. Acks server-blessed ops back into the local mirror; presents peers' transformed ops as they arrive on the broadcast channel.

Why it exists. The user types faster than any network round-trip. Without optimistic local apply, every keystroke would feel like 100ms+ of lag — the editor would be unusable. The IndexedDB queue exists so a tab close, a laptop sleep, or a transient network drop doesn't lose work that the user already saw on screen.

When it fails. If the client mirror diverges from server state (transform bug, custom build) the user sees their text un-write itself when the server pushes the canonical version — extremely high distrust signal. Mitigation: client checksums its mirror against snapshot every N ops and forces a re-snapshot on mismatch.

Peer Client (concurrent editor)Tiptap/ProseMirror + Y-CRDT mirror (same code as Author Client; different user, same docId)

Identical client code to client, but in this trace it's the *receiver* of the broadcast. Holds an open WebSocket terminated at ws-gw; the doc-server pushes every transformed op down that socket so User B sees User A's typing with sub-200ms latency. When User B types, they send back through the same path (peer-client → edge-lb → ws-gw → doc-server) and consistent-hash-on-docId routes them to the SAME actor pod that's already linearizing User A's stream.

Why it exists. The single architectural moment that distinguishes a collaborative editor from a single-user one. Modeled as a first-class node — not folded inside doc-server — so the broadcast leg becomes a visible graph hop and the *symmetry* of bidirectional editing is obvious on the diagram. Without surfacing peer-client, the diagram would suggest 'one user types into a server' and miss the entire point of the problem.

When it fails. If peer-client's WS drops mid-broadcast, the user sees a 'reconnecting…' banner. On reconnect, ws-gw consults doc-coord for current owner, doc-server resumes the stream from last_acked_seq. Worst case: User B's IndexedDB-queued ops conflict with User A's intervening ops; Yrs transforms them on the server side. Real failure is the *appearance* of inconsistency to User B — mitigated by the Figma-style 'reconnect = re-snapshot + replay' flow.

Snapshot CDNCloudflare / Fastly (signed-URL gated)

Caches binary-encoded ydoc snapshot blobs at the edge POP. Client opens a doc → first request is GET on a signed URL; if the edge has the snapshot for that (doc_id, snapshot_seq) tuple it returns immediately, otherwise it falls through to the snap-store origin. TTL 30s.

Why it exists. Doc-open is the noisiest workload — every page-load hits this path before any WS upgrade. Without an edge cache, every open paid an S3 round-trip + a TLS handshake to ws-gw — sub-second open is impossible at scale. The CDN absorbs ~95% of warm opens and lets the WS tier focus on real-time editing.

When it fails. On provider-wide CDN outage (Cloudflare 2022 incident class) the SDK falls back to api-gw → snap-store directly. Degraded UX (slower opens) but no data path failure. AZ-level POP outages are absorbed by anycast and invisible.

Anycast L4 LB + WAFKatran (BPF) + Envoy WAF + per-IP token bucket

Terminates TCP from the public internet, applies WAF rules and per-IP token-bucket rate limiting, then forwards to ws-gw (for WS upgrades) or api-gw (for HTTPS control plane). L4 only — does not terminate TLS, does not parse HTTP application semantics.

Why it exists. An anycast LB is the only practical way to deliver low-latency ingress globally without doing GeoDNS gymnastics. The L4-only design is deliberate: putting TLS at ws-gw lets us keep TLS state on the box that owns the connection, instead of duplicating it across the LB tier (which would 4× the OpenSSL session-cache memory pressure). WAF + rate limit is here, not at ws-gw, because we want to drop bad traffic *before* it consumes a TLS handshake.

When it fails. AZ-level edge-lb loss: anycast routes to the next-closest region transparently, ~50ms latency hit. Region-wide loss: anycast withdraws the regional prefix; clients re-resolve and connect to next region (BGP re-convergence ~30s). Worst case: BPF program regression takes the whole fleet — staged rollout + auto-rollback on health checks mitigates.

Collab WS GatewayGo + nhooyr.io/websocket, BoringSSL, consistent-hash (Maglev) by docId

Terminates TLS, accepts WS upgrades, looks up the doc's current owner pod via the lease cache (Redis 5s TTL, fall-through to doc-coord on miss), and forwards op frames bidirectionally between client and the owning doc-server pod over a long-lived gRPC stream. On the egress side, fans out broadcast frames received from doc-server to every subscribed WS for that doc — including peer-client. Emits reconnect-after hints during drain.

Why it exists. Two reasons separated from doc-server: (1) we need a tier that holds millions of long-lived WS connections — that's a different scaling axis than holding doc state. (2) Routing — every WS for the same doc *must* land on the same actor pod, and the pod identity changes on doc-server failover. ws-gw is the indirection layer that makes the (doc, owner) mapping changeable without dropping client connections.

When it fails. Single pod loss kicks 100K WS connections — clients reconnect with full-jitter exponential backoff. On a fleet roll, drain mode sends reconnect-after: jitter(60-600s) and holds existing conns 120s before SIGTERM; this is the load-bearing reconnect-storm mitigation. Without drain mode, a fleet roll = thundering herd against doc-coord and a multi-minute brownout.

Control API GatewayEnvoy + WAF + per-token bucket

Routes HTTPS requests for: doc CRUD (create/list/delete), permissions (share/unshare), version-history fetch, and minting presigned PUT URLs for the paste-import path — the client uploads the blob straight to snap-store, never through this tier. Validates capability tokens, applies per-token rate limits, fans out to auth, doc-meta, and snap-store as appropriate.

Why it exists. Separated from ws-gw deliberately: a bad deploy of the control plane (e.g. a permission API regression) must not be able to kick every active editor's WS. By splitting the planes we get independent rollout cadence, independent circuit breakers, and independent capacity sizing. Also: the control plane has different latency budgets — 300ms is fine for 'create doc'; 200ms p99 is the hard ceiling for 'broadcast keystroke'.

When it fails. All control-plane operations halt during full outage — users can't create docs, can't change perms, can't fetch version history. Editing existing docs continues unaffected (different gateway). Region-loss → traffic shifts to next region; brief client-side errors during the failover window.

Identity + ACLCustom OIDC/SAML + per-doc capability tokens (signed JWT, exp ≤ 30s)

Validates user identity via OIDC/SAML against the workspace's IdP. Mints per-doc capability tokens (signed Ed25519, embed (user_id, doc_id, scope, exp), TTL ≤30s). On ACL changes in doc-meta, increments the doc's acl_version and pushes access_revoked events to active sessions so they re-mint or get cleanly disconnected.

Why it exists. Real ACL enforcement requires that mid-session revocation actually take effect. A 24h JWT means a fired employee can keep editing for 24 hours after revocation — unacceptable for enterprise. ≤30s TTL bounds the revoke window to one mint cycle. The trade-off is mint QPS — every concurrent session holds a live token, viewers included, not just active typists: 3.75M sessions × 1 mint / 30s ≈ ~125K mints/s steady, and that number is what sizes this fleet.

When it fails. Auth tier outage: nobody can mint new tokens. Existing sessions continue editing for up to 30s, then their tokens expire and they're disconnected. The fix is auto-failover to a secondary auth region; capability tokens issued there validate against the same key (per-region keys would be too operationally hairy).

Doc-Actor Lease Coordinatoretcd v3 cluster (Raft) — 5-node, 2 us-east + 2 us-west + 1 eu (any single-region loss keeps quorum)

Stores three pieces of state under linearizable Raft: (1) (doc_id) → (owner_pod, lease_epoch, expires_at) — the authoritative ownership map. (2) (workspace) → home_region — owns regional rebind during DR. (3) Lease lifecycle (acquire / renew / revoke). Doc-servers acquire leases when first opening a doc, renew every 5s with batched RPCs, and surrender on graceful shutdown.

Why it exists. Without a strongly-consistent ownership map, two doc-server pods could simultaneously believe they own the same doc — and simultaneously append ops to the oplog with conflicting causal histories. The lease + epoch is the *only* reason single-writer holds under partial failure. It also owns home-region binding so DR is a strongly-consistent transition, not a runbook step.

When it fails. Leader loss → Raft re-election in 5s; lease state survives. Two-region partition — the side without quorum stops issuing new leases; existing leases continue until expiry. Loss of >2 nodes simultaneously → the cluster halts (better than split-brain). Etcd write-throughput ceiling at ~10K writes/s; if we ever exceed that we shard etcd by hash(doc_id) % N (designed but not built — current load is ~6 RPCs/s).

Doc-Actor Pool (single-writer per doc)Rust + Tokio, embedded Y-CRDT (Yrs), per-doc serialization queue, mimalloc

Materializes a doc on first WS connection (loads the latest snapshot from S3, replays the oplog tail to current head). Holds in-memory Y-CRDT state. For every incoming op: transforms against the in-memory state, BLOCKS on Kafka acks=all + min.insync.replicas=2 append fenced by the lease epoch, THEN acks the client AND broadcasts the transformed op to all peer WS subscribers. Snapshots to S3 every 1000 ops or 60s. Holds presence pub-sub for the doc. Self-evicts on sustained Kafka append failure to prevent ghost ownership.

Why it exists. This is the *spine* of the design — the load-bearing claim that makes 'real conflict-freedom' real. By keeping exactly one writer per doc (enforced by doc-coord lease + Kafka epoch fencing), there is no transform-divergence surface, no multi-master merge surface, no cross-replica reconciliation surface. OT and CRDT structures are used INSIDE the actor as data-structure choices, but they're never asked to do their hardest job (multi-master conflict resolution) because the architecture eliminates the multi-master scenario.

When it fails. Pod crash → up to ~80K open docs lose their actor briefly. doc-coord re-promotes them to surviving pods (RTO 5s); actors materialize from snapshot + oplog tail. Critically: ops the user *saw acked on another client* are durable in Kafka — they replay on the new actor. RPO = 0 on the durable boundary. The dangerous failure is partition-from-Kafka-but-not-doc-coord: actor self-evicts after 5s of failed appends to prevent piling up acks-it-can't-deliver.

Op Log (canonical, source of truth)Kafka 3.7 KRaft + KIP-405 Tiered Storage (S3 cold tier)

Receives one append per op from the doc-actor, with strict per-partition ordering and acks=all + min.insync.replicas=2 durability. A transactional producer keyed by a per-doc transactional.id prevents stale appends from a partitioned former owner: the successor's initTransactions() bumps the broker producer epoch and fences the predecessor (ProducerFenced). (Idempotence alone only dedups retries within one producer session and would not fence a second process.) Consumed by: doc-server (replay on cold open / handoff), dlp-scanner (DLP + audit). 30d hot in-broker; 7y cold via S3 tiered storage.

Why it exists. Two non-negotiables it provides: (1) causal order — a single per-doc partition gives strict total order without needing application-level vector clocks. (2) durability before ack — acks=all + min.insync=2 means an op the user sees acked on another client is in Kafka before they see the ack. Without those two, the design has no durable single-writer story; this is the spine of WAL-before-ack.

When it fails. Broker loss (single): ISR drop, leader re-election in ~6s, producers retry seamlessly. Multi-broker loss bringing ISR below min.insync.replicas=2: producer errors, doc-actor self-evicts after 5s. Region loss: cross-region MirrorMaker 2 lags by ~90s steady, alerted at 120s — region-loss RPO is 120s. Tiered storage S3 outage: hot-tier ops still work; cold replay (older than 24h) errors until S3 recovers.

Snapshot + Version StoreS3 versioned, KMS-encrypted, S3 Cross-Region Replication (CRR) to DR region

Holds binary-encoded Y-CRDT snapshots (every 1000 ops or 60s, written by doc-server). Holds version-history snapshots retained 7y per workspace policy. Holds DLP-scanner audit records under S3 Object Lock (write-once for compliance). Cross-region replicated via S3 CRR.

Why it exists. Bounds cold-open replay cost — without snapshots, opening a 6-month-old doc would replay every op since creation (potentially millions). With 60s snapshots, replay is bounded to ~1000 ops. Versioning gives 'restore previous version' for free. Object Lock gives a tamper-evident audit trail required for SOC 2 / HIPAA.

When it fails. AZ event: S3 native multi-AZ — invisible to us. Region loss: failover to DR region (us-west-2 mirror); cross-region replication lags by ~15min p99, so the last 15min of snapshots may be missing — oplog replay fills the gap. S3 5xx storms: doc-server retries with exp-backoff; if persistent, snapshot pipeline lags, alerted via snapshot_age_p99 > 300s.

Doc Metadata DBPostgres 16 + Patroni + Citus, sharded by workspace_id

Stores doc rows (id, workspace, owner, current_snapshot_seq, current_oplog_seq, home_region, acl_version), ACL grants, DLP quarantine flags. Sharded by workspace_id (32 shards) for tenant locality; sync followers for HA; cross-region async observer for DR + read-local-region.

Why it exists. Some data is fundamentally relational and transactional — ACL grants, quarantine flags, workspace membership — and doesn't fit the per-doc oplog model. Co-locating them in a sharded relational store lets us run real transactions (e.g. 'add a member AND grant role AND bump acl_version' as one atomic update). Sync followers are non-negotiable: losing a recent ACL revoke on failover is a security incident, not a 'lost write'.

When it fails. Primary loss: Patroni auto-failover in 30s. Sync followers ack means RPO=0. Region loss: cross-region observer promotes; up to 30s of recent control-plane writes lost (RPO 30s) — the design accepts this for the rarity of region loss. Vacuum brownouts: monitored via per-shard write p99; autovacuum tuned aggressively given churn on acl_version bumps.

Presence + Idempotency CacheRedis 7 Cluster (Streams + ephemeral keys, AOF every-write)

Hosts three disjoint key prefixes: (1) presence:{doc_id} — XSTREAM of cursor/selection events, TTL 30s, never persisted. (2) idem:{client_op_id} — 7-day SETNX dedup window absorbing client retries through reconnect. (3) lease-cache:{doc_id} — 5s TTL view of doc-coord's owner mapping so ws-gw doesn't hit etcd on every WS upgrade.

Why it exists. Each prefix exists for a different reason. Presence lives here because it's high-volume (12M events/s) and short-lived — it must NOT touch the durable oplog or the keystroke path slows down. Idempotency lives here because client retries through reconnect need a sub-millisecond dedup check; Postgres would melt. Lease cache lives here because hammering etcd with raw lookups would exhaust its write/read budget. Co-located on one cluster because the operational overhead of three separate Redis fleets isn't worth the marginal isolation.

When it fails. Single shard primary loss: failover in 15s, RPO ~1s on idem keys (AOF flush window). Worst case: idem loses a 1s window of dedup state and a client retry creates a duplicate op — Yrs idempotent-by-content-hash de-dups it on the actor. Whole cluster down: presence stops working (cursors freeze); idem cache cold, retries may double-apply briefly; lease-cache cold, ws-gw falls through to etcd directly (etcd absorbs the 100× spike for the duration of the outage).

DLP / Audit WorkerFlink + per-tenant policy bundle + immutable audit log writer (S3 Object Lock)

Separate consumer group on the oplog. Two concerns: (1) DLP — runs per-tenant policy (PII patterns, classification labels, share-out blocks) against op content; on a hit, writes a quarantine_op row to doc-meta which the doc-server enforces by rejecting future ops on the affected range and showing a banner to the editor. (2) Audit — writes immutable audit records to S3 Object Lock, keyed by (workspace_id, ts), retained 7y per SOC 2 / HIPAA policy.

Why it exists. Enterprise customers will not sign without (a) DLP enforcement that catches a leaked SSN before it propagates and (b) tamper-evident audit logs that can prove what was edited and when. Both run async because synchronous DLP on every keystroke would 5× the latency budget. Critically: DLP NEVER mutates the canonical oplog — the audit invariant is 'what was typed is preserved verbatim'; quarantine acts via a sidecar flag that the doc-server reads.

When it fails. Lag SLO 60s, paged at 300s. DLP brownout: quarantine flags arrive late; a sensitive op is visible to peers for the lag window. Audit-write outage: alerted at the per-write 5xx threshold; policy is to hard-fail the consumer (don't ack the offset) rather than skip — auditability cannot have gaps. Recovery: replay from offset, idempotent on (op_id).

Stage by stage

The same 10 stages the workspace walks, answered.

01Clarifications

What would you ask before drawing a single box?

Typical clarifications to surface:

  • Rich text or plain text? Affects op model (ProseMirror steps vs Yjs Y.Text). Assume rich text — paragraphs, formatting, embedded objects.
  • Offline editing required? Yes — laptop-sleep-then-wake within hours; reject pure-server-required designs.
  • Multi-device same user? Yes — laptop + phone editing concurrently is normal.
  • Scope note. This canonical is the real-time editing core. Comments, suggestions (track-changes), and @mention notifications are intentionally out of scope — they layer over the same oplog spine (add an op-kind discriminator; add a consumer group) without changing the write path.
  • Multi-region edits to the same doc? Yes for read; for write, home-region per doc (not active-active per op).
  • Latency target? p99 keystroke→broadcast intra-region < 200ms (the shared-cursor "feels live" bar). Cross-region eats the WAN RTT.
  • Retention? Operational oplog 30d hot, cold tier 7y for compliance. Snapshots versioned forever per workspace policy.
  • Authentication? SSO (OIDC/SAML), per-doc ACL, capability tokens with ≤30s lifetime so revocation actually takes effect mid-session.
  • Max doc size / max editors? ~50MB / ~50 concurrent writers per doc. Above 50 writers → broadcast-fanout mode.

Assumptions to state (Mid tier, 50M DAU):

  • Concurrent WS at peak: 3.75M (7.5% of DAU).
  • Active editors at peak: ~500K (13% of online).
  • Sustained ops/s system-wide: 2M → peak 8M ops/s.
  • Peak deliveries/s (fan-out): 12M/s (mean 1.5 editors/active-doc).
  • 30d hot oplog after RF=3 + LZ4: ~1 PB on disk (≈330 TB compressed pre-replication).
  • Snapshot every 1000 ops or 60s.
  • 5-region active-standby with home-region pinning per doc.

02Functional reqs

What must this system actually do?

  • Open a doc → see current contents (p99 < 300ms warm, < 1.5s cold).
  • Type → see your own keystrokes locally instantly (optimistic), see peers' keystrokes within p99 200ms intra-region.
  • Multi-cursor / multi-selection visible per editor (presence channel).
  • Version history — restore any point in the last 7y (a native side-effect of snapshots + oplog retention; not a separate service).
  • Offline edit → automatic merge on reconnect (Figma-style: re-snapshot + replay buffered).
  • Permission revoke → mid-session WS close within ≤30s (one capability-token TTL).

03Non-functional

What must it promise about speed, uptime and correctness?

  • Availability: 99.95% on the read+open path; 99.9% on the write path (writes can briefly degrade during a doc-actor handoff).
  • Latency: p99 keystroke→broadcast 200ms intra-region, 400ms cross-region. p99 doc-open 300ms warm / 1.5s cold. p99 reconnect-resume 1.5s.
  • Durability: Zero loss of any op the user observed on a second device. RPO=0 on the durable boundary (Kafka acks=all, MIN.ISR=2). Cross-region DR: RPO ≤ 120s on doc content (MirrorMaker-2 oplog replication, alerted at 120s lag); snapshots ride S3 CRR at ~15min p99 with the gap filled by oplog replay; doc-meta ≤ 30s via its async observer.
  • Consistency: Strong order per-doc (single-writer doc-actor); causal across docs is sufficient. Read-your-write within a doc is guaranteed by single-writer; across regions it is region-local.
  • Conflict-freedom: Real, not aspirational. The doc-actor is the linearizer; OT/CRDT structures used inside it are an implementation choice, not a multi-master substitute.
  • Scalability: Horizontal on every tier except the per-doc actor (which is, by design, a single writer). Doc-server pool scales by adding pods; doc-coord scales by sharding the etcd keyspace if a 5-node cluster bottlenecks.
  • Security: Per-doc capability tokens with ≤30s TTL; ACL-version-bump invalidates in-flight caps; mTLS service-to-service; KMS-encrypted snapshots; doc-server runs in workspace-scoped tenancy boundaries.

04Capacity estimation

How much load and data does this have to hold?

Walk-through (Mid tier, 50M DAU; Workspace tier 4×):

  1. Concurrent WS at peak. 50M DAU × 7.5% online + doc-open = 3.75M concurrent WS. Anchored against Slack's published 5M+ peak WS at smaller-than-Workspace scale.
  2. Active editors. 13% of online actively typing → ~500K active editors. The remainder are viewing or idle.
  3. Sustained ops/s. 500K × 4 ops/s/editor (after 250ms client-side coalesce of fast typing) = 2M ops/s sustained. Google's published 60K–200K coarser-grained ops/s for ~100M DAU lands in the same regime.
  4. Peak ops/s. Workday peak multiplier 4× → 8M ops/s peak.
  5. Deliveries/s. Mean 1.5 editors/active-doc → fan-out factor 0.5 above the producer → 12M deliveries/s peak. Discord pushes 26M WS events/s on 400 Erlang hosts — same order of magnitude.
  6. Doc-server pod count. Each pod holds ~250K hot docs (32GB heap, ~110KB live state per doc — ops since last snapshot + presence + WS refs). At 333K active docs (= 500K editors / 1.5 mean), 2 pods could fit it — but the math that matters is handoff blast radius. Six pods per region × 5 regions = 30 pods caps each pod's blast radius at ~80K open docs per crash (3.75M WS / 1.5 per doc = 2.5M materialized actors ÷ 30), restorable within RTO 5s.
  7. ws-gw pod count. 100K active WS / pod (Phoenix bench 2.3M idle, 5× drop active, halved for safety). 3.75M / 100K = 38 pods raw → 12/region × 5 = 60 pods with 1.5× headroom for region failover.
  8. Oplog partition count. 8M ops/s peak / max 4K ops/s/partition = 2K → 2048 partitions, RF=3, distributed over 60 brokers. Each broker handles ~33 partition leaders + followers — well within a c6i.4xlarge's 16-core / 4Gbps NIC budget.
  9. Oplog hot bytes. 2M ops/s × 200B × 86400 × 30d = 1 PB raw → ÷3 LZ4 → ~330 TB raw, × RF=3 = ~1 PB on disk per region. KIP-405 tiered storage offloads anything older than 24h to S3, so on-broker is ~30 TB/region hot.
  10. doc-meta size. 1B docs × 1KB metadata + 100M ACL grants × 200B = ~1 TB. 32 shards × 30GB each. Sized for ACL read/write QPS: ~50K reads/s (capability mint), ~5K writes/s (perm changes).
  11. Snapshot churn. 500K active × 1 snapshot/min × 50KB median = 36 TB/day of S3 PUTs. Request charges dominate storage: ~720M PUTs/day ≈ $108K/month at $0.005 per 1K PUTs, vs ~$23K/month for the ~1 PB steady state under the 30d checkpoint GC. The cost lever is PUT count, not storage class — a median doc at 4 ops/s hits the 60s timer with only ~240 ops, so stretching slow-churn docs toward the 1000-op bound cuts PUTs ~4× with the same replay ceiling. Intelligent-Tiering is a non-answer here (50KB median objects sit below its 128KB tiering floor).
  12. Cache working set. Two distinct sizing axes: (a) hot-doc lookup cache 1B docs × 2% hot × 50KB = ~1 TB would fit in 10 pods, but (b) presence absorbs 12M events/s peak (presence XADD is the dominant cost), so we run 128 pods to keep per-shard write CPU and AOF fsync queue healthy. Both prefixes co-reside.
  13. Capability-token mint QPS. Every concurrent session holds a live capability token — viewers included — so mint load tracks WS count, not typing: 3.75M sessions × 1 mint per 30s TTL ≈ ~125K mints/s steady. auth runs 60 pods → ~2.1K mints/s/pod against an HSM-backed signer's budget (hardware-keyed Ed25519 ≈ 5K signs/s/pod) — headroom for reconnect-storm re-mint spikes.
  14. Lease coordinator load. 250K active docs renewed every 5s with batched renewals (one Lease.KeepAliveOnce RPC per pod carries up to 10K leases) = 30 pods × 1 batched RPC / 5s = 6 RPCs/s total — three orders of magnitude under etcd's documented ~10K writes/s ceiling. Without batching, naïve per-doc renewal would be ~50K writes/s and require an etcd shard, which is why batching is non-negotiable.
  15. Cross-region. Each doc home-pins. ~20% of editors collaborate cross-region → ~100K editors paying WAN RTT (~80–150ms transatlantic), eating into the keystroke→broadcast budget. SLO has a separate cross-region p99 of 400ms.

Three scale-tier flips (the design must be ready for):

  • Single doc-process saturates at ~50 editors / ~500 ops/s for one doc. GC + per-doc serialization queue dominates. Above this → broadcast-fanout mode (writers on actor, viewers on edge-cached snapshot stream).
  • Don't shard a single doc. Above 100 concurrent editors, demote viewers to read-only or section-partition; never split write-authority (re-introduces OT divergence).
  • Presence flips to compacted Kafka above 10K events/s/doc so cursor noise doesn't poison the inline op path.

05API design

What does the outside world call, and what comes back?

POST /api/v1/docs
Content-Type: application/json
Authorization: Bearer <token>

{ "workspaceId": "ws_abc", "title": "Q4 plan", "type": "doc" }

201 Created
{ "docId": "doc_01H...", "homeRegion": "us-east-1", "snapshotUrl": "https://..." }
GET wss://collab.arc.ly/v1/collab/{docId}?lastAckedSeq={seq}
Authorization: Bearer <capability-token, 30s>

# Server upgrades to WS. First server frame:
{ "type": "snapshot", "seq": 12345, "snapshotUrl": "https://...", "epoch": 8 }
{ "type": "ops",      "seq": 12346..N, "ops": [...] }   # delta if lastAckedSeq known

# Client → server (continuous):
{ "type": "op",       "clientOpId": "uuidv7", "parentSeq": N, "op": {...}, "epoch": 8 }

# Server → client (continuous):
{ "type": "op-ack",   "clientOpId": "...", "seq": N+1 }
{ "type": "op",       "seq": M, "op": {...} }            # peer ops
{ "type": "presence", "userId": "...", "cursor": {...} } # ephemeral, never persisted
{ "type": "reconnect-after", "ms": 3000, "reason": "drain" }
{ "type": "access-revoked" }                              # forces clean WS close
GET /api/v1/docs/{docId}/history?cursor=&limit=100&from=&to=
# cursor: opaque, descending snapshotSeq; limit caps at 200; from/to: optional ts range ("versions around date X")
200 { "items": [ { "snapshotSeq": ..., "ts": ..., "by": ... }, ... ], "nextCursor": "..." }   # nextCursor null at end

06Data model

What gets stored, and what is it looked up by?

Op record (stored in Kafka oplog as one record per op, partition key (workspace_id, doc_id)):

fieldtypenotes
seqint64monotonic per-doc (assigned by the doc-actor)
epochint32doc-coord lease epoch — fencing token
client_op_iduuid v7sender-generated; idem dedup
parent_seqint64causal context — ack of last seen server seq
opbytesbinary-encoded Yrs update (transformed)
op_kindenumtext / format / share
author_idbigintfrom capability token, server-trusted
tstimestampserver wall-clock, advisory only

Doc-meta (Postgres, sharded by workspace_id):

tablekeynotes
docs(workspace, doc_id)current_snapshot_seq, current_oplog_seq, home_region, acl_version
doc_acl(doc_id, principal)role: viewer/editor/owner
quarantine(doc_id, op_id)DLP-scanner-written flag; doc-server enforces

Snapshot (S3 versioned):

  • s3://docs-{region}/ws={workspace}/{doc_id}/{snapshot_seq}.bin — binary-encoded Yrs state, zstd, KMS-encrypted.
  • Latest pointer in docs.current_snapshot_seq. Versioned bucket retains every snapshot for 7y per policy.

Why Kafka (KIP-405 tiered) for the canonical oplog and not Postgres? Per-partition strict ordering, a transactional producer keyed by a per-doc transactional.id whose epoch bump fences a stale owner (the WAL-before-ack contract), 30d hot retention with seamless cold tier, and compaction support if we later want a per-key snapshot view. Postgres would need partitioning logic + outbox tables + a separate truncate/cold-store pipeline that we already get for free here.

07High-level design

Which components handle a request, and in what order?

Architecture summary (request-flow, left to right):

  1. Client holds the local Y-CRDT mirror, applies ops optimistically, queues them to IndexedDB, and ships them over a WS to the closest edge.
  2. Edge L4 LB + WAF (Katran/Envoy) anycasts traffic to the nearest region; rate-limits, blocks the obvious nasties, terminates only at the gateway.
  3. Collab WS Gateway terminates TLS, consistent-hashes the doc to a doc-server pod via doc-coord, forwards op frames in either direction, and on drain emits reconnect-after: jitter(60–600s) to prevent storms.
  4. Control API Gateway carries doc CRUD, share/unshare, version-history reads, and mints the presigned URL for the paste-50MB-import path — the blob itself goes client → S3 directly, never through the app tier.
  5. Identity + ACL mints per-doc capability tokens with ≤30s TTL; ACL bump on the doc-meta forces token re-mint and pushes access-revoked on existing sockets.
  6. Doc-Actor Lease Coordinator (etcd Raft) owns (doc_id) → (owner_pod, lease_epoch). It is the only thing that can promote a new owner; the lease epoch becomes the fencing token on every oplog append.
  7. Doc-Actor Pool (Rust + Yrs) holds the in-memory doc, transforms ops via the embedded CRDT engine, appends transformed ops to Kafka with acks=all and the lease epoch, and only then broadcasts to peer subscribers. This is the load-bearing claim of the design.
  8. Kafka Oplog (KRaft + KIP-405) is the source of truth — strict per-partition ordering, transactional producer with per-doc transactional.id epoch fencing, 30d hot + 7y cold tiered.
  9. Snap-store (S3) holds checkpoints every 1000 ops or 60s, versioned, KMS-encrypted; cross-region async replicated for DR.
  10. Doc-meta (Postgres-HA + Citus) holds the durable rows for ACL, version pointer, DLP quarantine flags; sync-replicated so failover RPO=0.
  11. Presence + Idem Cache (Redis Cluster) is split into three key-prefix concerns: ephemeral presence (TTL 30s, never persisted), idempotency dedup window (7d, AOF-fsync), and lease lookup cache (5s).
  12. DLP / Audit worker is the one derived pipeline included in this core canonical: a separate consumer group on the oplog runs per-tenant policy + writes an immutable audit log to S3 Object Lock; the canonical oplog is never mutated. (Comment/mention notif-workers are the same shape — omitted here to keep the diagram focused on the real-time editing story.)

Data flow on keystroke (the contract path) — two users editing the same doc:

  1. User A's ingress: client → edge-lb → ws-gw → doc-server. ws-gw consistent-hashes by docId, so both users land on the same doc-server pod.
  2. Single-writer transform: doc-server transforms User A's op against the in-memory Y-CRDT state.
  3. WAL-before-ack (durability boundary): doc-server BLOCKS on a Kafka acks=all + min.insync.replicas=2 append, fenced by the lease epoch.
  4. Broadcast back: doc-server pushes the transformed op down its bidi gRPC stream to ws-gw, which fans out the WS push to every other subscriber for that doc — including User B (peer-client).
  5. User B sees User A's keystroke within the 200ms intra-region budget. Symmetric: if User B also types, their ops enter via peer-client → edge-lb → ws-gw → doc-server and land on the same actor, where the per-doc serialization queue interleaves A's and B's ops in causal order.

This is the moment of collaboration. The single-writer doc-actor is what makes it conflict-free in the strong sense — A's and B's concurrent ops never race against each other on different machines, never need a multi-master merge, never produce divergent server states. The diagram makes this visible: both client and peer-client fan into the same ws-gw → doc-server, and the broadcast leg (doc-server → ws-gw → peer-client) closes the loop.

Data flow on reconnect: client (last_acked_seq=V) → ws-gw → doc-coord (find owner) → doc-server (load snapshot if cold; replay oplog from V → head; push delta or full reload if V < oplog floor) → resume.

Data flow on doc-actor handoff (failover): detected crash (kubelet death, ws-gw's severed gRPC streams) → immediate doc-coord.LeaseRevoke, handoff in ~5s; silent partition → doc-server-A misses renewals (every 5s) and doc-coord expires the lease at the 15s TTL. Either way: doc-server-B acquires new lease (epoch+1) → calls initTransactions() on the per-doc transactional.id (bumping the broker producer epoch) → loads snapshot from S3 → replays oplog tail → starts accepting ops. doc-server-A's now-stale appends are fenced by Kafka (ProducerFenced/INVALID_PRODUCER_EPOCH); A self-evicts on the rejection.

08Deep dives

Which part breaks first, and what do you do about it?

1. Why server-authoritative single-writer beats both pure OT and pure CRDT here.

Pure OT (Google Wave, Mobwrite legacy): Ops travel in any direction; every node must transform every op against every other concurrent op. Subtly buggy (TP1/TP2 properties) — Wave shipped with documented divergence cases. The fix that took Atlassian and Google a decade was the same one we're starting with: a central server that linearizes ops and broadcasts the transformed result, so clients only ever apply server-blessed ops. We don't even call it OT internally — we use Yrs (Y-CRDT) as our state structure because the math is cleaner — but the authority model is OT-style central server.

Pure CRDT (Yjs P2P, Automerge merge-anywhere): Wins for genuinely offline-first / P2P scenarios (think local-first apps, no cloud). Costs: tombstone bloat (Yjs #321, Automerge #337 — every delete is a permanent record until causal-stable GC), vector-clock per op, and on multi-region active-active write the merge surface is real and surprising for users. We get none of CRDT's advantages — we have a cloud, we have a single linearizer per doc — and would pay all of the costs.

Server-authoritative single-writer (this design): The doc-actor is the linearizer. We use Yrs for the data structure because deletes-as-tombstones-with-GC are mathematically convenient, but only one process at a time mutates a given doc, and the lease+fencing-token machinery makes "only one process" a hard guarantee under partial failure, not a hopeful one. Conflict-freedom is real, not approximated.

2. WAL-before-ack and the durability budget.

The doc-actor blocks on Kafka's acks=all (with min.insync.replicas=2) BEFORE acking the client. This is non-negotiable. Consider the alternative: actor broadcasts to peers, then crashes before persisting. Other clients' UIs already show the op. On recovery, the canonical state is older — the user sees their text un-write itself on the peer's screen. This is the worst possible failure for a collab editor, far worse than a 30ms latency hit. So we eat the latency: typical Kafka acks=all in same AZ is ~5ms, cross-AZ ~15ms, comfortably inside the 200ms keystroke→broadcast budget.

3. The fencing-token flow that makes single-writer real.

doc-coord (etcd) issues (owner_pod, lease_epoch). doc-server-A holds lease epoch=8 for doc-X. Network blip — A misses 3 renewals; lease expires. doc-coord promotes B with epoch=9. Meanwhile, A still has WS connections and queued ops; A tries to append to Kafka. Here the mechanism matters: plain enable.idempotence=true does NOT fence a zombie — each producer session gets its own broker-assigned PID, and idempotence only dedups retries within one session by (PID, partition, seq); two different processes get two different PIDs and both would be allowed to write. The real fencing comes from a transactional producer keyed by a deterministic per-doc transactional.id (e.g. oplog-<doc_id>). When B is promoted it calls initTransactions() for that same transactional.id, which bumps the broker-side producer epoch and fences the previous holder — A's next append fails with ProducerFenced (surfaced as INVALID_PRODUCER_EPOCH). So the Raft lease authorizes B to take over; the shared transactional.id + epoch bump is what makes the broker actually reject A. A sees the rejection, self-evicts, drops its WS connections (clients reconnect to ws-gw, get re-routed to B). No double-apply. No ops lost (A's ops never committed, so they never existed in the canonical store; clients re-send via the IndexedDB queue).

The harder partition (doc-server can reach doc-coord but NOT Kafka). A's lease renewals succeed via e13 (different network path); but e14 (Kafka) is unreachable. Without WAL-before-ack, A would broadcast to peers and silently lose ops on recovery. With WAL-before-ack, A simply cannot ack the client — the user sees saving… and types into a buffer. A's safety mechanism: if e14 fails sustained > 5s, A proactively calls doc-coord.LeaseRevoke and drops all WS connections. Clients reconnect, ws-gw re-routes them to a pod with healthy Kafka connectivity, and the IndexedDB-queued ops replay against the new actor. The client UX during the gap is saving… followed by a brief reconnect — equivalent to a 5s flap, no data loss.

Why the oplog is the source of truth, not doc-meta. Earlier drafts of this design coupled the oplog append (e14) and a doc-meta cursor update (e16) into a transactional pair. We rejected that: under actor crash between the two writes, the cursor lags the oplog head and a senior reviewer caught it as a non-outbox masquerading as an outbox. Final design: current_oplog_seq is updated on snapshot, not per-op, via compare-and-set (UPDATE docs SET current_oplog_seq=$new WHERE doc_id=$1 AND current_oplog_seq < $new). The cursor is a hint that bounds cold-open replay; if it lags, replay starts a few hundred ops earlier and Yrs idempotently re-applies them. The oplog's per-partition ordering + transactional-producer epoch fencing is the only durable boundary — there is no second store to coordinate with.

4. Reconnect = re-snapshot + replay (the Figma move).

Old design: client tries to splice missed ops into local state, applying transforms client-side. Subtly buggy class — every causality-violation bug is here. New design: on any reconnect, client throws away local view, downloads fresh snapshot via signed URL through CDN, re-applies its IndexedDB-queued ops on top. p99 reconnect-resume budget 1.5s — well within human "I just refocused this tab" tolerance. Removes a whole class of bugs.

5. Hot-doc fanout (the all-hands doc problem).

A shared all-hands doc with 5,000 concurrent viewers (one author actually typing) saturates a single doc-actor. Fix: above concurrent_clients_per_doc > 200 OR ops_per_sec_per_doc > 500, the doc enters fanout mode. Writers stay on the actor. Viewers get demoted to a compacted Kafka topic + edge-cached snapshot — they see updates within snapshot cadence (60s) plus oplog tail. This is the Google Docs "view-only mode for >100 collaborators" rule and the Figma broadcast pattern, formalized.

6. Snapshot/oplog retention floor.

Oplog retention must satisfy oplog_min_seq < snapshot_seq AND oplog_min_age > max(reconnect_window, ops_replay_for_audit). We pick 30d hot because (a) reconnect window: laptop sleep up to 7d is realistic, doubled for jitter; (b) audit replays: forensic reads up to 30d. Below the floor, the client gets STALE_CLIENT + force reload. We never truncate ahead of snapshot fsync — the silent-corruption path.

7. Multi-region (home-region per doc).

Each doc is created in a home region (workspace's primary region + locality hint). All writes for that doc serialize there. Cross-region collaborators connect locally for low-latency control + presence, but op frames forward to home; their keystroke→broadcast eats one WAN RTT. Snap-store replicates async cross-region (~15min RPO) so DR + version-history reads are local. Doc-meta has a cross-region async observer per shard — read your workspace's ACL locally without a WAN hop. We are deliberately CP per doc; AP on multi-master CRDT was rejected because (a) the single-writer model already guarantees conflict-freedom, (b) no UX team wants to explain "your colleague's edits merged with yours in a way you didn't expect" to enterprise customers.

09Trade-offs

What did this design cost, and what breaks at 10×?

Trade-offs we accepted (and the alternatives we rejected):

  • WAL-before-ack adds ~10-15ms to every keystroke but buys RPO=0 on the durable boundary. Rejected: ack-before-persist (broadcast-then-Kafka). Failure mode of the rejected variant — clients see ops the server later loses — is unacceptable for an editor.
  • Single-writer doc-actor means a hot doc has a single-pod ceiling. Rejected: shard the doc by paragraph or section. Section-shard reintroduces cross-shard transform / OT divergence — exactly the bug class we're trying to extinguish. Hot-doc fanout (writer/viewer split) handles the real-world hot-doc shape (many viewers, few writers) without breaking authority.
  • Home-region per doc, not active-active per op. Rejected: multi-leader CRDT replication. Costs (tombstone bloat, surprising merges) > benefits (cross-region write latency reduction for the 20% of cross-region collabs).
  • Y-CRDT (Yrs) inside the actor instead of OT. Rejected: ProseMirror-style step OT. Yrs is faster (josephg's 56ms vs 5min Automerge benchmark in the Rust port), tombstone GC is well-understood, and we get a binary update format we can persist directly. We don't get CRDT's P2P advantage — but we don't need it.
  • etcd for the lease coordinator, not ZooKeeper. Cleaner Raft semantics, simpler client lib, mature in production. Spanner/CockroachDB rejected because we need simple per-key leases, not distributed transactions.
  • Capability tokens with 30s TTL (not 24h JWTs). Costs: more frequent token-mint RPS. Buys: revocation actually takes effect mid-session. The auth tier is sized to absorb the mint rate.
  • Peripheral product features excluded from this canonical (comments, suggestions/track-changes, @mention notifications). Costs: real-Docs is bigger. Buys: the diagram teaches the load-bearing story — WAL-before-ack + single-writer + lease-fenced Kafka epoch — without dilution. Each excluded feature layers over the SAME oplog with a new op_kind and a new consumer group; the write path never changes. The dlp-scanner in this canonical is deliberately kept as the reference example of that "add a consumer group" pattern.

Open questions for human review:

  • Is 30d hot oplog sustainable cost-wise at Workspace tier (4× Mid)? KIP-405 tiered should make this fine, but we should validate broker IOPS at peak rebuild scenarios.
  • Should presence flip to a separate Kafka cluster from the start, rather than co-residing in Redis? If presence events spike to 50M/s in pathological viral-doc events, Redis Streams memory pressure is a worry. We've sketched the flip but not load-tested it.
  • The 50-editor-per-doc cliff is theoretical. Real-world load tests need to validate it (and the broadcast-fanout transition logic) on production-class hardware.
  • Layering peripheral product features (comments / suggestions / @mention push notifications) back on top of this canonical is a one-node-per-feature exercise: add a new consumer group on the oplog (same shape as dlp-scanner), route to doc-meta and/or an external push tier, size the consumer lag SLO independently. Deliberately out of scope here to keep the diagram focused on the real-time editing spine.
  • doc-coord scaling beyond 5-node etcd: at ~10M lease ops/s, etcd's write throughput becomes the ceiling. Sketch: shard the etcd keyspace by hash(doc_id) % N across N etcd clusters. Not needed at Mid tier; queue for a 10× re-architecture review.

Primary sources

  • Figma — How Figma's multiplayer technology works (Evan Wallace)
  • Figma — Making multiplayer more reliable
  • Figma — Live migration of multiplayer servers
  • ProseMirror — Collaborative Editing (Marijn Haverbeke)
  • Apache Wave — Operational Transformation paper (Wang et al.)
  • Shapiro et al. — Conflict-free Replicated Data Types (CRDTs)
  • Yjs — Awareness, YATA internals
  • Notion — How we made Notion available offline
  • Atlassian — Administering Confluence Synchrony
  • Lamport — Time, Clocks, and the Ordering of Events
  • AWS — Exponential Backoff and Jitter (Marc Brooker)
  • Kafka KIP-405 — Tiered Storage

Now defend it

Reading a design is not the same as being able to hold one under questioning. The workspace asks the same questions an interviewer would, and the simulator disagrees with you when the diagram does not support the claim.

Work Collaborative Editor (Google Docs) yourself