A node keeps two things on its own disk. The queue log files hold every message, and they are also the write-ahead log: a message is written once, in its queue’s log, and that write is the one that makes it durable. The store holds everything else (queues, partitions, cursors, dedup rows, KV, timers) in RAM, and checkpoints it to an LMDB file once a second. There is no other database, and no write waits for the checkpoint.
This page covers when a write is durable, what a restart replays, the damage that makes a node refuse to boot, and how disk space comes back.
The data directory
$QUEEN_RAFT_DIR default /var/lib/queen/raft
store/data.mdb the store checkpoint (LMDB)
qlog/
LANES SHARDS the log layout, fixed when the directory is created
q0/ the system log: entries that touch no queue
q<queue id>/r00000001.qlog a queue's log files, rolling
r00000001.qidx the index of a sealed file
raft/ openraft only: vote, commit point, purge point, membership
snapshots/ a snapshot received from the leader, staged for the next boot
local.db dash.db node-local metrics and dashboard rows (not replicated)
groups/g1/ ... the same layout for each extra raft group (QUEEN_RAFT_GROUPS)A data directory belongs to the replicator that created it: the local replicator refuses a directory openraft wrote, and the reverse.
The queue logs are the write-ahead log
We built it this way to write each message once. An earlier design kept the payloads in a consensus log and copied them into segment files as entries applied, so every message reached the disk twice. With the queue logs as the log, the bytes that make a message durable are the bytes a pop reads back.
Each queue has its own log directory by default. A record starts with its length and an xxh3
checksum, and each covers every byte after itself, so a file describes itself: a scanner always
knows where the next record starts, and a record torn by a crash fails its checksum. A message
record holds one append, the messages of one push for one partition after dedup: the partition id,
the first offset, the count, the timestamp, the 16-byte hash of each transactionId, and the
payloads.
An entry of the log can touch many queues, and each touched log has to be able to prove, after a
crash, that it holds its part of the entry. So the writer puts a payload-free entry record into
every log the entry touches: the whole record in the lowest one, and a 61-byte stub that names it
in the others. Before the stub layout every touched log got the whole record. With 1,000 queues of
one partition each at 5,000 messages a second on one local node, the copies came to about 14,900
bytes of queue log per message, and the stubs brought it to 236
(benchmark-queen/2026-09-25-qlog-entry-stubs, runs q1000-r5k-copies-np-1 and
q1000-r5k-stub-np-1, no preallocation).
A payload of at least 512 bytes is stored zstd-compressed at QUEEN_RAFT_QLOG_ZSTD_LEVEL
(default 1; 0 turns it off), unless that saves less than a tenth. Compression is node-local: the
replicated entry stays raw, two nodes may store the same record differently, and every read
decompresses at one point.
The active file rolls at QUEEN_RAFT_SEGMENT_BYTES (64 MiB). A sealed file is immutable and gets a
.qidx index beside it, while the active file’s index lives in RAM; a stale or missing .qidx is
rebuilt by scanning its file.
QUEEN_QLOG_SHARDS=K, set before the directory is created, maps every queue onto K shared logs
(queue id % K), so one write fsyncs at most K files however many queues it spans. That is the
layout for many small queues: with one log per queue, a group that touches 10,000
single-partition queues fsyncs 10,000 files.
When a write is durable
For each group of entries the writer puts the payload records into their queues’ logs, the entry
records into every log the entries touch, and then fsyncs each touched log once (fdatasync on
Linux). Only then does the entry go to apply. In a cluster every follower does the same before it
acknowledges the append, and the entry commits once a majority has.
A write is durable when its entry is fsynced in the queue logs of a majority of voters, or on a single node, in its own.
server/src/rsm/qlog/mod.rs, server/src/rsm/store/The store is not on that path. Its keyspaces live in RAM, loaded whole at boot. Every
QUEEN_RAFT_DURABLE_EVERY_MS (1000) or QUEEN_RAFT_DURABLE_EVERY_BYTES (256 MiB), a checkpoint
writes the rows that changed into store/data.mdb, commits and syncs it, and records the applied
index in the same transaction. That checkpoint bounds what a restart replays. No answer ever waits
for it, which is why a store write costs a push nothing.
Crash recovery
- The store reopens exactly at its last checkpoint.
- Each queue log cuts a torn tail from its active file.
- Recovery merges the entry records of every log by index and replays, from the entry after the checkpoint’s index, the gapless run of complete entries. The first missing or incomplete entry and everything after it is cut from every log: that group’s fsyncs did not all land, so it was never acknowledged.
- Apply skips any entry at or below the store’s applied index, so replaying twice changes nothing.
An entry that touches three queues is written into three logs and records that it has three
copies, and it replays only if every copy is present and identical (the present-in-all rule,
planner/txn.rs). This is where transactions get their atomicity on disk. A transaction torn by a
crash is cut from every log and leaves no gap in the offsets, because its offsets were never
applied.
In a cluster a restarted node then catches up from the leader. A node that was away longer than
QUEEN_RAFT_PURGE_HOLD_S (600), or a new one, receives a snapshot instead: a copy of the leader’s
store checkpoint plus its queue log files. The node stages it under snapshots/, exits, and swaps
it in at the next boot, never while it serves.
Refused at boot
A node does not guess its way around damage. It refuses to start when:
- a sealed queue-log file holds a record that fails its checksum (sealed files were fsynced before the roll, so this is corruption, not a torn tail);
- a damaged record in the active file is followed by a verified record the store already counts as durable;
- the queue logs end before the last record the store recorded as durable;
- the copies of one entry disagree across logs;
- a store value fails its checksum. Every value carries an xxh3 over its keyspace, key and bytes, and the load at boot verifies every row, the key order and the row counts.
At runtime, a corrupt value that a checkpoint or the scrub meets ends the process. The scrub walks
the whole store every QUEEN_STORE_SCRUB_EVERY_S (21600) at QUEEN_STORE_SCRUB_ROWS_PER_S (2000),
to find damage on one node before a restart finds it on several. A cluster member that refuses to
boot is restored by wiping its data directory, so it rejoins from a peer. See
cluster.
Retention and reclaim
Retention is decided once, on the leader, and executed on every node, like every other change.
- Decide. A scanner on the leader walks every partition at
QUEEN_RAFT_RETENTION_SCAN_PER_S(50,000), at most one round perRETENTION_INTERVAL(5000 ms). The planner checks each proposal against the current state and logs two watermarks per partition:log_start, below which the payloads are gone, andtxns_start, below which the hash lists are gone. Hash lists are kept for the longest of the queue’s dedup window, its completed retention andQUEEN_RAFT_TXN_WINDOW_MIN_S(900 s), so dedup outlives the payload. A partition with no write forPARTITION_CLEANUP_DAYS(30), no retained message, no dead letter and no live lease is deleted. - Reclaim. Each node frees its own files below
txns_start, one sealed file at a time. A file with no live record is unlinked. A partly live file is rewritten with only its live records; the copy is made without holding the log’s lock (a sealed file is immutable) and swapped in under it, so appends do not wait for it. A per-queue log file is rewritten on its first dead message, a shared log file once half its message bytes are dead;QUEEN_QLOG_COMPACT_MIN_DEAD_PCToverrides both.
Reclaim never touches the active file, nor a file holding a record above the last store checkpoint
or, in a cluster, above the raft purge point. One pass examines at most 256 files, spends at most
200 ms and rewrites at most one file. An active file that holds data older than
QUEEN_QLOG_SEAL_AGE_S (600) is sealed so that reclaim can reach it.
Limits
Disk is freed by whole files: a partition’s bytes go when its file dies or is rewritten, and a
watermark that moves frees nothing by itself. RAM holds the whole store, since partitions,
cursors, KV keys and request ids are loaded at boot on every node, so memory grows with them.
Replication protects against losing a minority of disks; losing a majority for good can lose
acknowledged writes. A data directory does not move between replicators, and QUEEN_QLOG_SHARDS
and QUEEN_QLOG_LANES take effect only on a new directory.