Skip to content

Design Replica Backfill

Kadyapam edited this page Oct 10, 2026 · 2 revisions

Design: replica backfill — attaching a replica does not replicate what already exists

v0.7.0 (#400, PR #401).

The defect

Uploads are enqueued at exactly one site: L0Engine::register_and_upload, on seal. There is no other upload_tx.send in the engine.

That is correct for a store whose replica set never changes, and wrong the moment one is added. Attaching a substrate to an existing store replicates only parts sealed after the attach; every earlier part keeps its old replica count permanently. Nothing retries, because nothing considers it outstanding — the engine thinks it is finished.

Why it read as healthy

Measured on the prod server's embedded store, 2026-10-10, read back from the remote's own manifest rather than from the engine's view:

parts listed in the remote manifest 40
replica_count == 1 (local replica only) 39
replica_count == 2 (local + GCS) 1 — sealed after the attach

Every aggregate durability signal read healthy throughout:

replica_set_size 2
survives_node_loss 1
replica_domains_distinct 1
replica_single_point_of_failure 0

Those are computed from self.replicas — the declared replica set — not from the parts. survives_node_loss is a statement about configuration and it was being read as a statement about data. The only part-level signal was parts_under_replicated holding at a non-zero constant, which is indistinguishable from an upload backlog draining slowly. There was no signal anywhere for "these can never improve".

⚠⚠ The resulting state is worse than an empty remote — but measurably, not vaguely. write_manifest_to_all ships the manifest on every upload, so the bucket holds a durable manifest naming all 40 parts and the bytes of one. Measured against a GCS emulator in exactly this shape:

opened=true
replay_err="part shard-0-seq-...-000008: no reachable replica (1 listed)"
recovered=0 of 32

The open succeeds, because the manifest is complete — that is why the bucket reads as a usable store. The replay then fails loudly. So the remote is detectably broken rather than quietly wrong: a recovery attempt errors out, it does not return a short log and call it history.

⭐ That distinction is worth keeping straight. The dangerous outcome would be a silent partial recovery, and the engine does not do that — replay_all refuses the whole read rather than skipping an unreachable part. The guard in noetl/server tests/gcs_substrate_roundtrip.rs pins both halves, because an assertion that only said "recovered 0 records" would read as "the remote is an empty log", which is a different and unestablished claim.

The fix

pub fn backfill_under_replicated(&mut self) -> Result<usize>

Enqueues every manifest part whose replica_count() is below self.replicas.len() onto the existing background uploader, and returns the count. Pair with flush_and_wait_uploads() to wait. Returns 0 for a single-replica set — one replica makes no spreading claim.

Idempotent. put_if_absent returns Ok(false) for an already-present object and replicate_bytes records a ReplicaLocation for that case as well as for a fresh write. So a re-drive re-states the copies that exist rather than dropping them from the manifest — which matters because the uploader assigns p.replicas = locations outright.

UploadSource — the trap underneath

Manifest::durable_view stores local_path: None, and the durable view is what load_durable_manifest seeds the in-RAM manifest from at open. So after any restart every recovered part has no local file. A repair path that only knew how to fs::read(local_path) would have compiled, run, reported success and copied nothing — the original defect repeated one level down.

Hence the job carries a source:

enum UploadSource {
    LocalFile(String),                        // sealed in this process
    Replica { replica: String, key: String }, // read it back out of a holder
}

Separate counters, and why

backfill_uploads / backfill_upload_bytes, not uploads / upload_bytes / upload_lag_micros_total.

A backfilled part has no seal→durable interval — its seal happened in an earlier process. UploadJob.sealed_at is therefore Option<Instant> and None for a backfill. Folding these into upload_lag_micros_total would divide the real lag total by a larger denominator and quietly deflate the reported mean seal→durable latency: a fabricated sample, in a statistic a capacity decision reads.

record_replicated_lag needed no change and is worth stating as a guard: UnreplicatedTracker has no entry for a part sealed in a previous process, so on_upload_done returns None and no histogram observation happens. That is the correct behaviour by construction, not by accident — do not "fix" it by supplying a default.

The guards, and what they exist to prevent

crates/ehdb-l0/tests/replica_backfill.rs:

  1. attaching_a_replica_leaves_every_existing_part_short_forever — characterises the defect. Attach, then flush_and_wait_uploads() (the strongest thing a caller can do short of a backfill): all three parts are still single-copy and uploads == 0. Nothing is enqueued anywhere, so it cannot improve on its own. ⚠ Do not delete this as redundant; it is the statement that the seal-only enqueue is insufficient.
  2. backfill_replicates_history_and_the_new_replica_alone_can_serve_it — the fix, plus the assertions that uploads and upload_lag_micros_total stay 0 (no fabricated lag).
  3. a_single_replica_set_backfills_nothing — the call must not invent work.

RED controls: three mutants, two distinct kill sites. A no-op backfill and a LocalFile-only source both die at the enqueue assertion; writing the manifest to only the first replica dies at the cold-load assertion — which is what proves the headline clause is load-bearing rather than decorative.

Honest framing

A backfilled sealed part is durable across node loss. The unsealed tail stays RF=1 and replication is asynchronous, so recovery is a bounded-loss consistent prefix of sealed history — not "no data loss". Closing the tail is separate work (#394).

Consumed where

noetl/server calls this at engine open behind NOETL_EHDB_REPLICA_BACKFILL (default off — it is real egress on existing data). With a replica bucket armed and the flag off, the server logs that the remote is not a recoverable copy, so the healthy-looking gauges are not the only thing an operator sees. See that repo's deployment-specification page.

Related

Clone this wiki locally