Skip to content

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:

Inside one node: three paths through the same components A load balancer sends each request to any node, here narad-1. Inside the node, the HTTP API takes it. A produce goes straight into the node's ingress WAL, and the message, ord_123, is answered once it is fsynced there; the produce dispatcher then reads it back and commits it, to the broker engine's partition logs when this node owns the partition, or over QUIC to the owner on another node. A consume or an ack goes to the router, which looks the partition up in the node's own metastore replica and hands it to the broker engine when the partition is owned here, or forwards it to the owner. The fan-out runner reads the owned partition logs and commits child copies, here or on the child partition's owner. One node: narad-1 Load balancer HTTP API produce ord_123 Ingress WAL Produce dispatcher commit here consume, ack Router owned here or forward to the owner Metastore Raft replica routes Broker engine owned partition logs Fan-out runner reads logs, commits child copies Owners on other nodes over QUIC Inside one node: three paths through the same components A load balancer sends each request to any node, here narad-1. Inside the node, the HTTP API takes it. A produce goes straight into the node's ingress WAL, and the message, ord_123, is answered once it is fsynced there; the produce dispatcher then reads it back and commits it, to the broker engine's partition logs when this node owns the partition, or over QUIC to the owner on another node. A consume or an ack goes to the router, which looks the partition up in the node's own metastore replica and hands it to the broker engine when the partition is owned here, or forwards it to the owner. The fan-out runner reads the owned partition logs and commits child copies, here or on the child partition's owner. One node: narad-1 Load balancer HTTP API Metastore Raft replica routes Router consume, ack produce ord_123 Ingress WAL Produce dispatcher commit here owned here or forward to the owner Broker engine owned partition logs Fan-out runner reads logs, commits child copies Owners on other nodes over QUIC
A produce is answered once the node's own ingress WAL has fsynced it; the dispatcher then commits it here or over QUIC to the owner. Consumes and acks go through the router, which decides from the local metastore replica whether the partition is served here or forwarded.

Across the cluster, every node holds a full copy of the metadata and the partitions it owns:

Metadata on every node, each partition on one A Narad cluster of three nodes, narad-0, narad-1 and narad-2. Each holds a full replica of the metastore, the same topics, users and assignments on all three, kept in step by Raft from the leader, narad-0, which also runs the controller. Message data is different: each partition lives on one node only, orders/0 on narad-0, orders/1 on narad-1 and orders/2 on narad-2, so orders/1 is one copy, on one node. The nodes send each other commits and forwarded requests over QUIC. Narad cluster Raft: every node holds all metadata narad-0 metastore all topics, users, assignments orders/0 controller: leader only narad-1 metastore all topics, users, assignments orders/1 one copy, on one node narad-2 metastore all topics, users, assignments orders/2 QUIC: commits and forwards between any two nodes Metadata on every node, each partition on one A Narad cluster of three nodes, narad-0, narad-1 and narad-2. Each holds a full replica of the metastore, the same topics, users and assignments on all three, kept in step by Raft from the leader, narad-0, which also runs the controller. Message data is different: each partition lives on one node only, orders/0 on narad-0, orders/1 on narad-1 and orders/2 on narad-2, so orders/1 is one copy, on one node. The nodes send each other commits and forwarded requests over QUIC. Narad cluster Raft: every node holds all metadata narad-0 metastore all topics, users, assignments orders/0 controller: leader only narad-1 metastore all topics, users, assignments orders/1 one copy, on one node narad-2 metastore all topics, users, assignments orders/2 QUIC: commits and forwards between any two nodes
Every node holds the same metadata, kept in step by Raft, but each partition lives on exactly one node: 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:

One message, end to end The life of one message, ord_123, across five actors, in order but not to scale. 1: the producer sends POST /produce to narad-0, which fsyncs it into its ingress WAL and answers 202 Accepted. 2: narad-0's dispatcher commits it over QUIC to narad-2, the partition's owner, which appends it, fsyncs, reads it back and checks the CRC, then advances the high watermark, and ord_123 becomes visible. narad-2 confirms, and narad-0's WAL entry becomes reclaimable. 3: the fan-out cursor on narad-2 commits a copy to narad-1, the child partition's owner. 4: a consumer's GET /consume gets ord_123 from narad-2 with a receipt handle. 5: the consumer acks it with POST /ack and gets 204. Producer narad-0 accepting narad-2 owner narad-1 child owner Consumer 1 POST /produce 202 Accepted fsync into the WAL 2 commit over QUIC append, fsync, read back, CRC ord_123 high watermark: visible WAL entry reclaimable 3 fan-out copy child copy committed 4 GET /consume 200, receipt handle 5 POST /ack 204 One message, end to end The life of one message, ord_123, across five actors, in order but not to scale. 1: the producer sends POST /produce to narad-0, which fsyncs it into its ingress WAL and answers 202 Accepted. 2: narad-0's dispatcher commits it over QUIC to narad-2, the partition's owner, which appends it, fsyncs, reads it back and checks the CRC, then advances the high watermark, and ord_123 becomes visible. narad-2 confirms, and narad-0's WAL entry becomes reclaimable. 3: the fan-out cursor on narad-2 commits a copy to narad-1, the child partition's owner. 4: a consumer's GET /consume gets ord_123 from narad-2 with a receipt handle. 5: the consumer acks it with POST /ack and gets 204. Producer narad-0 producer's POST /produce: fsync into the WAL, 202 1 narad-2 commit over QUIC: append, fsync, read back, CRC 2 ord_123 high watermark: visible narad-0 WAL entry reclaimable narad-1 fan-out copy committed 3 narad-2 consumer's GET /consume: 200, receipt handle 4 narad-2 consumer's POST /ack: 204 5
The producer waits for one fsync on the node it reached; everything after the 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:

  1. Delivery contract: what Narad promises, and what each failure does to your messages.
  2. Linearizability: that contract as a property, and how it is tested.
  3. Produce path: from a produce request to a durable, visible record.
  4. Storage engine: segments, frames, the high watermark and retention.
  5. Consume path: leases, acks and how they reach the disk, and consumes that cross nodes.
  6. Fan-out engine: cursors that copy a parent's log to its children.
  7. Metastore and Raft: the replicated metadata, and the stale-replica problem.
  8. Networking and security: the HTTP and QUIC planes, authentication and authorization.
  9. Cluster lifecycle: bootstrap, joins, crash recovery and topic incarnations.
  10. 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/narad Wiring: config, boot order, the join loop, startup reconcile. serve.go is the table of contents for the whole process
internal/transport/httpserver HTTP routes, auth middleware, handlers. Thin on purpose
internal/cluster Everything node-to-node: router, produce dispatcher, fan-out runner, QUIC RPC client and server, leader confirmation
internal/broker/ingress The produce WAL: accept, replay, checkpoint, compaction
internal/broker/messaging The engine: produce commit, consume, ack, fan-out slab reads
internal/broker/runtime Partition-log registry, offset committer, orphan sweeps, lifecycle
internal/consumer The in-flight lease table: reservations, nonces, acked-ahead sets
internal/persistence/storage Partition log engine: segments, frames, flusher, retention, high watermark
internal/persistence/wal The generic segmented WAL under ingress
internal/persistence/metastore Raft and the bbolt state machine: topics, members, users, assignments
internal/remote, 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/security/remotecred New in v3.2.0. Sealing and opening remote passwords: key derivation, AES-256-GCM, fingerprints, the secret strength rule
internal/domain/* Pure types: topic, user, records, remotes
internal/platform/* 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