Skip to content

Storage

What a node keeps on disk: queue log files that are the write-ahead log, a store held in RAM and checkpointed to LMDB, the recovery rules, the damage a node refuses to boot on, and reclaim.

Updated View as Markdown

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.

How an entry that touches two queues reaches the disk. The writer puts the payloads into each queue's own log, the whole entry record into the lowest of the touched logs and a 61-byte stub naming it into the other, then fsyncs each touched log once. Only then does the entry go to apply, which updates the store in RAM. Once a second a checkpoint writes the changed store rows to store/data.mdb, off the answer path.one entryorders and receiptsqlog/q12/payloads + entry recordqlog/q31/payloads + 61-byte stubstore, in RAMcursors, KV, timersstore/data.mdbcheckpoint, once a secondwrite, fsync oncewrite, fsync oncethen applyoff the answer path
A message is written once, in its queue's log, and that write is the one that makes it durable. The store checkpoint only bounds what a restart replays; no answer waits for it. Source: 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

  1. The store reopens exactly at its last checkpoint.
  2. Each queue log cuts a torn tail from its active file.
  3. 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.
  4. Apply skips any entry at or below the store’s applied index, so replaying twice changes nothing.
Recovery on one node. The store reopens at its checkpoint, with entries up to 102 applied. The queue logs hold entries 103 and 104 complete, entry 105 with one of its copies missing, and entry 106 after it. Recovery replays 103 and 104 and cuts 105 and 106 from every log: the fsyncs of that group did not all land, so it was never acknowledged.101102103104105106replayed the gapless runcheckpoint the store's applied index
Replay runs from the checkpoint through the gapless run of complete entries. The first incomplete entry and everything after it are cut from every log.

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.

  1. Decide. A scanner on the leader walks every partition at QUEEN_RAFT_RETENTION_SCAN_PER_S (50,000), at most one round per RETENTION_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, and txns_start, below which the hash lists are gone. Hash lists are kept for the longest of the queue’s dedup window, its completed retention and QUEEN_RAFT_TXN_WINDOW_MIN_S (900 s), so dedup outlives the payload. A partition with no write for PARTITION_CLEANUP_DAYS (30), no retained message, no dead letter and no live lease is deleted.
  2. 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_PCT overrides 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.

Navigation

Type to search…

↑↓ navigate↵ selectEsc close