Skip to content

Consistency Invariants

Kadyapam edited this page Sep 16, 2026 · 3 revisions

Consistency invariants, per tier

What each tier actually guarantees — and, for each guarantee, what forces it to hold and how you would find out if it stopped. A guarantee with no enforcement and no detector is a description, not an invariant.

Per-tier invariants as of 2026-08-30; the transition section (§5) re-measured 2026-09-16. Live modes: eventlog primary, projection / kv / object shadow, NOETL_EHDB_VECTOR not set.

Summary

tier ordering ack durability can it serve? live mode
event log total order per shard (global_sequence) one local fsync; substrate async ✅ wired primary
projection fold order = event-log order derived; no independent ack ✅ wired shadow
KV fold order; clock-free derived ❌ no serve path shadow
object immutable parts; content-addressed write-once ❌ no serve path shadow

⚠ "Can it serve?" is SERVE_WIRED_TIERS = ["eventlog", "projection"]. KV and object accept primary as configuration and cannot serve as capability.


1. Event log

I-EL-1 — total order per shard. Every record has a global_sequence assigned by that shard's writer, ascending within the shard.

  • Forced by: a single writer per shard assigning keys under one engine lock (append_writer_assigned).
  • Detected by: out_of_order_appends — a counter that is 0 by construction on the writer-assigned path. Non-zero means some producer is appending out of order again. This exists because such a record "lands behind any follower cursor and is silently never delivered": the loss class was made observable rather than fixed away.

I-EL-2 — at most one writer per shard. ⚠⚠ Not currently enforced.

  • Forced by: nothing stronger than StatefulSet replicas: 1, which is an orchestration preference, not a mutual-exclusion primitive. A partitioned node shows as Terminating while its process keeps appending.
  • Detected by: now detectable, still not enforced. #330 put Invariant F in the storage contract in shadow: a write from an epoch below the shard's highest accepted epoch is counted in ehdb_fencing_stale_observed_total and still succeeds. ehdb_fencing_stale_refused_total stays 0 until enforcement is enabled, and the gap between the two is exactly what the flip would change.
  • Elected by: #331 — per-shard Lease election issuing a monotonic epoch, wired but not authoritative. Single-writer still rests on replicas: 1.
  • ⚠⚠ Ordering hazard: enforcing before the election issues real tokens is an outage, not a degradation — with no election every writer's epoch is 0. See the four-gate plan.
  • ⚠ Still true: this tier is already primary on prod, so all of the above is remediation of a live gap.

I-EL-3 — an acknowledged append is durable on one disk. Holds.

  • Forced by: an fsync before the ack. The prod writer runs FlushPolicy::CallerDriven and pays one fsync per batch itself (group commit); append_batch returns only after it.
  • ⚠ Not implied: durability beyond that one disk. See I-EL-4.

I-EL-4 — an acknowledged append reaches the durable substrate. ⚠⚠ Bounded by volume, not by time.

  • Forced by: a background uploader that runs only when a part seals, and should_seal() triggers on seal_max_bytes (8 MiB) or seal_max_records (1024) — there is no age trigger. An idle shard's records are never uploaded. The system is least durable where it is least active.
  • Detected by: ehdb_l0_unreplicated_age_seconds{shard} — the age of the oldest acknowledged record not yet durable, measured from the append (#328, live on the writer's /metrics since worker v5.125.0). With ehdb_l0_unreplicated_records{shard} and an ehdb_l0_replicated_lag_seconds histogram. Every shard is pinned to a row, so an idle shard reads 0 rather than vanishing, and ehdb_l0_durability_sample_ok states whether the scrape sampled at all.
  • ⚠⚠ Do not use instead: upload_lag_micros_total is measured from sealed_at, so records waiting in an unsealed part contribute zero — it reads healthy in exactly the failing case. And ehdb_feed_shard_lag / ehdb_feed_total_lag are consumer backlog, not replication lag, despite being the names an alert author reaches for first.
  • Bounded by: an age-based seal trigger (#329), merged and default-off. ⚠ The flag alone is inert — should_seal() is only consulted on append, and the shard it protects takes none; a timer must drive seal_aged_parts().
  • ⚠⚠ Worth knowing: in prod the substrate is the same volume, and #332 now gives that a name — FailureDomain, resolved from the device id rather than the path, with validate_replica_domains refusing a shared or nested root. ⚠ Prod as it stands would fail that check, so turning it on at open before fixing the layout is a startup outage by construction. NOETL_EHDB_TIER_SERVICE_DIR=/data/eventbus/ehdb-tier sits inside NOETL_EVENT_BUS_WRITER_DIR=/data/eventbus, one PVC. LocalFsSubstrate is the only substrate implementation. So the upload buys no independent failure domain today.
  • Remediated by: spec: durability window (#322).

I-EL-5 — a crash does not lose fsync'd records outside the manifest. Holds.

  • Forced by: active-part replay on open.
  • Detected by: recovered_active_records. Read a rising count as a report of hard kills, not an error — a clean shutdown seals, so a non-zero value means the process did not exit cleanly.

2. Projection

I-PR-1 — fold order equals event-log order. A projection is a deterministic fold; replaying the same events yields the same state.

  • Forced by: the fold consuming the shard's ordered stream.
  • Detected by: canonical_state_digest parity between a materialised state and a re-fold.
  • ⚠ A digest comparison can agree with itself. An in-path verdict that materialises and then re-folds from the same source agrees by construction; only a cross-store comparator sees tier-vs-Postgres divergence.

I-PR-2 — a projection is rebuildable from the log alone. Holds — but only if the mirrored payload carries every field the fold reads.

  • ⚠ This failed once, silently. The mirror did not write context, which apply_event reads, so a tier-sourced fold was guaranteed to digest-differ by construction. Repointing the fold at the tier without fixing the payload would have alarmed on every execution. Payload version is now pinned (mirror_payload_version).

I-PR-3 — recovery covers completed executions, not only in-flight ones. Holds since the fold reads the durable tier.

  • ⚠ Coverage was ~0 by construction before that: the spine holds in-flight work, and the comparator asks about completed work.

3. KV

I-KV-1 — the fold is clock-free. No wall-clock input, so the same events yield the same value on any node at any time.

I-KV-2 — ❌ no serve path. kv is absent from SERVE_WIRED_TIERS. It can be set to primary and will not serve. The guard exists specifically because an earlier state "told operators a serving tier was inert".


4. Object

I-OB-1 — parts are immutable and content-addressed. A part is written write-once; re-publishing identical bytes is idempotent (newly_written distinguishes a genuine write from an idempotent re-publish at the same length and digest).

I-OB-2 — immutability removes the need for consensus in replication. N-way copy of immutable objects cannot conflict, so replication is the HDFS block-replication model, not a replicated log. ⚠ This is why the per-shard-Raft plan was retired — and it is a statement about replication, not about write ownership, which still needs I-EL-2.

I-OB-3 — ❌ no serve path. As I-KV-2.

I-OB-4 — a reader never re-pulls a reclaimed segment. Holds even in the crash window where the reclaim watermark is committed but the shared objects are not yet deleted, because a cold-load materialises only segments above the watermark.


5. What a reader observes during a shadow → primary transition

⚠ This is the section the cutover questions actually turn on, and it is about visibility, not ordering. Ordering is a property of one tier; visibility is what a reader sees while two tiers disagree about who is authoritative.

The three states a tier passes through

state who answers a read what the shadow is for
off the incumbent, always nothing is written
shadow the incumbent, always the tier is written and compared; a divergence is a metric, never a served byte
primary the tier — only if it is in SERVE_WIRED_TIERS the incumbent becomes the fallback

⚠⚠ The middle state is the one that misleads. A tier in shadow is fully written and fully compared and still serves nothing. Watching its append rate or its store size tells you the mirror is alive; it tells you nothing about whether a flip would be correct. Only the parity verdict does.

⚠⚠ And the third state is not reachable for two of the four engines. Setting NOETL_EHDB_KV=primary or NOETL_EHDB_OBJECT=primary promotes nothing — the flip is accepted as configuration, tier_serves_primary returns false, and the promotion is logged as primary_not_wired. See Architecture — the four engines.

What is dual-written, and what that costs a reader

A shadow tier is written in addition to the incumbent, never instead of it. So during shadow there is no window in which a reader can see a value that only the shadow holds — which is exactly why shadow is safe, and also why it proves less than people expect.

Measured on prod, 2026-09-16, on noetl-server-rust-embedded-0:

setting value what it means for a reader
NOETL_RESULT_MINT_AUTHORITATIVE true the canonical _uri is minted as the authoritative locator
NOETL_RESULT_STORE_DUAL_WRITE false results are not being written to both stores
NOETL_EHDB_PROJECTION_READ_SOURCE wal projection reads fold the WAL
NOETL_EHDB_PROJECTION_SERVE_ON_BEHIND true a projection that is behind still answers
NOETL_OBJECT_STORE_BACKEND gcs the object substrate is GCS, not Postgres

⚠ NOETL_EHDB_<TIER> mode variables are not set on the prod pods at all — the live modes come from code defaults, not from a manifest. Anyone reasoning about "what prod is configured to do" should read the defaults, not the deployment.

⚠ SERVE_ON_BEHIND=true is a visibility decision with teeth: a reader can be served from a projection that has not caught up, so projection reads are not read-your-writes against the event log.

What parity actually checks — and what it deliberately does not

The comparator answers one question: do the authoritative store and the shadow agree about the same key at the same version? Three things follow that routinely surprise people.

1. Only three things count as divergence.

DIVERGENCE_KINDS = ["digest", "size", "missing_object"]

2. superseded and unmirrored are observations, not divergences. A shadow holding an older version of a key that has since been rewritten is expected, not broken. Counting it as divergence would make every busy key look corrupt.

3. Arrival order is counted, not judged. The event-log comparator reports arrival_reordered — records that arrived out of tier order — as a number, not a divergence (#346). An earlier version treated reordering as disagreement and produced a false alarm: 8 of 74 "divergent", which fell to 0 of 10 once the comparator stopped confusing arrival order with disagreement. ⚠ The lesson generalises: a parity number is only as trustworthy as its definition of "different", and the first definition was wrong.

⚠⚠ What parity does not tell you: that the shadow could have served the read. Parity compares stored bytes. Serving additionally requires a read path, a fallback, and a recovery story — none of which parity exercises.

The recovery ladder, and why it is part of the visibility story

A reader's worst case is not a wrong answer but no answer. Recovery folds from, in order:

RECOVERY_FOLD_SOURCES = ["spine", "tier", "postgres"]

⚠ The last rung matters for a cutover argument: a tier that cannot answer no longer ends recovery, because Postgres is still behind it (noetl/server#441). Audited 2026-09-16 — every event-log read path still has Postgres authority behind it; none is served from the tier alone. That is what makes the current shadow posture reversible.

⚠⚠ Conversely, it is what a flip would spend. Promoting a tier to primary is the step that removes the authority behind a read path, and the audit above is only true while that has not happened.

Reading checklist before any flip

  1. Is the tier in SERVE_WIRED_TIERS? If not, the flip is a no-op with a log line.
  2. Does parity report 0 divergences over a window that contains real traffic — and is the comparator's definition of divergence the right one?
  3. Is the incumbent still authoritative behind every read path the flip touches?
  4. What does a reader see during the flip — and can it be reversed without a write that only the new primary holds?

What is not an engine

  • catalog — a projection with a StoreTier. A thing gains a StoreTier when it gains a store; it does not thereby gain an engine.
  • vector — a vectorized projection. Bounded cosine over an execution-scoped candidate set; no ANN index, no StoreTier, absent from SERVE_WIRED_TIERS.

See Architecture — the four engines.

Related

Architecture — the four engines · Backend Configuration · #321 election + fencing · #322 durability window · #324 scope + de-risking gate

Clone this wiki locally