Architecture overview¶
Learn how Narad is put together: what every node runs, where data lives, and how one message travels from producer to consumer.
The Understand pages describe the code on master. Where v3.0.1 behaves differently, the text says so; see which release these docs describe.
The whole system¶
Every Narad node runs the same binary with the same components. Nodes differ only in which data they own and whether they currently lead Raft.
Inside one node, a request takes this path:
Across the cluster, every node holds a full copy of the metadata and the partitions it owns:
orders/1 is one copy, on narad-1.Five design ideas¶
1. Any node accepts any request; ownership decides where data lives. Clients talk to any node through the load balancer. Each partition has exactly one owner, the node whose disk holds its data. A node forwards what it cannot serve locally to the owner over QUIC. See Networking and security.
2. Produce writes a local log first.
A produce is fsynced into the receiving node's ingress WAL and answered 202 at once. A background dispatcher then moves it to the partition's owner, which commits it durably. The client waits for one local fsync; delivery to the owner is asynchronous and retried until it succeeds. See Produce path.
3. Metadata is replicated by Raft; message data has one owner. Topics, users, assignments and fan-out links live in the Raft-replicated metastore: every node holds a full replica, and one node leads. Message data is not replicated: one owner, one copy, fsynced and verified. See Metastore and Raft and One copy per partition.
4. Consume is a queue with leases. Each message is reserved on its own, even in a batch consume (v3.1.0), which takes up to 100 in one request. A reservation is a lease that lasts the visibility timeout, held in the owner's memory, with a durable committed frontier behind it. Acks move the frontier forward, and a background committer writes it to disk; a crash only means redelivery. See Consume path.
5. Fan-out tails the parent's log. A child topic is fed by a cursor on the owner of each parent partition. The cursor reads committed parent records in bulk, commits them to the child, and keeps its own durable position, so no parent message is skipped and none is copied twice except by at-least-once retries. A delay child adds a due-time gate to the same cursor. See Fan-out engine.
One message, end to end¶
A produce is accepted by any node and committed on the partition's owner; once it is visible, a fan-out cursor copies it to child topics and a consumer takes it from the owner:
202 Accepted happens in the background, and ord_123 becomes visible only when the owner advances its high watermark. The steps are in order, not to scale.Design principles¶
Two rules appear in every subsystem. Both came out of chaos testing, where each was learned from a bug that lost data.
- Destroying data needs the Raft leader's confirmation. No node deletes data (topic directories, cursor files, WAL records) on its local view alone, because a freshly restarted replica can be far behind while believing it is current. Every deletion asks the leader first, and a node that is the leader must pass a Raft barrier before it trusts its own state.
- Every failure keeps the data. When a check cannot complete (no leader, a peer unreachable, a barrier failed), the answer is always to keep the data and try again later.
The bugs that taught these rules are in Cluster lifecycle.
Reading order¶
The pages in this section follow the path of a message, then the cluster around it:
- Delivery contract: what Narad promises, and what each failure does to your messages.
- Linearizability: that contract as a property, and how it is tested.
- Produce path: from a produce request to a durable, visible record.
- Storage engine: segments, frames, the high watermark and retention.
- Consume path: leases, acks and how they reach the disk, and consumes that cross nodes.
- Fan-out engine: cursors that copy a parent's log to its children.
- Metastore and Raft: the replicated metadata, and the stale-replica problem.
- Networking and security: the HTTP and QUIC planes, authentication and authorization.
- Cluster lifecycle: bootstrap, joins, crash recovery and topic incarnations.
- Rebalance and decommission: moving a partition between nodes without losing a record.
Codebase map¶
If you are about to read the source, start from this map:
| Package | What lives there |
|---|---|
cmd/ |
Wiring: config, boot order, the join loop, startup reconcile. serve.go is the table of contents for the whole process |
internal/ |
HTTP routes, auth middleware, handlers. Thin on purpose |
internal/ |
Everything node-to-node: router, produce dispatcher, fan-out runner, QUIC RPC client and server, leader confirmation |
internal/ |
The produce WAL: accept, replay, checkpoint, compaction |
internal/ |
The engine: produce commit, consume, ack, fan-out slab reads |
internal/ |
Partition-log registry, offset committer, orphan sweeps, lifecycle |
internal/ |
The in-flight lease table: reservations, nonces, acked-ahead sets |
internal/ |
Partition log engine: segments, frames, flusher, retention, high watermark |
internal/ |
The generic segmented WAL under ingress |
internal/ |
Raft and the bbolt state machine: topics, members, users, assignments |
internal/, internal/remote/sink |
New in v3.2.0. Remote replication: the credential cache, the address guard, the target checks, and the send path of a remote child (lanes, chunks, the gate, the answer classifier) |
internal/ |
New in v3.2.0. Sealing and opening remote passwords: key derivation, AES-256-GCM, fingerprints, the secret strength rule |
internal/ |
Pure types: topic, user, records, remotes |
internal/ |
Config, metrics, partitioner, network utilities |
One produce, by function name¶
The same journey as the diagram above, with function names you can search for, from the handler to the disk:
POST /v1/topics/orders/produce?key=k
└─ messaging.Produce (transport/httpserver/handlers/messaging/produce.go)
└─ Engine.AcceptProduce (broker/messaging/produce_accept.go)
├─ resolveAcceptedProducePartition: hash the key, no liveness
│ check
└─ ingress.Manager.AcceptProduceWithTopicID
(broker/ingress/produce.go)
└─ wal.Log.AppendWith: staged into the group-commit
buffer, blocks until the shared fsync lands
<- the 202 line
... in the background, woken as soon as the record is durable ...
└─ ProduceDispatcher.run (cluster/produce_dispatcher.go)
├─ read, then place (cluster/produce_dispatch.go): read the newly
│ durable records, queue each on its (topic, partition),
│ reroute dead owners
├─ launch, startCommit, runJob (cluster/produce_commit.go): at
│ most one commit in flight per partition, local or to the
│ owner over QUIC
│ └─ Engine.CommitAcceptedProduceBatch
│ (broker/messaging/produce_commit.go)
│ ├─ storage.Log.AppendBatchOwned: keyed-envelope records
│ └─ commitDurable, then storage.Log.CommitDurable:
│ fsync, CRC read-back, advance the high watermark
│ <- consumers can see it
└─ advanceCheckpoint: first seq not yet committed, stored; the
WAL compacts behind it
Each deep dive follows the same pattern: the concept first, then the constants and function names behind it.
Next steps¶
- Delivery contract: the promises these components add up to.
- Produce path: the first stop on a message's journey, in detail.