Skip to content

Design Tail Replication

Kadyapam edited this page Oct 10, 2026 · 2 revisions

Design: replicating the unsealed tail (B3)

v0.8.0 (#394, PR #404), with the lock boundary corrected in v0.8.1 (PR #405).

The gap, measured rather than argued

A2/A3 made sealed history survive node loss. A part is local-only until it seals, so the loss window was the seal interval — seal_max_age (B2) bounds that window and cannot close it.

Read off the prod server while designing this:

NOETL_EHDB_SEAL_MAX_AGE_SECS               = 900
noetl_ehdb_oldest_unsealed_age_seconds     = 719   ← on ONE disk
noetl_ehdb_unreplicated_oldest_age_seconds = 0     ← uploads keeping up

Together those two gauges are the whole story: nothing is waiting to upload, and ~12 minutes of records exist only on the local device.

worst-case loss window
before B2 unbounded in time — a shard that goes quiet never seals
after B2 the seal interval — 900 s on prod
after B3 the replicator tick — 15 s default

⚠⚠ Not zero, and the tick is the number to quote. A record appended just after a tick is unreplicated until the next one. No remote ack gates the local append: this is not consensus and not a quorum. Saying "the tail is replicated" without naming the interval overstates it exactly the way "no data loss" would.

Shape

L0Engine::replicate_tail() copies records still held in an active writer to the replica set as write-once, per-batch objects. cold_load_replicated replays them.

Per-batch, not a re-upload of the active part. The part grows, so re-uploading it every tick is O(size) per tick and quadratic over the part's life: a 7 MiB part re-uploaded every 60 s for 7 h is ~1.5 GB of puts for ~7 MiB of data. A batch per tick writes ~the data volume, once.

No new read path. Recovered records go through the ordinary append path, so they land in a fresh active part and are served by the tail read that read_partition_after_limited already performs:

// The active (unsealed) hot buffer for this shard — the tail, so it is
// only worth reading when the limit has not already been reached.
if out.len() < limit {
    if let Some(writer) = self.writers.get(&shard) { ... }
}

That discovery is what shaped the design — B3 needed no reader at all.

Three invariants, and what each exists to prevent

A tail object is not a part

It lives under tail/<dataset>/shard-<N>/, never parts/, has no PartMeta, no sparse index, and never enters the manifest. If one did, plan_retention could drop it as though it were a part, and a read would try to resolve it through a sparse index it does not have.

The overlap is expected, not an error

The ai-meta#335 shape: a part that seals after its tail objects were written makes them redundant, so the same record legitimately sits in both. Dedup is by sort key against the durable manifest's max_sequence, then across objects. tail_objects_superseded counts the redundant ones — expected non-zero in steady state, so a permanent zero alongside non-zero tail_batches means the dedup is not being exercised rather than that it is working.

Appending recovered records in ascending sort-key order is required by the ascending-contract canary.

A failing tail put fails nothing

(#394 D3.) The record is already durable locally; failing an append because a replica is unreachable converts a durability feature into an availability regression — the same trade A2 made when it chose to open at RF=1 rather than refuse. The watermark is left in place so the next tick re-drives exactly those records, and tail_replication_failed counts it, because a replicator failing every put and an idle one both leave tail_batches flat.

⚠⚠ The lock boundary (v0.8.1)

replicate_tail does prepare → upload → commit in one call, so the remote put runs while the caller holds the engine lock. For an embedded engine that is a latency defect rather than a style one: in noetl/server the same lock is taken by shadow_append, which runs inside emit_events on the live write path, so every append colliding with a tick waits out the put — p50 77 ms, p99 234 ms measured against GCS in-region.

The uploader thread has always avoided this and start_uploader says why — "the substrate writes happen OUTSIDE the lock so a slow store never blocks appends/reads". So a latency-sensitive driver uses the split:

I/O safe under the lock
prepare_tail_batches() none — clone + frame ✅
upload_tail_batch() the put ❌ a free function; borrows nothing from the engine
commit_tail_batch() / fail_tail_batch() none — moves the watermark, or deliberately does not ✅

replicate_tail is kept and implemented on top of these, so a test or CLI needs no change and the same code is covered either way.

⚠ This was found by reading the new driver against the lock the append path takes — not by an incident. If you add another substrate-touching method the server drives on a timer, ask the same question first.

D2 — cleanup (v0.9.0), and why the refusal is the design

Owner-approved 2026-10-10: delete-after-sealed. Without cleanup, tail objects accumulate forever — ~1,440/day at a 60 s tick against prod's measured 314 appends/hour. Done carelessly, it deletes the only off-box copy of an event.

reclaim_superseded_tail_objects is a free function, so the listing, the sizing and the deletes all run with the engine lock released. It takes watermarks as data from L0Engine::contiguous_durable_watermarks, which reads the in-RAM manifest and is therefore safe under the lock.

⚠⚠ "Contiguous" is load-bearing

The watermark walks each shard's parts in ascending min_sequence and advances only while they are durable, stopping at the first non-durable part.

The obvious alternative — max(max_sequence) over durable parts — skips a local-only part in the middle, and a durable part beyond that gap would then authorise deleting a tail object whose records exist nowhere off-box. That is not a cleanup; it is the destruction of the only remote copy of those events. The conservative answer costs some retained objects; the other answer costs data.

Pinned by a_tail_object_is_retained_when_a_middle_part_is_still_local_only, which builds a durable / local-only / durable run by refusing one part's upload. Two mutants, both compiling, both caught by that test alone:

mutant result
naive max-over-durable watermark [(0, 24)] where [(0, 8)] is correct — it skipped the gap
delete regardless of watermark "NOTHING may be deleted while a middle part is local-only"

Anything not understood is kept

Unparseable keys, other datasets, other prefixes. An unparseable key is a thing this code does not understand, and the safe response to not understanding an object in an event-log bucket is to leave it there. Counted as unparsed rather than skipped silently.

Four series, and one of them is a gauge on purpose

tail_objects_reclaimed · tail_reclaim_bytes (sized before the delete, so the figure is measured rather than estimated) · tail_reclaim_failed · tail_objects_retained.

⚠ Retained is a gauge and is reported deliberately: a reclaimer deleting nothing and one with nothing to delete both leave reclaimed flat. Retained rising while reclaimed stays flat is a stalled seal or upload, not a quiet cleanup.

The consumer's default is ON

noetl/server defaults NOETL_EHDB_TAIL_CLEANUP to true — the opposite of its other flags — because omitting cleanup does not leave the system as it was, it leaves a bucket growing forever, while the operation itself is safe by construction. Proven against a real GCS emulator: the object is deleted from the bucket and the records are then recovered from GCS alone with the local store gone.

Both halves, again

Like seal_max_age, the flag is inert without a driver: tail_replication: true replicates nothing unless something calls replicate_tail (or the split) on a timer. noetl/server guards this with the_tail_replicator_has_both_halves and the_driver_does_not_hold_the_engine_lock_across_the_remote_put.

The guards

crates/ehdb-l0/tests/tail_replication.rs — 10 tests. The load-bearing one is the RED control: without_tail_replication_the_unsealed_records_are_lost proves the loss is real (16 of 21 recovered with the flag off). Without it every green result is unfalsifiable, because a cold load that silently found the records another way looks identical.

Then: the tail survives a destroyed local root and cold-loads from the replica alone (21 of 21); tail objects stay out of the manifest; an overlap yields each record exactly once; a refusing replica is counted and retried with the watermark unmoved; a decoy under another dataset's prefix stays invisible; prepare touches no substrate while upload against the same refusing replica fails.

Mutants, all compiling, with distinct kill patterns: a no-op replicator kills four tests but correctly leaves the RED control passing; unwiring recovery from cold_load kills only the read-side tests; disabling the dedup watermark kills only the overlap test; making prepare do the upload kills only the no-I/O test.

Six series

tail_batches · tail_records · tail_bytes · tail_replication_failed · tail_recovered_records · tail_objects_superseded.

Counted separately from the uploader's, for the same reason the backfill counters are: these have no seal→durable interval to contribute, and folding them in would divide a real lag total by a larger denominator.

Related

Clone this wiki locally