Dropbox / Google Drive — a worked solution
Chunking, dedup, sync, conflict resolution.
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 Dropbox / Google Drive workspaceThe problem
Build a file-sync and storage service like Dropbox or Google Drive. A user drops a file into a synced folder on one device; it must appear, byte-identical, on every other device they own and on every device of everyone they've shared the folder with — quickly, durably, and without clobbering a concurrent edit. The product is two systems wearing one trench coat: a content plane that stores the bytes (immutable, content-addressed, deduplicated 4 MB blocks) and a metadata plane that stores the filesystem (a journaled, ordered, strongly-consistent index of which blocks compose which file). Clients sync by following a monotonic cursor over a long-poll loop, and divergent edits are resolved by keeping both — the "conflicted copy."
This canonical models the main sync flow: upload (chunk → dedup → store → commit), download, receive-a-remote-change, and conflict. It deliberately omits the adjacent sprawl (preview/thumbnail rendering, full-text search, sharing-invite UX, billing) to keep the spine sharp.
The reference architecture
What each component is for
- Sync Client (Desktop / Mobile)
Watches the local filesystem, splits each changed file into fixed 4 MB blocks, SHA-256-hashes each block to form an ordered blocklist, and uploads only the blocks the server doesn't already have. Maintains three trees — remote, local, and last-synced — and a per-namespace cursor (the journal ID, JID). It long-polls for changes, pulls the authoritative delta by cursor, and runs a deterministic Planner that computes the operation sequence to converge the trees.
Why it exists. The client is the reconciler, not a thin uploader. We considered a dumb client that PUTs whole files and lets the server diff; rejected because (1) without client-side chunking + hash-omission you re-upload gigabytes for a one-byte change, and (2) convergence (which of two divergent edits wins, what becomes a conflicted copy) is only decidable where all three trees are visible — and that's the device. Dropbox's Sync Engine 3 / 'Nucleus' rewrite (Python → Rust) exists precisely because the original couldn't make this reconciliation deterministic and testable.
When it fails. A clock-skewed or buggy client can resolve a concurrent edit as last-write-wins on a bad timestamp → silent lost update; the fix is to order by the server-assigned monotonic JID, never device wall-clock. A client that re-syncs a conflicted copy as a fresh edit spawns conflict-of-conflict loops. Detection: per-device conflict-creation rate, cursor-reset rate, distinct-blocks-per-file anomaly.
- Edge / Block CDNDropbox edge POPs + CDN
Terminates TLS near the user and caches block downloads at the edge. Because blocks are immutable and content-addressed (the URL IS the SHA-256), a cached block can never be stale — the edge serves a hot 4 MB block without ever touching the origin. Absorbs the bulk of read egress (download:upload runs ~3:1).
Why it exists. Block egress is the single largest bandwidth line item (~300 GB/s peak); serving every download from Magic Pocket origin would saturate the storage fabric and pay cross-region transit on popular blocks (installers, shared team assets). Content-addressing makes the edge trivially correct — no invalidation logic, infinite TTL — so the cache hit is free correctness, not a consistency risk.
When it fails. An edge POP outage shifts its block traffic to origin → Magic Pocket read load + cross-region egress spike; download p99 climbs but nothing is lost (content is durable at origin). A signed-URL signing-key rotation bug 403s downloads fleet-wide. Detection: edge hit-ratio, origin-offload bytes/s, 403 rate on block GETs.
- API GatewayEnvoy + WAF + TLS terminator
L7 ingress for every sync RPC and block transfer. Terminates TLS, runs the WAF ruleset, enforces per-device and per-user token-bucket rate limits BEFORE forwarding (so a runaway client can't storm the journal), validates the device token with the auth service, and routes by intent: metadata RPCs to the sync/metadata service, block PUT/GET to the block service, long-poll registrations to the notification service.
Why it exists. Considered inlining auth + rate-limit inside each service; rejected because (1) a crash in the sync service would skip the rate limiter and let a retry storm reach the journal, and (2) rate-limit + ACL state is per-user and centralizing it at the edge beats replicating it to every service replica. The gateway is also where we shed load and quarantine a misbehaving device or namespace without touching the data services.
When it fails. Gateway pool exhaustion reads to users as 'sync is slow' but is usually upstream backpressure. Detection: gw_active_streams, downstream_5xx_rate, rate-limit-reject rate. Mitigation: drain mode + per-zone failover; circuit-break the offending device/namespace rather than 503-ing everyone.
- Auth / Namespace ACLOAuth2 device tokens + ACL service
Validates the device's OAuth2 token and resolves namespace membership: which namespaces (personal root, shared folders, team spaces) this device may read or write, at what permission level. Returns a signed envelope (device_id, member_of[], membership_epoch) the gateway and services trust for ~60 s without re-checking.
Why it exists. File sync is multi-tenant on shared namespaces; the authority for 'can this device commit to folder X' must be separate from the sync logic so a CVE in the sync service can't bypass sharing, and so a permission change propagates from one place. Folding ACLs into the metadata service would put a 100M-edge membership graph on the commit hot path.
When it fails. Auth down → writes fail-closed (503) rather than fail-open; reads with an unexpired envelope survive 60 s then 401. A stale membership cache lets a just-removed collaborator read for up to the cache TTL. Detection: auth_5xx_rate, envelope_cache_miss_rate, epoch-skew at commit. Microsoft Teams' Feb-2020 outage (an expired auth certificate took Teams down for hours) is the lesson: auth must not depend on an artifact that can silently expire — rotate ahead of time and alert on the countdown.
- Sync / Metadata ServiceNucleus server-side (Rust)
The metadata brain. Serves commit (apply a filesystem mutation → allocate the next per-namespace JID → write the journal rows + blocklist), list_folder / list_folder_continue (paginate the delta since the client's cursor), and longpoll registration. On commit it compares the client's base cursor against the namespace's current JID; if another device already advanced past it on the same file, it does NOT overwrite — it materializes a conflicted copy as a second entry.
Why it exists. Sync correctness is a server-arbitrated, ordered-journal problem: someone must assign the single monotonic JID that defines 'what changed since' and adjudicate concurrent edits. That arbitration can't live on the client (no global view) or in the storage engine (no domain logic). We keep read (list) and write (commit) in one service because they share the cursor/JID machinery; the block bytes are a separate plane entirely.
When it fails. If commit returns success but the journal write didn't durably land, the client's cursor advances past data it can't see → phantom-delete on the next sync. If conflict detection is wrong, concurrent edits silently overwrite (lost update). A buggy conflict path that re-commits the conflicted copy spawns conflict-of-conflict loops; the guard is a per-namespace conflict-rate circuit breaker. Detection: commit_5xx, conflict_create_rate, cursor-vs-JID skew, list_folder p99.
- Metadata CacheRedis (cluster mode)
Caches the hot read shapes the metadata service serves: a namespace's current JID, recent list_folder pages, and file-metadata rows. Absorbs the ~95% of metadata reads that would otherwise hit the sharded journal, keeping list_folder p99 in single-digit ms.
Why it exists. list_folder is read-dominated and bursts hard during notification fan-out (every poked device pulls a delta at once). Without a cache that read amplification lands directly on the journal shards. The cache is a strict performance optimization — a cold cache is correct, just slower — so it never becomes the source of truth.
When it fails. Cold cache (mass eviction / failover) sends list_folder straight to the journal shards → read amplification, especially mid-fan-out — risk of a metadata meltdown. Mitigation: request coalescing + early-recompute on hot keys to avoid a stampede. Detection: cache_hit_ratio, journal read QPS, list_folder p99.
- Filesystem Journal (SFJ / Edgestore)Sharded MySQL/InnoDB + Panda (2PC + HLC)
Stores the filesystem as metadata — file/folder entities, blocklists (the ordered block hashes that compose each file, NOT the bytes), and a monotonically-increasing Journal ID per namespace. The Server File Journal is what makes 'what changed since JID N' answerable. Panda layers an ordered, transactional KV with hybrid-logical-clock timestamps and MVCC over thousands of MySQL shards, splitting data into ~100 GB ranges so a hot shard can be split rather than outgrowing a machine.
Why it exists. A single MySQL primary caps at ~5–10 k writes/s and ~1e9 rows; a trillion-file system has ~1e13 metadata rows — it blows the row ceiling by four orders of magnitude, forcing thousands of shards. We need ordered, transactional, strongly-consistent metadata (atomic cross-folder moves, exact deltas) that a leaderless eventual store can't give. Edgestore's lesson was that by 2019 an outlier shard 'could outgrow a single machine with no way to split it' — Panda's range-splitting is the fix.
When it fails. One giant shared namespace or a millions-of-files user creates a hot shard → per-shard lock-wait, noisy-neighbor stalls for everyone on it; mitigation is namespace-range sharding + whale isolation + write quotas. A schema migration that locks or unindexed-scans a trillion-row store is Dropbox's 2014 outage class — online/dual-write migrations + canary-one-shard + host-targeting guards are mandatory. Detection: per-shard write p99, lock-wait, replication lag, oldest_unreplicated_lsn.
- Block ServiceBlockServer (stateless)
Handles block transfer. On PUT, it looks the block's SHA-256 up in the dedup index: on a hit it skips storage entirely and just bumps the reference (free cross-user/cross-version dedup); on a miss it writes the 4 MB block to Magic Pocket and registers hash→location. On GET, it resolves the hash to a location and streams the bytes, verifying the stored checksum before serving.
Why it exists. The bytes plane is deliberately separate from the metadata plane: blocks are large, immutable, and content-addressed, while metadata is small, mutable, and ordered — they have opposite storage, consistency, and scaling shapes. A combined service would couple a 300 GB/s egress workload to a transactional journal. Splitting them lets each scale independently and keeps a download storm off the commit path.
When it fails. If the dedup index returns a false positive and we skip the byte-compare, a commit references a block whose bytes differ from what the user uploaded → silent corruption, the worst page. Mitigation: verify bytes, quarantine on mismatch. A giant or rapidly-rewritten file floods PUTs; per-file commit-rate throttles + content-addressing (only changed blocks move) contain it. Detection: read-time checksum-mismatch counter, dedup false-positive rate, PUT p99.
- Block Index (hash → location)Sharded MySQL-backed KV
The dedup decision store: maps each block's SHA-256 to where it lives (cell, bucket, checksum) and tracks a reference count. Answers two hot questions — 'do we already have this hash?' (the dedup check that lets the client skip an upload) and 'where are this hash's bytes?' (the locate for a download).
Why it exists. Dedup and locate need a key-by-content lookup with a very different access pattern from both the filesystem journal (ordered, range-scanned) and the block store (large blobs). Folding hash→location into the journal would bloat it and couple dedup correctness to filesystem transactions; folding it into Magic Pocket would put a billion-key hot index inside the blob layer. It earns its own tier.
When it fails. Index corruption or a hash collision maps two distinct blocks to one address → cross-file corruption; the defense is byte-verify at the block service plus a background index↔store consistency scrub. The subtler killer is the GC orphan-resurrect race: block H drops to refcount 0 (last referrer deleted) and is queued for GC at the SAME moment another user's dedup hit re-references it — without the tombstone+grace+CAS fence above, the sweep deletes bytes a freshly-committed file points at, silently losing a 'durable' file. Losing a shard makes blocks on it un-locatable until the follower is promoted. Detection: index/store scrub-mismatch count, tombstone-resurrect rate, refcount-vs-store divergence, shard availability, existence-check p99.
- Magic Pocket (Block Store)Magic Pocket — OSDs / buckets / volumes / cells
The exabyte content plane. Stores immutable, SHA-256-addressed blocks (≤4 MB) aggregated into 1 GB buckets, grouped into volumes that are replicated or erasure-coded across physical OSDs (Object Storage Devices), grouped into ~50 PB cells each driven by a soft-state Master that handles placement, repair, and GC. Recent writes are replicated for fast durability; aged blocks migrate to Reed-Solomon erasure coding to cut storage cost.
Why it exists. At ~1 EB, storage economics dominate: owning the stack and erasure-coding cold data instead of paying S3 saved Dropbox ~$75 M (per its pre-IPO S-1) and let it run exotic hardware (SMR drives). Content-addressed immutability is what makes this layer simple — no updates, no versioning, perfect cacheability — so all the mutability/ordering complexity stays in the metadata plane.
When it fails. OSD/disk failures trigger RS repair; a whole-region rebuild can saturate the inter-node network and starve foreground reads (a cascading failure) — mitigation is rate-limited, durability-prioritized repair under a network budget with capacity headroom. A cell Master bug mis-places or over-GCs blocks (the GC orphan-resurrect race is owned by the GC sweep — see the dedup index). The only write RPO exposure is a block ack'd locally but not yet cross-zone-replicated (~1 s window) when its zone dies: recent writes carry RPO ≈ cross-zone lag, older data is RPO 0. Detection: durability-margin (live fragments vs repair threshold), rebuild throughput vs network budget, OSD-down count, read p99.
- Notification Service (long-poll fan-out)epoll long-poll fleet
Holds every active device's long-poll connection and, when a namespace it cares about changes, wakes exactly the devices subscribed to that namespace so they pull the delta. It consumes namespace-change events off the notification stream and maps them to held connections. It carries NO file data — only the poke 'namespace X advanced, come pull.'
Why it exists. Polling at 120 M devices would hammer the metadata tier for mostly-empty deltas; a held long-poll turns 'tell me when something changes' into a near-zero-cost wait. This must be its own tier because connection-holding (memory + file descriptors for tens of millions of idle sockets) is an utterly different resource profile from the request/response services — you scale it on connection count, not QPS.
When it fails. A shared folder with tens of thousands of members produces a fan-out storm — N pokes → N simultaneous list_folder pulls → metadata read amplification that can bleed onto unrelated users (the thundering herd). Mitigation: per-namespace coalescing + fan-out rate caps + client jitter + degrade-to-polling under load shed. A box failure drops its connections; clients reconnect (stateless). Detection: pokes/s vs commits/s ratio, fan-out amplification, connection count, reconnect rate.
- Change Notification StreamKafka
The durable fan-out backbone between the write side (a committed journal mutation) and the read side (waking devices). The outbox relay publishes one event per namespace advance, keyed by namespace_id; the notification service subscribes and fans out to held connections. Decouples commit latency from notification latency.
Why it exists. Without a broker, the metadata service would have to know about and synchronously push to every notification box — coupling commit throughput to fan-out and making delivery lossy on any notification-tier blip. A partitioned log gives ordered-per-namespace, replayable delivery and lets the notification fleet scale and restart independently. It carries the poke, never file bytes.
When it fails. A poison/oversized event can wedge a partition and stall pokes for every namespace on it (consumer crash-looping on one offset) — mitigation is a dead-letter path that skips after K retries. Consumer lag delays 'my other laptop sees it' but never loses the change (the cursor still reflects truth on the next pull). Detection: per-partition consumer lag, poison/DLQ rate, ISR shrink.
- Journal CDC Relay (outbox)Debezium-style logical-replication tail
Tails the filesystem journal's committed log (logical replication / change-data-capture) and publishes one namespace-change event to the notification stream per advance. It is the bridge that turns a durable commit into a fan-out poke — the outbox pattern: a notification is emitted if and only if the journal mutation actually committed.
Why it exists. If the metadata service published to the stream itself, in the same breath as the DB write, you get the dual-write problem: commit succeeds but the publish fails (or vice-versa) → a change with no poke (other devices never sync it) or a poke with no change (a phantom). Tailing the COMMITTED log makes the poke exactly-once-after-commit and removes that race. CDC also means the metadata service never blocks on the notification tier.
When it fails. If the relay stalls past the journal's WAL-retention window, the replication slot's bytes pile up and threaten the primary's disk (Debezium's classic failure) — and on restart it may have to backfill from the WAL archive. A stalled relay delays every downstream poke (sync feels frozen) but loses nothing (the journal is truth). Detection: replication_slot_lag_bytes, oldest_unpublished_JID_age, relay consumer lag.
- KMS / URL-Signing AuthorityManaged KMS (envelope encryption + URL signing keys)
The key authority for two things the data path depends on: SSE-KMS envelope encryption (the block service unwraps a data-encryption key before writing/reading block bytes) and the signing/verification keys for the time-limited, signed block-download URLs the edge honors. Keys are fetched and cached, not called per-chunk.
Why it exists. Modeled as an external dependency on purpose: it is the highest-blast-radius thing the content plane leans on (a signing-key rotation bug 403s every download fleet-wide; a KMS outage blocks encrypted reads/writes), and treating it as a first-class boundary forces a circuit-breaker, aggressive caching, and an expiry-countdown SLI rather than an implicit assumption. Folding it into auth would hide that blast radius.
When it fails. A bad signing-key rotation or a KMS hard-outage past the cache window 403s/stalls ALL block downloads fleet-wide — the single most-cited 3am page. Mitigation: cached keys + circuit breaker + a key/cert validity-countdown SLI that alerts ≥7 days before expiry. Detection: signing-403 rate, kms_call_error_rate, key-validity-remaining (days).
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:
- Chunking scheme? Fixed-size (simple, boundary-shift problem) vs content-defined (Rabin/FastCDC; robust to inserts). Dropbox ships fixed 4 MB — defend the trade-off.
- Dedup scope? Per-user, per-team, or global cross-user dedup? Global maximizes savings but raises a privacy/side-channel question (can you probe whether a block exists?).
- Consistency for "see my own change"? A device that commits then lists must see its own write (read-your-writes), or the reconciler diverges.
- Conflict policy? Merge (only for known formats) vs keep-both/"conflicted copy" (the universal default for opaque files).
- Sharing fan-out scale? A 50 k-member shared folder changes the notification design entirely.
- Multi-region? Single-region metadata with DR, or active-active? Drives RTO/RPO and the conflict surface.
- Offline editing? Yes — that's the whole point, and it's why conflict resolution is load-bearing.
Assumptions to state:
- 500 M registered users, ~12% weekly-active, ~2 devices each ⇒ ~120 M long-poll connections.
- ~1 EB logical, ~4 MB blocks, ~35% dedup savings.
- Sync-notify "other laptop sees it" target: sub-second p50, single-digit-second tail.
02Functional reqs
What must this system actually do?
- Upload a file: chunk into blocks, upload only blocks the server lacks, commit a blocklist.
- Download a file: resolve its blocklist, fetch blocks (edge-cached).
- Detect changes since last sync via a per-namespace cursor (list_folder / list_folder_continue + long-poll).
- Resolve concurrent edits: keep both as a "conflicted copy"; never silently overwrite.
- Atomic move/rename independent of subtree size.
- Share a folder (namespace) with other users; their devices sync it too.
- Deduplicate identical blocks across files and across users.
03Non-functional
What must it promise about speed, uptime and correctness?
- Durability: 12 nines on stored blocks. A committed file is never lost. This is the headline guarantee.
- Availability: 99.9%+ on sync; reads (downloads) higher. Sync may briefly lag; it must not lose or corrupt.
- Latency: metadata commit p99 < 30 ms; block download (warm) p99 < 250 ms; sync-notify delivery p99 < 1 s.
- Consistency: strong, read-your-writes on the metadata journal (cursor correctness). Block reads are read-after-write by construction (content-addressed immutability).
- Correctness: no silent lost updates (order by server JID, not client clock); no silent corruption (verify bytes even on a dedup hit).
- Efficiency: chunking + dedup must materially cut both bandwidth (only changed blocks move) and storage (RS-coded cold data).
04Capacity estimation
How much load and data does this have to hold?
Anchored on published Dropbox figures (>1 EB stored, >500 M users, ~90% on Magic Pocket, ~4 MB blocks, RS erasure coding, Edgestore >10 M req/s).
- Devices: 500 M × 12% × 2 = 120 M active devices holding long-poll connections.
- Metadata commits: ~50 k commits/s peak (uploads + renames + moves + deletes). Each touches ~3–5 rows (file row, namespace journal row, block pointers) ⇒ ~200 k row-writes/s.
- Sync-notify deliveries: 50 k × mean fan-out 4 = 200 k pokes/s (≈12 M/min — matches "tens of millions of notifications").
- Block uploads: 50 k files/s × (2 MB avg / 4 MB block) ≈ 25 k block-writes/s. Downloads ~3× ⇒ 75 k/s ⇒ egress 75 k × 4 MB ≈ 293 GB/s ≈ 2.4 Tbit/s peak (~26 PB/day). The edge absorbs most of this.
- Storage: 1 EB logical → ~0.65 EB physical after ~35% dedup. 3× replication = 1.95 EB raw; RS(6,3) = 0.975 EB raw — a 50% saving (~1 EB of disk), the entire economic case for owning the store.
- Metadata size: ~3×10¹² files × ~4 rows ≈ 1.2×10¹³ rows ≈ 3–10 PB on SSD.
Where the choices flip (these drive the diagram):
- One MySQL primary → thousands of shards. A primary caps at ~5–10 k writes/s or ~1e9 rows. The row ceiling breaks first: 1.2e13 rows / 1e9 ≈ ~12 000 shards (write-QPS would need only ~40). → range-sharded journal + Panda.
- 3× replication → erasure coding. Flips on cost, not capacity, and only above ~100 PB: at 1 EB the 3.0× vs 1.5× gap is ~1 EB of disk (hundreds of millions of dollars). Hot/new blocks stay replicated; cold blocks re-encode to RS(6,3).
- One notification box → a fleet. A long-poll box holds ~500 k–1 M idle conns. 120 M / 700 k ≈ ~170+ boxes minimum → run a few hundred.
- Hot-namespace single-partition ceiling. Pokes are keyed by
namespace_idacross 256 partitions to keep a namespace's changes ordered — but that pins a 50 k-member whale folder's entire feed to ONE partition's consumer throughput. Per-namespace coalescing (one poke per commit-burst, not one-per-member) is what keeps the whale's poke rate under a single partition's consume ceiling; without it, "isolate the hot namespace to one partition" degrades into "bottleneck the whale on one partition."
05API design
What does the outside world call, and what comes back?
# 1a. Pre-flight: ask which block hashes the server LACKS (metadata/dedup-index lookup, no bytes)
POST /2/files/blocks/check { "hashes": ["h1","h2","h5",...] } → { "missing": ["h2","h5",...] }
# 1b. Upload ONLY the missing blocks (the check precedes upload — append never carries a block the server already has)
POST /2/files/upload_session/append (per missing 4 MB block, keyed by block SHA-256)
# 2. Commit the blocklist as a filesystem mutation
POST /2/files/commit
Idempotency-Key: {device_id}:{client_op_id}
{ "path": "/Reports/q3.xlsx", "blocklist": ["h1","h2",...], "base_cursor": "JID:8841" }
201 { "id": "id:f9...", "rev": "JID:8842", "content_hash": "..." }
# A replayed key returns the ORIGINAL 201 (the journal keeps key → JID for the retry window) — without it,
# a retried commit whose first attempt landed re-arrives with a stale base_cursor and conflicts with itself.
# If another device advanced past base_cursor on this file:
201 { "id": "id:f9...", "conflict": true, "name": "q3 (Alice's conflicted copy 2026-06-06).xlsx" }
# 3. Enumerate / detect changes by cursor
POST /2/files/list_folder → { entries:[...], cursor:"JID:8842", has_more:true }
POST /2/files/list_folder/continue { cursor } → next page …
POST /2/files/list_folder/longpoll { cursor } → blocks until change or timeout → { changes:true }
# Stale/garbage cursor → 409 reset → client must re-list from scratch.
# 4. Download blocks (capability flow; key IS the content hash, but the bare hash is NOT the authz)
# The metadata response mints a signed, expiring block URL AFTER the namespace ACL check:
# → { "block_url": "/2/blocks/{sha256}?exp=1717689600&sig=..." }
GET /2/blocks/{sha256}?exp=...&sig=... → 4 MB bytes (immutable; cache forever)
# The edge verifies the signature (KMS-fetched key) before serving — authz rides the signature, not the hash.06Data model
What gets stored, and what is it looked up by?
Metadata plane — the filesystem journal (sharded MySQL / Panda), range-sharded by namespace_id:
| entity | key | notes |
|---|---|---|
| namespace | namespace_id | personal root / shared folder / team space; owns a monotonic JID |
| file/dir entry | (namespace_id, path or file_id) | name, parent, file_id (stable across moves), latest rev (JID) |
| journal row | (namespace_id, JID) | the ordered log: "at JID N, this changed" — the cursor target |
| blocklist | file_id → [block_hash…] | ordered block hashes (NOT bytes) |
| membership | (namespace_id, user_id, epoch) | ACL + monotonic epoch for commit-time fencing |
Content plane — the block store (Magic Pocket), keyed by block_hash:
| store | key | value |
|---|---|---|
| block | sha256(block) | immutable ≤4 MB bytes (replicated → RS(6,3) when cold) |
| block index | sha256(block) | → (cell, bucket, checksum, refcount) — the dedup/locate map |
Why two stores: the journal is small, mutable, ordered, and transactional; blocks are large, immutable, content-addressed, and erasure-coded. Opposite shapes — coupling them couples a 300 GB/s egress workload to a transactional journal.
07High-level design
Which components handle a request, and in what order?
Read/sync path (top): client ⇄ notification-service long-poll; on a change the device pulls the delta client → gateway → metadata-service → metadata-cache → (miss) metadata-store, then fetches bytes client → CDN → (miss) block-service → block-store.
Write path (bottom): upload blocks client → gateway → block-service → dedup-index → (miss) block-store, then commit client → gateway → metadata-service → metadata-store (allocate JID; invalidate the metadata cache).
Fan-out (async): metadata-store → outbox-relay → notification-bus → notification-service → poke devices. The outbox relay tails the committed journal so a poke is emitted iff the commit durably landed.
The eight authored traces walk exactly these journeys: CDN-hit download, CDN-miss download, receive-remote-change (wake + delta-pull), upload-new-file (blocks + commit), upload-dedup-hit, commit-conflict, change-fan-out, and the cursor-reset re-index storm. The KMS / signing authority (kms) is a control-plane dependency of the content plane (block-service unwraps SSE-KMS data keys; the edge verifies signed block URLs), modeled as kind:external so it carries a circuit breaker and an expiry-countdown alert rather than living as an implicit assumption.
08Deep dives
Which part breaks first, and what do you do about it?
1. Fixed 4 MB blocks vs content-defined chunking. Fixed offsets suffer boundary-shift: insert one byte at the front of a 1 GB file and every block re-hashes → full re-upload. CDC (rsync rolling checksum → Rabin → Gear/FastCDC) cuts on content so an insert only shifts one chunk. Dropbox chose fixed 4 MB anyway: most files are < 4 MB (sub-block; whole-file identity dominates dedup), and fixed blocks keep the blocklist and index small. CDC's win is concentrated in append-heavy large files; the metadata cost wasn't worth it for the average mix.
2. Dedup correctness — verify bytes even with SHA-256. The index says hash H exists, so we skip the upload and point the new file at the stored block. If that "exists" is a false positive — pointer rot, a truncated block, or an adversarial collision (cf. SHAttered) — we've silently corrupted a file. The hash is an index hint, not proof of equality: for correctness-critical writes, byte-compare against the stored block before collapsing, and recompute the checksum on read so bit-rot is caught, not served.
3. Conflict resolution — three trees, keep both. Nucleus models remote, local, and last-synced trees and a Planner that computes the op sequence to converge them. A conflict is when local AND remote both diverge from last-synced. For opaque files there's no safe merge, so we keep both and rename one "conflicted copy" (editor + phrase + date). The trap: a buggy client that re-syncs the conflicted copy as a NEW edit spawns conflict-of-conflicts unboundedly — defend with deterministic conflict naming (a re-observed conflict maps to the same name, not a new file) + a per-namespace conflict-rate circuit breaker.
4. Notification fan-out storm. One edit in a 50 k-member folder → 50 k pokes → 50 k simultaneous list_folder pulls. The metadata tier melts first and can bleed onto unrelated users on the same shard. Coalesce/debounce pokes per namespace (one wake per burst), cap fan-out rate, jitter the client long-poll, isolate whale namespaces onto their own shards, and degrade to slow-polling under load shed.
5. Cursor invalidation → re-index storm. A cursor-format change that 409s every cursor at once sends 120 M devices into a full re-list simultaneously → metadata meltdown. Version the cursor and migrate it forward without invalidating; if invalidation is unavoidable, stagger re-syncs with server-controlled jitter and admission control. (Adjacent to GitHub 2018, where a 43 s partition forced expensive fleet-wide reconciliation.)
09Trade-offs
What did this design cost, and what breaks at 10×?
- Fixed 4 MB blocks over CDC. Simpler index, smaller blocklists, good enough for the average file; we accept worse dedup on insert-heavy large files. Rejected: per-file CDC (Rabin/FastCDC) — better dedup, but the variable-chunk index overhead and complexity weren't justified by the corpus.
- Two planes (content vs metadata) over one unified store. We accept the operational cost of two systems + an outbox bridge to get independent scaling and blast-radius isolation. Rejected: one store for both — couples 300 GB/s egress to a transactional journal.
- Keep-both conflict over merge. Universal correctness for opaque bytes; we accept duplicate files the user must reconcile. Rejected: automatic content merge — unsafe for anything but known structured formats.
- In-region sync replication + async cross-region DR over active-active. RPO 0 where it matters, no per-commit cross-region latency, and a simpler conflict surface. Rejected: active-active metadata — multi-master conflict on the journal is a nightmare and a split-brain risk.
- RS(6,3) cold / replicated hot over uniform replication. 50% raw-disk saving at equal fault tolerance. We accept RS read-repair complexity (reconstruct 6 fragments on failure), which is why hot/new data stays simply replicated until it cools.
Primary sources
- Dropbox — Inside the Magic Pocket
- Dropbox — Streaming File Synchronization (journal + cursors)
- Dropbox — Rewriting the heart of our Sync Engine (Nucleus)
- Dropbox — Panda: petabyte-scale transactional KV over sharded MySQL
- Dropbox — Reintroducing Edgestore
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 Dropbox / Google Drive yourselfBuild the primitives this design leans on
Each one is an animated curriculum that constructs the system from scratch.
- Build Build a Bitcask-style KV storeThe simplest possible KV store that still works: an append-only log on disk + an in-memory hash index. Build it from first principles and feel which trade-offs every later store inherits.
- Build Build an LSM-tree storage engine (LevelDB / RocksDB style)The simplest possible storage engine that gives you BOTH ordered reads AND more keys than fit in RAM, by accepting a deal: write to RAM at memory speed, log to disk for safety, then merge sorted files in the background forever.
More in Object Storage, File Sync & Media Delivery
Durable bytes at rest and bytes in flight: an S3-style object store, a Dropbox-style sync client, and adaptive video streaming.
- Build Build an S3-style distributed object storeEleven nines of durability over disks that fail weekly. Build the object store from first principles — flat keyspace, immutable objects, erasure coding instead of replication, eventual consistency turned strong, multipart upload, lifecycle and tiering — and feel why every modern data lake sits on top of something shaped exactly like this.
- YouTube / Netflix StreamingABR, CDN, transcoding, hot/cold — the egress, fan-out, and popularity Pareto.