Replication
Crewlet’s own tracker and knowledge base are replicated state machines (the embeddings and each node’s usage history ride the same framework as compacted domains). One ordered stream per domain is the write-ahead log; every node applies it into its own SQL database; the checkpoint commits in the same transaction as the rows it covers. That last clause is the whole design: a node’s position and its rows can never disagree, because they are written together or not at all.
This page is for operators. It covers what a write means, what an acknowledgement promises, how far behind a node can be, and what the design does not promise.
Three numbers that are not the same number
Section titled “Three numbers that are not the same number”Confusing these is the single most common misreading of this system, so they are stated first.
| Number | Where it lives | What it is for |
|---|---|---|
| the arbitration anchor | statelog_anchor | The position of the last record this node consumed on a subject, whatever that record then did. It is what the broker arbitrates a new write against. |
version | a column on every object row | The object’s accepted state. It is what a caller’s if_match compares against. |
| the barrier comparison | computed | MAX(version, scoped_through) — what a read barrier compares, because a record can change an object without being about it. |
The anchor and the version are equal only while every accepted record produces rows, and an apply gate is by definition the rule that makes them differ. An evicted node’s append is accepted by the broker at sequence 102 and dropped by every applier: the object’s version stays at 100 on every node while the broker’s last message on that subject is 102. A writer that formed its expectation from the version would be refused, re-read 100, and burn its whole round budget — for at least the trim’s age floor, and unbounded above while any retention term blocks.
scoped_through is the third. A record bumps version only on its own
subject; every other row it writes carries scoped_through instead. A barrier
that compared only version would let a read past a record that had already
changed the row it was about to return.
A write has three outcomes
Section titled “A write has three outcomes”Not two. applied, pending, unknown — and a caller does something
different with each.
applied— the record is durable at its position and this node has applied it. Read it back and you will see it.pending— the record is durable at its position and this node has not applied it yet. The write succeeded. Do not retry it; retrying publishes a second record. What is not yet true is that you can read it back here.unknown— nothing can be established about the record. It may be on the log and it may not. This is the only outcome a retry is correct for, and the retry carries the same operation id so the ledger collapses a duplicate.
A lost acknowledgement is resolved from the log before it is ever called
unknown: the write reads the newest record on its subject, and a record that
carries its own operation id is its answer — applied where this node has
applied it and no gate dropped it, even where this node’s ledger has lost the
operation’s row, and refused under the gate that dropped it otherwise — the
domain’s own, or one of the log’s: a rule a reanchor placed on the checkpoint
(abandoned, overtaken). A gate that drops a record by who wrote it (an
eviction) is asked about the
node the record names, which is not always the one asking: an append the broker
collapses onto another node’s copy of the same operation is answered by that
copy.
A gesture made of several records in order — a cross-project move, a merge
of duplicates, a dependency with its mirror, a subtree’s removal or restore —
treats a step answered unknown as the end of the walk, not as a step that
landed. Nothing after it is written: no descendant follows a root whose own
move or removal is unknown, a mid-move or mid-merge mark is not taken down, and
a dependency’s mirror is not written over an authored edge nobody can vouch for
(a mirror whose own outcome is unknown is reported one-sided, the state the
tracker duty repairs). The gesture fails naming the operation id, and running
it again under that id answers each step that landed from the ledger and
finishes the rest. Every step is named by the task it writes rather than by its
place in the walk, because a re-run reads its list afresh — a restore’s is what
is still in the trash — and a step named by position would carry another task’s
operation id.
What “under that id” means depends on who is asking. A seat derives its
ids from its turn, the call’s arguments and how many different calls to the
same tool it made first, so it runs a gesture again by repeating the call with
exactly the same arguments — before calling that tool with anything else. The
operator’s MCP has no turn, so every write’s answer carries its op_id and
a call that brings it back is that operation again
(/operator/mcp)
— and only that one: the id names the call’s tool and a digest of its
arguments, and brought back with any other it is refused before anything is
written, rather than half-answered from the first call.
The purge and node-gate routes take theirs as ?op_id=. A caller that sends
neither starts a new operation, which finishes nothing: a create repeated that
way files a second item. A create whose item landed and whose dependencies did
not says so — the item is named, and its dependencies are finished on it with
update_work_item rather than by filing it again.
What a retry is judged by: the instant its operation was minted
Section titled “What a retry is judged by: the instant its operation was minted”Every node keeps an operation ledger — which operations its own applier
has applied, and where — and it is the ledger, not the broker, that makes a
retry safe. The broker also collapses a repeated operation id, but only inside
its two-minute duplicate window, and a turn re-run after a crash or a caller
repeating an unknown routinely comes later than that.
A retry of an operation the ledger records is answered, not decided again:
it comes back applied, at the position the first copy landed, however long
afterwards it arrives. Deciding it again would mean deciding against rows that
already hold it — a retried purge would find its task gone and refuse, an
update conditioned on a version would find that version moved by its own first
copy and refuse as stale. A retried create comes back as the task the first
copy filed, under the key it took — the new item’s id is derived from the
operation, so a turn re-run after a crash that makes the same call files
one work item, not two. Only a create that took its key number and never filed
its task is finished on a fresh number, which leaves a gap in the project’s
numbering (ENG-8 is skipped) rather than two items sharing a key.
That holds for a re-run’s calls only where they are the same calls. A turn’s operation ids are derived from its unit of work, the verb, the object, the call’s arguments and how many different calls to the same tool the run made before it (see the turn engine), and the model behind a re-run is sampled again: a title or a comment it words differently is a different operation, so that re-run files a second item or posts a second comment. The engine cannot tell a reworded call from a new one, and a crashed run leaves no record of its calls to compare against.
The same holds for a write whose acknowledgement was lost and which the ledger then answers: unless the broker acknowledged this very attempt at the position the ledger names, the record there may be an earlier copy’s — decided on another node, against other rows — so the write is answered as that copy, and anything the attempt computed in its own decision is not used. A key is then read from the item’s row, or minted again on a fresh number. A retried create on a node that was behind used to be answered as its own stale decision, and filed its item under the number another item already held.
An operation id names one write, and a ledger row answers only for a write
to the object it landed on. The same id sent with a write to a different object
is refused op_reused, naming where it landed — whether the ledger finds
it before anything is decided or the broker’s duplicate window collapses the
second append onto the first record — rather than answered applied at a
position on an object the second write never touched. A different write takes
a fresh id.
The ledger travels inside a snapshot. A node that adopted one holds a row for every operation its donor applied, so a retry of one is answered from it exactly as the donor would have answered it, and an operation neither holds never applied — so a turn woken by work that was queued before the join files its writes like any other, however old the work is.
What the ledger can do is lose rows, and it keeps a watermark saying how far
back it may have. Each node deletes its rows thirty days after applying them,
and records how far back each pass deleted. An operation minted before the
watermark whose row is gone is answered unknown — before anything is
published — rather than decided a second time, so a retry a month on is never
applied twice. It is answered unknown rather than refused, too: a retry’s
rows already hold what its first copy did, which is exactly what makes a
re-run refuse (a move finds its task already in the target project, a create
finds its object already there), so on such an operation “no” would be the
ledger’s silence read as an answer. Only a refusal about the node itself — it
holds a record it cannot read, or the object is deleted for good — stands
either way. A create can say more than unknown, because its item’s id is
derived from the operation: one the ledger cannot vouch for is answered from
that item’s own row — the item the first copy filed, under the key it took,
where this node holds it, and unknown with no number minted where it does
not, since the first copy may be on the log beyond what this node has applied.
A gesture that writes several records in order and meets such a step stops
there and says so: re-running it under the same id on this node stops at
the same step every time, because the row that step needs is the one the loss
took, and it is a node whose ledger lost nothing that far back that can finish
it. The thirty days are sized for the slowest real retrier, a seat
that only runs on a schedule carrying an operation id across a long weekend; a
write that is retried later than that is told unknown, and the operation’s
own record, if it landed, is on the log. The watermark travels with the ledger,
so a node that adopts inherits its donor’s.
What decides it is the instant the operation was minted, and that instant
travels inside the operation id itself: every id the engine mints is a
time-ordered one whose leading bits are its mint time, followed by a name
saying what the operation is (01a0…-7…-….update-<task>). A retry reuses the
id, so it reuses the instant — including a turn re-run on another node, whose
writes derive their ids from the unit of work and the instant it began rather
than from the run. The exception is a unit of work older than the ledger — a
trigger dispatched, or a turn resumed from a coding run, more than twenty-nine
days after the work began: an id minted at that start would be one no node’s
ledger can vouch for any more, so that attempt mints at its own instant and
records it in the coordination store, and a later attempt at the same work
inherits it while it is recent enough to (see
Turn Engine). An
operation id the engine did not mint carries no instant
and is read as older than everything the ledger ever lost: on a node whose
ledger ever lost a row, a write under it answers unknown unless that node’s
ledger holds its row — which is why the routes that accept an operation id
from outside, a purge’s and a node gate’s retry, refuse one the engine did not
mint.
The instant is compared with the watermark, which another node’s clock may have set — the sweeping node’s, which a joining node inherits with its donor’s ledger — so the fleet’s clocks are assumed to agree to within the margin the sweep leaves, its thirty days. It is the same kind of assumption the trim’s age term rests on. Keep them synchronised.
pending is the outcome an ordinary busy fleet produces most often under load:
the applier is 16 seconds into a bulk apply and a small write’s five-second wait
for its own record expires. Five seconds is the applier’s own stall grace
divided by twelve — many times an ordinary batch commit and its linger, and
short enough that a caller holding a request open learns “durable, unresolved
here” rather than waiting out a node that has stopped applying. The record is
safe. A caller that treated pending as a failure would double every write it
made during a burst.
applied is not permanent
Section titled “applied is not permanent”Two gates can drop a record after a node has applied it, and both are deliberate:
- The eviction gate. A record written by a node the fleet evicted before the record’s own position is dropped everywhere. A node that applied it before learning it was evicted will drop it on replay. Every record names the node that published it and the generation it was decided in; the write path is handed both and refuses to append a record that does not carry them, because a record naming nobody is one this gate can never drop.
- The deletion gate. A record about a task a purge destroyed — a turn’s record of what it spent on the task included — or about a purged page applies nowhere, for ever. This is what stops a redelivery months later resurrecting rows an operator deliberately removed.
Neither is a bug being worked around: they are what make an eviction and a purge mean something on a system where the log outlives the decision.
Every node that can write is one the trim waits for
Section titled “Every node that can write is one the trim waits for”No check on a write decides whether its node may make it, because every node
that can is already one the trim waits for. Only a data node runs the logs and
writes to them — a node without the data role holds no estate and reaches it
through a data node
(how a request reaches the estate)
— and every live data node is in every log’s
counted set from the moment it starts, at
position zero until its first heartbeat reports one. So the trim never removes
a record that a node which may still write has not applied, and the one way
out of that set while a node is still running — an eviction — is the eviction
gate’s.
Who a refusal names
Section titled “Who a refusal names”A write refused evicted that names a position is on the log, and it holds
its operation id for the log’s duplicate window (two minutes for every
shipped log): the broker collapses the same id, sent again inside it by any
node, onto that record, and the answer is evicted again — the resolution
judges the record by the node that wrote it, never by the node asking. So
another data node takes the write under a fresh operation id, or under the same
one once the window has passed; neither can apply twice, because the record in
the way applies nowhere. An evicted refusal with no position was made before
anything was appended, and another node takes the write under the same id at
once.
When the record in the way is another node’s — the asking node’s append
was collapsed onto a copy that node wrote — the refusal names that node as the
copy’s writer, because the reason is then its standing and not the asking
node’s. The asking node passed its own checks before it appended, so it is the
one that takes the write, once the window has passed; a node gesture refused
this way says so in its hint and offers the same operation id again rather
than another node. That holds for every gate that drops a record for what its
writer was or did — evicted, abandoned and overtaken — and for none
other: a copy dropped because its task or page
was purged is refused deleted and names no writer, because the purge refuses
the asking node’s own write exactly as it did the copy, now and on every retry.
A record whose scope meets a deferred scope is deferred too
Section titled “A record whose scope meets a deferred scope is deferred too”A record this build cannot decode is retained, not dropped: the bytes are the only copy, and a rolling upgrade puts records on the wire the older half has never heard of. What follows is the rule that keeps that safe:
If a record’s declared scope meets a deferred record’s scope, it is deferred too.
So an object’s rows on any node are always a prefix of that object’s own record sequence. There is no state in which record 5 has applied and record 3 has not. That is what lets a reader say “this node is behind” rather than “this node has a hole in the middle of this object’s history”.
What an acknowledgement promises
Section titled “What an acknowledgement promises”stream.sync decides it, and the default is the strong value.
| Setting | An acknowledged write has reached |
|---|---|
always (default) | the disk, on every replica that acknowledged. |
a duration, e.g. 30s | the replicas’ memory, and their disks within that window. |
The five failure classes, and which the design survives at each topology:
| Failure | Single node | Clustered, one host | Clustered across hosts | External NATS |
|---|---|---|---|---|
| F1 the process is killed | survived | survived | survived | survived |
| F2 the host loses power | survived only with sync: always | survived for any one host | survived for any one host | whatever that cluster promises |
| F3 the disk is lost | not survived | survived | survived | whatever that cluster promises |
| F4 the filesystem corrupts | not survived | survived for any one host’s disk | survived for any one host’s disk | whatever that cluster promises |
| F5 the site is lost | not survived | not survived | only if the hosts are in different sites, which the engine does not know and does not claim | ask your NATS operator |
On stream.type: nats the engine refuses stream.sync rather than
warning: it is a server option on a server the engine does not run, and a field
that named a mechanism it cannot reach would let an operator hold a durability
belief nothing delivers. Ask the NATS operator for sync_interval: always
instead; crewlet validate says it cannot check it.
A declined stream.sync is refused on a leaf (stream.leaf.urls) for the
same reason: a leaf’s broker runs no JetStream and has no file store at all —
the members it joins decide what an acknowledged write has reached, in their
own stream.sync. A satellite’s durability is the members’ row of the table
above, never its own.
Above the floor and below it — two different regimes
Section titled “Above the floor and below it — two different regimes”This is the sentence the backup schedule hangs on.
Above the trim floor, the log holds every record on R replicas and every node holds the applied rows. Losing a node loses nothing.
Below the trim floor the log may already hold nothing — the trim has licensed deleting it — and each node’s own database file is the only copy of that history — N of them, independent, none replicated. Losing history below the floor takes all N disks, and it is covered only by the backup gate.
That is why the trim refuses to advance past a floor no backup has reached. The backup schedule is a correctness input, not hygiene. See Retention.
A node resumes from its rows, never from its reader
Section titled “A node resumes from its rows, never from its reader”Each node reads each log through a durable consumer of its own on the broker, named after the node. That consumer’s position is a second, weaker number than the checkpoint: it moves after the commit, and the broker keeps it when the node’s database does not. The checkpoint is what the node resumes from — but resuming can only drop what arrives below it. Nothing on the node’s side can make the broker hand over again a record it has already delivered and been told was applied.
So at every boot the node compares the two, and a consumer that disagrees with
the rows is deleted and created again at the checkpoint. Each rebuild is logged
as jetstream_domain_consumer_rebuilt, naming the drift:
drift | What it means | What it would have cost |
|---|---|---|
acknowledged_past_checkpoint (a WARN) | The broker was told records past the checkpoint were applied, so this node’s replicated database is older than its reader. The database was deleted, or restored from a backup. | The node never hydrates, and every write on the missing objects refuses behind. On the vector log, where nothing checks for gaps, the skipped records’ rows are just missing. |
delivered_past_checkpoint | Records past the checkpoint were handed to a process that has since stopped. | The node stays behind until the 30-second ack window returns them. |
in_flight_for_a_gone_reader | Deliveries at or below the checkpoint are held for a process that has since stopped. | Each one holds a slot of the in-flight ceiling for 30 seconds. When they fill it, nothing new arrives. |
behind_checkpoint | The rows moved without the reader, which is what adopting a snapshot at boot does. | Every record in between is delivered only to be dropped. |
A clean restart logs none of these: a node that applied and acknowledged everything it was handed keeps its consumer as it is.
A node below the trim floor is not ahead of its rows either, although its
reader reports records past the checkpoint as read: the broker moves a reader
over the records the trim removed, to the first one the log still holds. The
node reads the log’s first sequence to tell the two apart. If it cannot read
it, it judges the raw positions and rebuilds, and the
acknowledged_past_checkpoint warning then says the log’s first sequence was
unreadable, so it may be either case.
A rebuild is a delete and a create on the broker, and on a fleet whose broker is not answering, either one can fail. What happens then depends on the drift:
behind_checkpointandin_flight_for_a_gone_reader: the node asks the broker which reader it now holds, keeps that one, and logsjetstream_domain_consumer_rebuild_failedas aWARN. Neither drift has handed over a record past the checkpoint, so keeping the reader costs what the table says and loses nothing. The next boot rebuilds it.acknowledged_past_checkpointanddelivered_past_checkpoint: the boot fails and names the domain. A node left on that reader could wait for ever.- Any drift where the broker holds no reader after the failure, or cannot say which it holds: the boot fails too. A reader the node guessed at could be one that no longer exists, and every read through it would fail while the node reported itself up.
A rebuild touches only this node. No peer reads through its consumer, and no
retention term reads it either: the trim works from the positions each node
publishes from its own rows. What the rebuild makes safe is losing the
replicated database while the broker keeps its estate. A node started on an
empty crewlet-replicated.db replays each log from the broker, or adopts a
peer’s snapshot where the trim has already removed the start of one. It no
longer waits on records its old reader already acknowledged.
Replication lag is two positions
Section titled “Replication lag is two positions”Not one number. Every node publishes both:
seq— the record it has consumed: its checkpoint, the contiguous prefix it has moved over.applied_through— the prefix it has actually applied, which is lower whenever a record was retained rather than applied.
Beside them rides stream_created_at, the creation instant of the stream
the node’s rows are keyed to. A generation cannot say which stream a sequence
is on — a stream deleted and remade keeps the generation and counts from 1
again — so this is what lets a reanchor
tell a peer further along the lost stream from one that came up on the new one.
And checkpoint_stored_at names the record the checkpoint stands on — the
broker’s instant for it — which a sequence alone cannot: after a broker
restored from an older copy is written past a node’s rows, the log holds
another record at the same sequence. It is what lets a reanchor tell a peer
whose history is the log’s from one holding history the log lost.
Folding them into one would make a node that is applying nothing while its
position advances look identical to one that is fully caught up. crewlet retention status prints seq and applied_through per node and per domain,
beside the generation the position is in (GEN) — a position from a
generation the log has left prints left gen N in place of a lag, since its
sequence compares with nothing the log holds — and marks a node that reported
log_diverged. The two instants are on the retention answer’s JSON
(stream_created_at, checkpoint_stored_at) rather than in the table.
Lag does not move a node’s seats, at any size. A node that is behind keeps every seat it holds and claims no new ones until it is level — see a copy that is behind, and a copy that is wrong for the six states that take a copy out of service, none of which is a distance — and none of which moves a seat either: the node’s seats read the estate from the other data nodes. What lag does bound is how fresh an answer a read can ask for (Read Consistency).
The bulk-apply degradation, priced
Section titled “The bulk-apply degradation, priced”One applier per domain, one goroutine, one totally ordered log. The measured drain is at least 2 000 rows/s and the steady-state load of the reference company is 0.244 rows/s — 0.012 %. What binds is not the average; it is the burst.
That floor is deliberately conservative, and it is the number every figure below is derived from. An applier writes a record’s child rows — tags, watchers, relations, dependency mirrors, the inverted index’s postings — as multi-row inserts chunked to the engine’s probed bind-parameter limit, and to at most 1 000 rows a statement, rather than one statement per row. The pinned driver accepts more parameters than the probe’s own 32 766 ceiling, so the 1 000-row cap is what binds: a seven-column row and a three-column row both batch 1 000 to a statement, and an 8 000-row apply is 8 statements rather than 8 000. Measured unloaded that shape drains about four times faster than one statement per row; measured on a loaded CI runner under the race detector, closer to 1.5 times. The floor above is the loaded, contended figure, so a fleet sized against it has margin rather than a number it has to hope for.
One maximal bulk update is about 31 800 rows, which is about 16 seconds of applier occupancy on every peer. For those 16 seconds two things degrade fleet-wide without tripping any alarm:
- an unrelated single-object write’s five-second wait expires and returns
pendingwith its position and its lag — correct, and designed for; - a barriered read answers
behindwith a computedretry_after_secondsof about 16.
Both are the system working. crewlet retention status publishes the longest
apply transaction and the longest batch actually observed, so this is a
measured property of your fleet rather than a surprise.
The occupancy is per domain, not per node. Every domain’s applier writes the same replicated database, and a node hands that database’s write lock to its writers in the order they asked for it. A bulk update commits one transaction at a time — at most 4 000 rows, about two seconds at the measured drain — and a waiting writer takes the lock as soon as the one in front of it commits. So a bulk update in the work tracker delays the knowledge base’s apply by the transactions already queued ahead of it, at most one per other writer on that node, never by the whole 16 seconds.
A writer that does not reach the front within store.busy_timeout_seconds
fails retryably and rejoins the line, logged as store_tx_retry naming the
knob. With three domains applying and a bulk update in flight (the usage
domain’s transactions are a node-day’s handful of rows, and weigh nothing
here outside a replay), the default
five seconds is close to the three transactions a fourth writer can
legitimately wait behind — so that log line on a node doing bulk work is the
signal to raise it rather than a fault. Raise the knob before you widen
anything else: a longer per-waiter bound lets one stuck holder block the line
for longer, and this way the line still drains in order.
No apply transaction is ever aborted by a commit elsewhere in the database,
and none is ever re-run because of one. Every write transaction takes the
lock when it begins rather than at its first write, so there is no window
between an applier’s read and its write for another commit to land in. The
crewlet.statelog.apply.tx.aborts counter in
Metrics reads zero on a healthy node for that
reason, and a non-zero count is the retry budget being spent on something
else.
What a rolling upgrade blocks
Section titled “What a rolling upgrade blocks”A record at a version this build cannot decode is retained, and everything whose scope meets it is retained too. That is bounded rather than open-ended: the deferral grace is 30 minutes, past which the alarm fires and names the position and the scope.
So a rolling upgrade should finish inside that window. An upgrade that stalls
half-done leaves the old nodes holding records they cannot apply and refusing
reads about the objects those records touched, and writes to them — with
deferred naming exactly what to do, which is finish the upgrade. What a record
touched includes where it sits: a record about a knowledge space’s settings
covers every page in that space, a page purge covers the whole space it
re-files the page’s children in, and a task purge covers every task it
rewrites — its dependents, the blockers it waits on, the tasks related to it
or referencing it, and its subtree — each by name, or by its project once
there are too many to name.
A record is written at the lowest version that can apply it whole, never at the newest the build knows, so what an older node holds back is exactly the objects whose shape changed — the next section lists them. Every barrier and every generation record stays at version 1 for good: an older node retaining those would hold a deferral for every linearizable read, or never make the transition a reanchor announced. So does every record that installs a gate, and for a harder reason: a node that deferred an eviction would go on applying everything the evicted node appends, so a gate record no build can decode stops that build’s applier instead, and its shape may only ever grow by addition.
The upgraded node applies what it retained at its next boot, before its applier consumes anything new: every retained record it can now read, in log order, each released in the transaction that applied it. A record whose scope met a retained one is applied after it, which is what makes the objects’ rows a prefix of their history again rather than a hole. A record still above the new build’s version stays retained, with everything it covers, until a build that reads it boots.
Which records an upgrade holds back
Section titled “Which records an upgrade holds back”Only the ones that need the newer build. A record is stamped with the lowest version that can read it: 1 when it carries nothing a later build added, and the version of the newest field it carries otherwise. So during a rolling upgrade an old node applies everything the new nodes write except the records that use something new, and it retains those whole rather than apply them with the new part dropped — which is what would leave its copy of that object different from its peers’ for good. An upgrade that adds no record field holds nothing back at all.
Today no record in any domain carries a field past the base format, so every record is written at version 1; the first field a later build adds raises exactly the records that carry it.
Values the engine computes are recomputed once
Section titled “Values the engine computes are recomputed once”Some columns are not copied out of any record but computed from the history a
node already holds: in the tracker, how often a task was reopened (reopens),
how much of a project’s open work has been started (task_counts.active), the
hand-off count on each history row, and when a task last changed
(updated_at). The apply maintains each one as it goes. When the rules that
compute them differ from the ones a node’s rows were derived under — a build
that adds such a column or changes how one is computed, a node that just
adopted a snapshot from a peer on a different build, or a checkpoint a reanchor
created — the node’s next boot recomputes them from the rows it holds, in the
same transaction that records which rules the rows now follow, before it
applies anything new. It happens once per change, on every node; the
statelog_rederived log line names the domain, the rule versions it moved
between and how many rows it wrote.
Two compacted domains: the embeddings and each node’s day
Section titled “Two compacted domains: the embeddings and each node’s day”Two domains do not keep a history at all. Their streams keep one message per subject — a keyed table rather than a log — so a node that joins replays the current value of each object and nothing before it, a gap is a coverage number rather than a fault, and neither gates a node’s seats or claims that two nodes hold identical rows.
| Domain | Stream | One subject per | Written by | Kept for |
|---|---|---|---|---|
| vectors | CREWLET_TRACKER_VECTORS | embedded source | the fleet’s one embedding duty | 90 days per message; the rows until the source is forgotten |
| usage | CREWLET_USAGE_LOG | (node, company day, seat or schedule) | every node, for its own days only | 181 days |
usage is what makes history fleet-wide. Each node derives its own day —
the spend by phase, worker, model and provider slot, the turns that ended and
how (failed, reviewed, first pass, sent back, a duration histogram), the pages
its seats read, the schedules it fired — from its own event log, every 15
seconds while the day moves, and publishes each object’s whole cumulative value.
Every node applies every node’s days, so any node answers for the fleet and a
node that leaves takes nothing with it: its days stay in every peer’s rows, and
a node that joins later replays them from the stream. Because the node is part
of the subject, no two nodes ever write one object and nothing is arbitrated;
an apply replaces the object’s rows.
A day is cut on the company’s clock (timezone). A node re-derives today and
yesterday on every tick, so the turn that ended a second before midnight
reaches the fleet, and on every start, so a day a node was down across still
arrives. History older than 181 days leaves in the same transaction that
applies the day that makes it old, on every node alike — there is no sweep.
The decision is ADR-0020.
The six capacity ceilings
Section titled “The six capacity ceilings”- The broker’s storage limit —
stream.store_max_bytes. Every ceiling below is a reservation checked against this one number, so it is the ceiling above the ceilings. Unset, the embedded broker takes three quarters of the free space onstream.store_dirwhen its JetStream comes up, measured once. Divide it when more than one engine shares a filesystem: free space bounds their sum, not each of them. - Each log’s byte ceiling:
stream.tracker_log_max_bytes,stream.tracker_vectors_max_bytes,stream.pages_log_max_bytesandstream.usage_log_max_bytes, sized together as below. A full log refuses appends rather than shedding old records; see Retention. On the tracker and pages logs ordinary writes are refused a sixteenth short of it, the rest being kept for gate records so that an eviction can still unpin a full log. - The trim floor — how far back the log can be replayed from, which is what bounds how long a node may be away.
- The store’s own size — every node is a full replica, so the corpus is held N times.
- Applier occupancy — the 16 seconds above.
- The snapshot repository — one artefact per node, sized in Retention.
How the byte ceilings are sized
Section titled “How the byte ceilings are sized”A byte ceiling is a reservation. The broker grants it in full when it creates the stream, before a single record is written, and refuses to create a stream whose ceiling it could not honour. So the four logs compete for one number, and a node sizes them together, once, when it creates their streams:
| Step | What happens |
|---|---|
| What the broker can grant | Read from the broker itself. An embedded broker’s limit is stream.store_max_bytes where you set one, and otherwise three quarters of the free space on the volume holding stream.store_dir, counting what its own streams already hold there; an external one’s is the NATS account’s JetStream limit. What counts against it is the ceilings already granted, not the bytes stored. |
| The logs’ share | Half of that, with the ceilings the logs’ own streams already hold counted as theirs, so a restart divides the same half the first boot did. Where the broker states no limit, or it cannot be read, the share is half of the stream volume’s free space instead, and nothing is added to it: a reservation never spends free space, so that figure already contains what the logs hold. The other half is for everything that reserves nothing: every mailbox, every coordination bucket and the snapshot a joining node reads. |
| Each log’s ask | Its Tier A field when you set one. Unset, the mutation log asks for a quarter of the stream volume’s free space (4..64 GiB); the knowledge base’s log for a quarter of the mutation log’s ceiling — the one you set, or the one it derived — because both hold a trailing window of records and the knowledge base’s grows at about a quarter of the rate (1..16 GiB from a volume alone); the vector changelog for the same quarter of the volume the mutation log derives, whatever the mutation log is set to, because its peak is the whole corpus rather than a window of it; and the usage log for a fixed 1 GiB — a count of node-days, which more disk does not grow. |
| The fit | A log whose stream already exists takes the ceiling it holds off the logs’ half first, whatever its field says now. A ceiling you set for a log being created comes off next, and is never scaled. The unset ones being created share what is left in proportion to what each asked for, none goes below 1 GiB, and none is created above what it would get if no log existed yet — the figure every later boot reports its stream against. |
A stream that already exists keeps its ceiling. Sizing decides what a
missing stream is created with and nothing else: a booting node never rewrites
a running stream’s configuration, and the broker never re-checks a reservation
it has already granted. A log created larger than today’s sizing would make it
boots as it is, and the node logs jetstream_stream_capacity_differs with both
numbers — its ceiling, and what this sizing would create it with if no log
existed yet. Its reservation was granted when it was made, so this is
harmless. To reclaim it (or to raise any log), use
crewlet retention set-capacity.
And it counts at that ceiling when another log is created beside it. A log
a new version adds, or one whose stream was deleted, is sized from what the
existing logs leave of the half, not from what they would ask for today — so a
log being created fits inside what the existing logs leave of their share,
past it only by the 1 GiB floor and a ceiling you set for it. What the existing
logs already hold is not reduced: when they hold more than the share (a ceiling
set and later unset, a log created while the volume had more room, one raised
with crewlet retention set-capacity), the logs reserve that much past it too,
and a log created beside them gets the floor. When a log is created below what
it would have had on an empty broker, the node says so with
statelog_ceiling_short_of_fit, naming each log it created short, the ceiling
it got (created) and the one it would have had (fit), and what each
existing log holds (held); once the node is up, crewlet retention set-capacity gives an existing log’s reservation back and raises the new one.
It is never created above that fit, even where the existing logs leave more:
every later boot reports its stream against that figure, and a log created
past it would be reported as a capacity difference on every one of them.
A ceiling the broker would not report is counted as absent. Before it
sizes anything a node reads every existing log’s ceiling from the broker — all
four at once, each asked again every two seconds for up to thirty while nobody
answers. A read still unanswered after that, or one the broker refused, counts
that log as one this boot would create. A log that does exist keeps the
ceiling it holds, but every log this boot creates is sized without knowing
what it holds, so it can come out larger or smaller than a boot that read it
would make, and keeps that ceiling. The node logs statelog_ceiling_unread
naming the stream; once it is up, crewlet retention set-capacity changes
either log.
A boot that still cannot reserve a log says why. When even the floors, or a ceiling you set, do not fit, the node refuses to boot with an error naming the stream, the bytes it needed, the bytes the broker had left and what sets that limit, and the Tier A field the ceiling came from:
engine: the broker refused to reserve the pages log's ceiling:CREWLET_PAGES_LOG needed 1073741824 bytes and the broker had 536870912 bytesleft to reserve (… of its …-byte limit already reserved), and that limit isstream.store_max_bytes where you set one, and otherwise three quarters of thefree space on the volume holding stream.store_dir (/var/lib/crewlet/stream),counting what the broker's streams already hold there.stream.pages_log_max_bytes is unset, so the ceiling was derived and scaledinto what the state logs that already exist leave of their share of thebroker, and it goes no lower than 1073741824 bytes. Give the broker more room;the state logs that already exist keep the ceilings they were created with,and no Tier A setting changes them: …The remedies are the ones it lists. Give the broker more room: raise
stream.store_max_bytes where you set one, or, where you did not, free space
on that volume (a first boot needs at least 5⅓ GiB free there, three quarters
of which is the four 1 GiB floors), which the broker measures again when the
node next starts. Or, when the refused log’s ceiling is
above the 1 GiB floor, set its field to a smaller ceiling, and the refusal says
so when that applies.
A log that already exists cannot be shrunk to make room from here. Its ceiling
changes only through crewlet retention set-capacity, which runs on a node
whose state logs are up, and every mode starts them, maintenance and seal
included: a node refused here cannot run it.
Reading a node’s state-log lines
Section titled “Reading a node’s state-log lines”Every line the state log writes about its own work carries
component=statelog, so filtering on it narrows a node’s log to its
replication:
| Line | Level | What it says |
|---|---|---|
statelog_apply_retrying | WARN | The applier hit a failure it retries in place. Written once, when the run of failures starts. |
statelog_apply_faulted | ERROR | The same failure has outlived the retry budget (30 seconds): this node’s rows have stopped moving, its reads refuse and it stops serving the estate — its seats read it from the other data nodes — until a retry succeeds. Written once per run of failures, when it crosses the budget — not on every retry. While it lasts, the node’s status and every refused read name the current error, and crewlet.statelog.apply.retries counts the attempts. |
statelog_apply_recovered | INFO | A retry succeeded and the run of failures is over, with how long it lasted (after) and the last error it saw. A failure after it starts a new run, written again from statelog_apply_retrying. |
statelog_applier_stopped | ERROR | The applier stopped for good — a gate this build cannot read, a hole that will not close, a recreated stream, a record written in a generation this node never entered — naming the stream, the position its rows froze at and why. Every read of that domain refuses from then on, and for the tracker or the knowledge base the node also stops serving the estate: the seats it holds stay and read it from the other data nodes, and a new one is admitted only on a data node whose copy is sound. What resumes it is a build that can read what this one could not, at its next boot; for a recreated stream, crewlet retention reanchor of that one stream, which resumes it in place with no restart; and for a log a peer re-anchored, the snapshot this node adopts on its own. Written once per stop. |
statelog_adopted | INFO | The node replaced its replicated database with a peer’s snapshot, naming the donor, the artefact’s sha256 (the donor’s statelog_snapshot_sent carries the same one) and when it was taken (taken_at), which is how old the history it installed is. |
statelog_record_gated | WARN | A durable record this node’s applier dropped: it applies on no node, and this line is its only witness. Each node’s applier writes it once per record, when the transaction that drops it commits. It names the domain, the record’s position and kind, the node that published it (writer, empty for a record that names none) and the gate that dropped it. The domain’s own gates: evicted — the writer was evicted below this position and not readmitted; deleted — the record is about a task (its turns’ records included) or a page a purge destroyed, whose marker holds every writer’s record on it for ever. The framework’s: abandoned for a record written in a generation a reanchor skipped because only an evicted peer held it (retention), and overtaken for one a node wrote in the old generation after a restored reanchor’s own record, before it learned of the move (retention). A record more than one gate holds is logged under the first the applier asks — the framework’s two, then the writer’s eviction, then the object’s marker — so a write refused over the same record can name another (statelog_write_gated names the one that holds it now). Counted by crewlet.statelog.records_gated under the same gate, which the records_gated alarm reads. |
statelog_write_gated | WARN | A write refused because what it appended applies nowhere, naming the domain, subject and op_id, the gate it was refused under, a position and a writer. When writer names a node other than the one whose log this is, the record at position is that node’s copy of the operation, which this write’s append was collapsed onto — not a record this node published — and under every gate but deleted the refusal is about that node and names it (above). When writer is this node, the record at position is its own, or another operation’s — the newest on the subject — at which the gate holds this node, so whatever it appended under the operation applies nowhere. The gate is the one that holds the record now, which can differ from the one the applier logged: once a task or page is purged, a record an eviction dropped on it is refused deleted. It is a refusal, counted under crewlet.statelog.publish.refusals by its reason, and not a second drop: the drop is the statelog_record_gated line each node’s applier writes once, which crewlet.statelog.records_gated and the records_gated alarm count. |
statelog_publish_unknown | WARN | A write could not tell whether its record landed. The operation id is in the line; retry under that id, never a fresh one. |
statelog_write_unvouched | WARN | A write was answered unknown rather than published or refused (a refusal field says what the decision refused), because its operation was minted (minted_at) before this node’s operation ledger may have lost rows to the ledger’s thirty-day sweep, and the ledger holds no row to say whether it already landed. Retrying on this node answers the same; a node whose ledger lost nothing that far back can answer it, and the operation’s own record — if it landed — is on the log. |
statelog_reanchor_started, statelog_reanchored | WARN | A generation transition of ONE domain. Both name the domain (domain), the one stream it moved (stream), the new generation, the stream’s live creation instant (stream_created_at), the case (case: recreated, followed from its first surviving record; restored, followed from its end; or abandoned, followed from this node’s own checkpoint with the records of the generation an evicted peer held void) and the new checkpoint (cursor); the start also names the instant the rows were keyed to before (keyed_to), this node’s checkpoint (position) and where the log ends (last_seq) and the generation its rows stood at (from_generation — every generation strictly between it and the new one is abandoned), and the completion gives the stream’s high-water mark before the reanchor (prev_last_seq_seen). A restored reanchor the operator ran with -discard names, on the start as discarding and on the completion as discarded, the sequence of the newest record written after the restore that it applied on no node (0 when it discarded none). No other domain’s checkpoint moves, and the domain’s applier resumes without a restart. |
The snapshotter, the donor and the adopter write under the same component. The
lines the engine writes around those loops — statelog_stream_recreated and
statelog_below_the_floor among them — carry component=engine, because the
component names the code that wrote a line rather than what it is about. Four
of them are about a peer’s reanchor: statelog_generation_passed (WARN) is
the heartbeat finding a log re-anchored past this node’s generation, naming
both (generation, fleet_generation), from which point the domain refuses;
statelog_behind_a_reanchor (WARN) is the join asking the fleet for a
snapshot at the new generation, naming the domains and the generations it asks
at; statelog_generation_holder_unread (WARN) is a beat that could not read
whether the node that opened that generation is evicted, which decides whether
the refusal names the adoption or a reanchor, and leaves it naming what it did;
and statelog_generations_unread (WARN) is a beat that read the
positions register but could not establish which generation the fleet is on —
the trim floors unread, or an eviction record unreadable. Such a beat judges
neither that nor the truncation below, and leaves both verdicts where they
were: the peer whose rows hold what a log lost is usually the peer that then
re-anchors past it, so judging the truncation alone would lift that fence on
the very beat that could not yet say the passed verdict replaces it. See
a node a peer re-anchored past.
Two are about a broker restored from an older copy:
statelog_log_diverged (ERROR) is the log holding, at this node’s
checkpoint, another record than the one it consumed there — the restored log
was written past this node’s rows — written once, at boot or by the heartbeat,
from which point the domain applies nothing past its checkpoint and refuses
until it is re-anchored
or its rows are replaced; and statelog_checkpoint_unverified (WARN) is a
heartbeat that could not read the record at the checkpoint to compare it, which
the applier asks again before it applies past it. And three are the same
restore seen from a node whose rows are the copy’s age:
statelog_log_truncated (ERROR) is the heartbeat finding a peer whose rows
hold records the log lost — past its end, or reporting log_diverged — naming
it (peer, peer_seq, last_seq, peer_diverged), from which point the
node’s writes of that domain refuse log_truncated while its reads go on;
statelog_log_truncated_cleared (WARN) is that ending, once the peer has
re-anchored, been rebuilt or been evicted; and statelog_truncation_unread
(WARN) is a beat that could not establish it either way, which leaves the
verdict where it was — as does a statelog_generations_unread beat, which does
not ask.
Three things CI cannot prove
Section titled “Three things CI cannot prove”Stated here because the alternative is implying a guarantee nobody measured:
- Durability under real power loss. The fsync counting and the fault-injection suites test the code path; a real machine losing power at a real disk’s worst moment is not reproducible in CI.
- Behaviour of an external NATS cluster. The engine does not run it and cannot assert its settings.
- Search capacity under production load. The scan benchmark has a concurrency axis, but the supported corpus at your concurrency is a projection until it runs on your hardware.
What the design does not promise
Section titled “What the design does not promise”N nodes applying one ordered log deterministically produce N identical copies including of an apply bug. A replica set protects against machine loss and against nothing else.
The residual is real and is not engineered away: a bug that corrupts silently and is noticed after the backup window has rolled costs history. That is why the trim refuses to advance past a floor no backup has reached — the backup gate is not ceremony — and why the restore test has a cadence rather than being a thing somebody did once.
See also
Section titled “See also”- Read consistency — the four levels and the twelve refusals.
- Retention — the trim, the six terms, snapshots and the join runbook.
- Backups & restore — the artefact, and what its interval actually means.
- Running a fleet — topologies, seat placement, rolling upgrades.
Part of Crewlet. Generated from crewlet/crewlet main at f665f5a. This is not the current version — see the latest docs.