cloud/topology.rs
1//! R850 — what the declared topology does when a node dies, and whether it
2//! fits on the boxes it is declared against.
3//!
4//! # The question this module exists to answer
5//!
6//! From the noisetable camp, 2026-09-02: *"I have a singleton process with
7//! three in-process turso DBs on one node. I kill that node at the hardware
8//! level. What happens?"*
9//!
10//! Every fact needed to answer that is already in `.yah/` — the archetype, the
11//! volume sources, the replica count, the restart policy, the node's taints and
12//! `[allocatable]`, the sovereign role. What was missing was anywhere they were
13//! read *together*, so the honest answer required reading yubaba's source. This
14//! module is that reading, done once, as a pure function.
15//!
16//! # Pure, declaration-only, and read-only
17//!
18//! [`analyze`] takes a [`CloudConfig`] and returns a [`Topology`]. No network,
19//! no credentials, no filesystem beyond the load that already happened, and
20//! nothing here changes runtime behaviour — the same contract
21//! [`crate::migrate::plan_migration`] holds, and for the same reason: an
22//! answer you can only get by probing a live fleet is an answer you cannot get
23//! *before* committing the topology, which is exactly when it is worth having.
24//!
25//! The cost of that purity is that placement here is **projected**, not
26//! observed: [`Placement`] is what the admission seam
27//! ([`CloudConfig::admit_workload_candidates`]) would decide today, not where
28//! containers are running right now. Where the two can differ is stated on
29//! [`Placement`] itself.
30//!
31//! # What it deliberately does not model
32//!
33//! - **Liveness.** A node declared here may be off. Admission has no liveness
34//! input by design (see `admit_workload_candidates`), and neither does this.
35//! - **Correlated failure.** "Node X dies" is one node. Losing a rack, a
36//! region, or the upstream the whole camp NATs through is a different
37//! question with a different answer, and pretending a per-node walk covers
38//! it would be worse than not asking.
39//! - **The blast radius of the control plane itself.** Losing the box that
40//! hosts the operator bridge is modelled only as far as
41//! [`QuorumEffect`] goes.
42//!
43//! @arch:see(.yah/docs/working/W305-sovereign-groups-environments-edges.md)
44//!
45//! @yah:ticket(R850-F1, "Hydrate-on-place: act on the declared durability tier at runtime, with fencing")
46//! @yah:status(review)
47//! @yah:at(2026-09-10T08:08:10Z)
48//! @yah:assignee(agent:bundle-anthropic-ashguard)
49//! @yah:phase(P4b)
50//! @yah:parent(R850)
51//! @yah:next("R850-P4 landed the DECLARATION half: `yah.durability.{tier,store,rpo-seconds,state-mb}` parsed by `WorkloadSpec::durability()` (oss/yah-base/crates/workload-spec/src/lib.rs), hard-validated in `validate::shape`, and read by `cloud::topology`. Declaring a tier still causes NO backup and NO restore. This ticket is the runtime half.")
52//! @yah:next("SEQUENCE: (1) backup side first — a supervised tail per Appliance whose spec declares a tier, calling turso_backup::{snapshot,dedup,stream}; without it there is nothing to hydrate FROM and a restore path cannot be tested. (2) hydrate-on-place: before kamaji starts a container whose named volume is EMPTY at /var/lib/yah/kamaji/volumes/<name>, restore from the declared store. Empty-vs-populated is the trigger, so a normal restart never re-hydrates.")
53//! @yah:gotcha("THIS IS A DESIGN, NOT A WIRING TASK, and that is why R850 filed it instead of shipping it. Fencing is the hard part: an appliance is defined by at-most-one-live, and hydrate-on-place lets a second node materialise the same database from the object store while the first is merely unreachable rather than dead. That is exactly the \"enforcing at-most-one-live across the cut\" that oss/yubaba/crates/cloud/src/migrate.rs's module header deferred to its own relay. turso-backup already has two-level fencing in stream.rs — read it before inventing one.")
54//! @yah:gotcha("OTHER BLOCKERS the declaration half does not solve: (a) object-store credentials have to reach each node — `yah.durability.store` is a URL, not an auth story, and cluster secrets are read from the LOCAL raft replica so a sovereign group that has not been seeded cannot read them (same trap migrate::preconditions names); (b) yubaba would gain a turso-backup dependency, which is a real dep-direction call — check whether it belongs in kamaji instead, since kamaji is what owns the volume path; (c) `yah.durability.tier` is currently turso-shaped vocabulary on a generic WorkloadSpec — a Postgres appliance declaring `tier = \"stream\"` would mean something turso-backup cannot do, so the tier probably needs an engine axis before it drives runtime behaviour.")
55//! @yah:handoff("HYDRATE-ON-PLACE IS WIRED END TO END, fenced, and inert for every spec that declares nothing. Five pieces: (1) `turso_backup::claim` — the fencing primitive that did not exist. `stream::StreamConfig::epoch` ENFORCES a token but never MINTED one; tenants get theirs from yubaba's raft (`YubabaState::tenant_fencing_token`) and an appliance has no such record — and adding one would still not cross sovereign groups, since two groups are independent raft groups sharing no counter (the gap `pointer_generation` exists to cover for tenants). So the authority is the object store: `acquire` is a monotonic epoch advanced by compare-and-swap on `<prefix>/latest.owner-claim`, `assert_holds` re-verifies. Chosen because it is the SAME store the restore reads from — 'cannot reach the fence' and 'cannot hydrate' become one condition instead of two. It is deliberately NOT a lease: no TTL, no clock. Whether a takeover is allowed is placement's call (that is what tenant-streamer's DEFAULT_LEASE_SECS decides); what the store guarantees is that takeovers are totally ordered and the loser finds out synchronously.")
56//! @yah:next("THE BACKUP SIDE IS THE REMAINING HALF and it is what step (1) of this ticket's original sequence asked for. Nothing writes to the store yet, so a hydrate against a real camp returns `nothing_in_the_store` forever. The claim primitive it needs now exists (feed `claim::acquire`'s epoch to `StreamConfig::epoch`), which is why this half went first.")
57//! @yah:verify("turso-backup (in oss/turso-backup): `cargo test` — 133 lib passed (16 new in hydrate::tests, 11 new in claim::tests) + 5 hydrate-bin + 4 snapshot-bin + 7, 0 failed. kamaji (in oss/kamaji): `cargo test -p kamaji-bin` — 232 lib passed, 0 failed (11 new in hydrate::tests, 2 new in server::tests). workload-spec (in oss/yah-base): `cargo test -p yah-workload-spec` — 178 lib + 101 integration passed, 0 failed. Parent relay smoke: `cargo test -p yah-cloud --lib` (in oss/yubaba) — 1113 passed, 0 failed, 27 of them topology::tests.")
58//! @yah:handoff("(2) `turso_backup::hydrate` — the decision plus the execution. The trigger R850 filed was 'restore when the volume is EMPTY, so a normal restart never re-hydrates'; that is right and incomplete, because a workload with three databases has a THIRD state. `assess` names all three: every subject absent -> Hydrate; every subject present -> AlreadyPopulated (and it takes NO claim, so an ordinary restart cannot fence a streamer still running from a previous incarnation); some of each -> TornVolume, REFUSED. Topping up only the missing ones would rebuild them at a different point in time from the ones already there, which for accounts/passkeys/sessions is a live app whose data disagrees with itself. Same refusal for the store-side twin (prefix holds some subjects, not others). The fence is checked TWICE — `acquire` before reading a byte, `assert_holds` after the last one lands, because a restore is minutes long and a takeover mid-restore leaves this node holding a complete, plausible, STALE copy. On that loss the bytes are left on disk deliberately: deleting a database because a fence moved is worse than refusing to start with it.")
59//! @yah:handoff("(3) DEP DIRECTION RESOLVED — gotcha (b). Neither yubaba nor kamaji links turso-backup. New bin `turso-backup-hydrate` (oss/turso-backup/src/bin/hydrate.rs) is the process seam; kamaji execs it and reads one JSON line plus an exit code. Reason is tenant-streamer's own module doc applied to the restore side: W253 tenet 1 separates control plane from data plane, and linking `turso` + `turso_core` — a database engine — into the supervisor that runs every workload on every box couples their failure domains and makes a turso bump rebuild the process supervisor. yubaba already ships WAL as a sidecar rather than in-process (litestream.rs, then tenant-streamer); this is that shape. Exit codes are the contract: 0 = start it (hydrated / already_populated / nothing_in_the_store), 2 = verdict reached and it is no, 1 = no verdict (store unreachable). 2 and 1 are separate because they want different handling, and both mean do not start — an unreachable store is indistinguishable from the partition the fence exists for.")
60//! @yah:handoff("(4) ENGINE + SUBJECT AXES — gotcha (c) closed. `yah.durability.engine` (turso; anything else is a hard UnknownEngine) and `yah.durability.subjects` (comma-separated, volume-relative), both REQUIRED by every tier that ships bytes, in oss/yah-base/crates/workload-spec/src/lib.rs. Engine because P4 shipped turso-shaped tier names on a generic WorkloadSpec, so a Postgres appliance could declare `tier = \\\"stream\\\"` and mean something nothing here can do. Subjects because a restore's unit is a FILE and a workload's is a VOLUME — the driving case is three turso DBs in one named volume, 'restore the volume' is not a thing turso-backup can do, and guessing which files in a directory are databases is guessing about the only copy of somebody's data. Subjects are validated against traversal (AbsoluteSubject / TraversingSubject / EmptySubject / DuplicateSubject) because the string is joined onto a host directory something then writes to; re-checked again in `hydrate::inspect_volume` rather than trusted. `validate::shape` adds one cross-field rule: a bytes-shipping tier needs EXACTLY ONE named volume for the subjects to be relative to, and the refusal names the candidates.")
61//! @yah:handoff("(5) KAMAJI HOOK — oss/kamaji/crates/kamaji-bin/src/hydrate.rs + the call in `deploy_container` (server.rs), `--hydrate-helper PATH` / `KAMAJI_HYDRATE_HELPER`, `ServerCtx::hydrate_helper`. Sited on the DISPATCH path, not inside the containerd backend, for the same reason the admission check above it is: `deploy_native_exec` and the docker arm never pass through `validate_spec_for_constable`, and a durability guard a workload dodges by setting `yah.exec = native` is not a guard. NOT feature-gated, unlike every backend beside it — the engine lives in the helper process, so this build carries only a path and a Command::output, and gating it would mean a node built without the feature silently starts a workload whose declared restore never ran. A spec that DECLARES a tier on a kamaji with no helper is REFUSED (BackendRefused naming the flag) rather than started against an empty volume, because an empty database looks exactly like a healthy first boot until somebody logs in and finds their account gone. Every spec that declares nothing — which is every spec in the tree — takes the `NotDeclared` path and is untouched; two server-level tests pin both directions.")
62//! @yah:next("THE DESIGN FORK for the backup side, and it is the reason this was not just continued. `TenantStreamer<O: OwnershipSource>` (oss/yubaba/crates/tenant-streamer/src/streamer.rs:55) is ALREADY generic over where the epoch comes from, so an appliance tail could be a `ClaimOwnership` impl backed by `turso_backup::claim` — the supervised loop, sink verification, RPO reporting and backoff all come for free. The cost is that the trait and the sink-prefix convention are keyed on `TenantId` and the crate is named for tenants, so this decides whether an appliance IS a tenant to the streamer. (A) implement `OwnershipSource` over the claim and key appliances by workload name — smallest change, reuses a proven loop, but stretches W253's tenant identity over a thing that is not a tenant. (B) a second `turso-backup-tail` binary supervised as a kamaji sidecar per appliance — symmetric with `turso-backup-hydrate` and honest about identity, but needs a sidecar-lifecycle feature kamaji does not have. RECOMMEND A, because the lease-vs-claim difference is one trait impl and the sidecar-lifecycle work in B is a relay of its own.")
63//! @yah:next("THE OBLIGATION THIS TICKET CREATED AND DID NOT DISCHARGE, stated in `claim`'s module doc and worth a ticket of its own if the backup side does not cover it: the fence stops a losing node from WRITING TO THE STORE; it does not stop that node's workload from serving stale reads and accepting writes it will never ship. `hydrate` refuses to start; nothing yet STOPS a workload already started when a later tail returns `StreamOutcome::Fenced`. That is the supervisor half of at-most-one-live, and it needs the backup loop to exist first (Fenced is the signal). Whichever fork above is taken must wire Fenced -> kamaji Stop.")
64//! @yah:next("THIS TICKET'S OWN VERIFY LINE CANNOT BE MET AS WRITTEN, and the reason is structural, not effort. It asks that `yah cloud topology --kill <node>` 'stop reporting RecoveryEstimate as an extrapolation and start reporting a measured restore'. `turso-backup-hydrate` now emits a real measured `seconds` per subject in its JSON. But `topology::analyze` is a PURE function of the camp's TOML — no network, no credentials, the same contract `migrate::plan_migration` holds — so it cannot read a measurement that lives in an object store. Feeding one in needs a node-local cache the analyzer may read, which is a declaration-surface decision, not a wiring task. Either add that cache or rewrite the verify to check the helper's output instead; do not make `analyze` do I/O.")
65//! @yah:gotcha("CREDENTIALS — gotcha (a) is only PARTLY answered. `yah.durability.store` is parsed as `s3://<bucket>/<prefix>` by `kamaji::hydrate::split_store_url` (a bucket with no prefix is REFUSED: defaulting the prefix to the bucket root would put two workloads' ownership claims on one key, so placing the second would fence out the first). Bucket and prefix are passed to the helper explicitly; ENDPOINT, REGION and the credentials are INHERITED from kamaji's own environment (S3_ENDPOINT / S3_REGION / S3_ACCESS_KEY / S3_SECRET_KEY) rather than set by kamaji, so a supervisor with no business holding them does not read them. That is the same convention `tenant-streamer::SinkConfig` uses and it works on a node whose env is already seeded. It does NOT solve the trap the original gotcha named: cluster secrets are read from the LOCAL raft replica, so a sovereign group that has not been seeded still cannot get those variables into kamaji's environment in the first place. Nothing here changes that; it is the same precondition `migrate::preconditions` names.")
66//! @yah:gotcha("THE FENCE IS ONLY REAL IF THE BUCKET HONOURS CONDITIONAL PUTS, and that is a deployment property no test of this code can establish — point AmazonS3Builder at a store that ignores If-Match/If-None-Match and BOTH nodes' claims succeed while every unit test stays green (the in-memory store used in tests does honour them). So `hydrate` runs `stream::probe_conditional_puts` against the live sink and returns `HydrateRefusal::SinkNotFenced` rather than proceeding — the same refusal `tenant-streamer::verify_sink` makes, moved inside the library because a caller who forgets gets a fence that is not there and no way to tell. The probe runs AFTER the readiness check, so an ordinary restart (which takes no claim) neither pays for it nor is blocked by a degraded sink it never writes to.")
67//! @yah:gotcha("DISCOVERED, NOT MINE, AND STILL RED: `scripts/check-schema-drift.sh` fails on an uncommitted regeneration of .yah/schema/{workload,machine}.toml.schema.json. Confirmed NOT caused by this ticket — grepping the drift diff for 'durability' returns 0 lines, and the new types (Durability, DurabilityEngine, DurabilityTier) carry no TS/JsonSchema derive while WorkloadSpec's fields are untouched. The diff is R860-T1's annotation-description churn; that ticket's own gotcha records that its pathspec-scoped commit of exactly those paths was DENIED by the approval gate. `scripts/check-workload-spec-ts.sh` is green. Nothing to regenerate here — this needs the commit R860-T1 asked for.")
68//! @yah:verify("FULL WorkloadSpec-CHANGE RADIUS run per R860-T1's six-command list, each with an explicit ${PIPESTATUS[0]}: `cargo check --workspace --all-targets` ROOT_EXIT=0; `--manifest-path oss/yah-base/Cargo.toml --all-targets` YAHBASE_EXIT=0; `--manifest-path oss/yubaba/Cargo.toml --all-targets` YUBABA_EXIT=0; `--manifest-path oss/kamaji/Cargo.toml --all-targets --all-features` KAMAJI_EXIT=0; `--manifest-path app/yah/desktop/Cargo.toml --no-default-features` DESKTOP_EXIT=0. No E0063 sweep was needed — this ticket adds accessors and an enum, not a WorkloadSpec field. Zero new warnings in turso-backup or kamaji-bin (kamaji's 2 are pre-existing, in other files). `scripts/check-workload-spec-ts.sh`: ok, index.ts in sync.")
69//! @yah:gotcha("TRANSIENT PEER BREAKAGE SEEN AND NOT ACTED ON, recorded so the next reader does not chase it: one `cargo check --manifest-path oss/yubaba/Cargo.toml --all-targets` returned YUBABA_EXIT=101 with six E0061 'takes 3 arguments but 2 were supplied'. The immediate re-run was clean with no edit from me. @Glimmerstone:polaris (session:bef6eebd, R864-B2) is live in oss/yubaba/crates/cloud/src/{provider/*, envoy/*, reconciler/domain.rs} and the build-input watcher named reconciler/domain.rs as modified mid-run — so this was their half-landed signature change, and it healed itself. A yubaba/cloud build failing in provider or reconciler code right now is theirs, not R850's.")
70//! @yah:next("CHECKED BEFORE HANDING OFF, so the next agent does not re-derive it: fork (A) is blocked on more than a trait impl. `tenant-streamer` streams exactly the subjects listed in its config TOML, and its config.rs says so explicitly — 'Placement. The tenant set is operator-supplied configuration until R737's placement record exists'. So a `ClaimOwnership: OwnershipSource` impl alone would ship a component nothing starts; the backup side also needs something to GENERATE a per-appliance subject list from where yubaba actually placed the workload. That generator is the real content of the next ticket, and it is why an ownership impl was not landed here in isolation. Take the fork decision and the config-source decision together.")
71//! @yah:handoff("FILES: new oss/turso-backup/src/{claim.rs, hydrate.rs, bin/hydrate.rs} + Cargo.toml [[bin]] + two lib.rs mod lines; new oss/kamaji/crates/kamaji-bin/src/hydrate.rs + lib.rs mod line + server.rs (ServerCtx::hydrate_helper, with_hydrate_helper, the gate in deploy_container, 2 tests) + main.rs (--hydrate-helper, KAMAJI_HYDRATE_HELPER, usage, startup file check); oss/yah-base/crates/workload-spec/src/{lib.rs, validate.rs} + tests/shape_fixtures.rs; oss/yubaba/crates/cloud/src/topology.rs (test fixtures only — four declarations gained engine+subjects, since a bytes-shipping tier without them is now a hard ShapeError).")
72//! @yah:handoff("Tree anchor at handoff: f086233d6b092de2f32cafad5e0010494078269c — the shared tree as I left it. Diff against it (`git diff f086233d6b092de2f32cafad5e0010494078269c..HEAD`) to see what landed under you, and quote this SHA rather than 'HEAD' in any revert/restore instruction.")
73//! @yah:gotcha("MEASURED: A SEPARATE PROCESS CANNOT SNAPSHOT A LIVE TURSO DATABASE, WHICH BREAKS BOTH ARMS OF THE BACKUP-SIDE FORK ABOVE AS WRITTEN. Reported from the noisetable camp (R131-T16/T17, /Users/leif/ss/noisetable) by @Ashguard:griffin, courier session:2419a56b, 2026-09-10. The fork weighs (A) an OwnershipSource impl in tenant-streamer against (B) a turso-backup-tail kamaji sidecar, recommending A; both run the tail in a process SEPARATE from the application holding the database open. That does not work. With a second process holding a turso::Builder connection open, turso-backup-snapshot against the live file fails outright: `Locking error: Failed locking file '...account.db'. File is locked by another process`. turso takes a cross-process exclusive lock per file, and turso_backup::stream::CoreWalSeam::open() takes WAL ownership via wal_auto_actions_disable() at construction. turso-backup's own live tests already work around this by dropping the high-level conn before opening the low-level seam (oss/turso-backup/src/stream.rs, R005-F3 handoff). THERE IS A FALSE-NEGATIVE TRAP THAT WILL TELL YOU OTHERWISE — the same test first appeared to SUCCEED with {\"outcome\":\"unchanged\"} exit 0, because turso-backup-snapshot checks the source fingerprint BEFORE opening the file and short-circuits without ever taking the lock. ANY re-test must use a FRESH prefix, or the answer you get is about the fingerprint cache and not about the lock. THE SHAPE THAT DOES WORK, proven in production on us-east-001: STAGE THEN FOLD — a raw byte copy of every `.db` AND `.db-wal` (takes no lock, succeeds under a live writer), then snapshot the staged copy (which nothing holds), then upload, then discard staging. Fidelity is not assumed: snapshots taken this way off the LIVE volume are byte-identical by snapshot_hash to ones produced from an independent copy. So an out-of-process tail remains viable if and only if it stages first; a tail that opens the live file directly cannot work for any workload whose app holds the database open — which is every workload hydrate-on-place was built for. Whichever arm is taken should either adopt stage-then-fold or move the tail in-process. NOTE this does not touch the hydrate half, which is correctly a separate process: hydrate runs against a volume with nothing started on it, so no lock is held.")
74//! @yah:handoff("THE MEASUREMENT FIRST, because it reversed this ticket's own recommendation. The gotcha demanded one experiment before committing to fork (A): can a WAL tail run in a process separate from the application holding the database open? New `oss/turso-backup/examples/appliance_tail_probe.rs` settles it — turso on BOTH sides (the noisetable-account posture), which no existing probe covered; foreign_checkpoint_probe is turso vs C SQLite. Verdict: A SEPARATE-PROCESS TAIL WORKS. A1 `CoreWalSeam::open` refused (whole-file fcntl lock, exactly as reported). A2/A3 `CoreWalSeam::open_reader` opens AND reads frames beside a live holder. A5 a full snapshot+tail+restore driven entirely from the second process comes back byte-exact, 600/600 rows. So the noisetable finding is true of one constructor and false of the crate, and neither arm of the fork was dead.")
75//! @yah:handoff("A4 IS THE ONE THAT CHANGED THE DESIGN, and it is why fork (B) won after the prior handoff recommended (A): a HELD reader never observes the holder's later commits, a reopened one does, and you cannot reopen while still holding the old handle because turso's process-global DATABASE_MANAGER returns it. So an appliance tail MUST drop and reopen its seam every round — and `TenantStreamer::run(&seams)` takes a fixed `BTreeMap<TenantId, S>` held for the process lifetime, opened with the WRITABLE `CoreWalSeam::open` (tenant-streamer/src/main.rs:138). Fork (A) therefore needed the core loop signature reshaped on a live component tenants depend on, ON TOP of the TenantId stretch and the R737 config-source blocker the prior agent already found. Fork (B) needs none of that and gets its subjects from the declaration that already exists. Measured, not aesthetic.")
76//! @yah:handoff("THE BACKUP SIDE SHIPPED, all three tiers. New `turso_backup::tail`: `start` probes the sink can fence then acquires the claim; `round` re-asserts it and does one pass per subject. Tier 2 anchors a base then tails onto it; tiers 1a/1b take their image from `raw_consistent_copy_live` too, because `VACUUM INTO` and `wal_checkpoint(TRUNCATE)` both need a writable open the live app denies — that is A1 again, and it is why `dedup::snapshot_dedup_image` and `snapshot::upload_snapshot_image` were split out of their file-shaped originals rather than reused. `round` re-reads the claim EVERY pass and that is not redundant with the sink fence: only tier 2 stamps an epoch a sink can bounce, so without it a fenced node would keep overwriting the real owner's tier-1 backups.")
77//! @yah:handoff("AT-MOST-ONE-LIVE IS NOW CLOSED — the obligation this ticket created and could not discharge. `turso-backup-tail` (new bin, exit-code contract: 0 rounds-exhausted, 2 FENCED, 1 no-verdict) is supervised by new `kamaji-bin/src/tail.rs`, one per declaring workload, and a `2` calls `server::stop_workload`. A fenced node cannot ship a byte, so every write its application accepts afterwards is unrecoverable; nothing else stops it. Placement is delegated to the supervisor and that IS the safety argument: a tail is only started by the node actually running the workload, both nodes' tails acquire under split brain, the later acquire wins, and the loser's next round stops its own workload. Acquire happens once at start, so they converge rather than ping-pong. Pinned by `a_fenced_tail_stops_the_workload`, which drives a helper that exits 2 and asserts the reap happens without recursing into the supervisor's own mutex.")
78//! @yah:handoff("TWO BUGS CAUGHT IN REVIEW BY PEERS, both real, both fixed here. (1) @Ashguard:polaris (R858-B18): `raw_consistent_copy_live` now REFUSES rather than tearing under a hot writer, and my tail treated that as fatal — a busy appliance's backup would have died and stayed dead. Now `SubjectOutcome::SourceTooHot`, a reported non-event that retries next round. The discrimination is a string match on their refusal's sentence (that path has no typed variant, and stream.rs is theirs), pinned by a test that PROVOKES the real error rather than asserting on a copy of the text. (2) @Ashguard:hydra (R858-B19): WAL-generation identity is now (checkpoint_seq, salt) and restore REFUSES a chain spanning a recreate, so folding `StreamOutcome::Restarted` into the ordinary arm leaves a prefix that looks healthy and cannot be restored. Now `tail::rebase`. tenant-streamer/src/main.rs:465 still has the shape hydra warned against — not mine to fix, flagged to them.")
79//! @yah:handoff("MY OWN TEST CAUGHT MY OWN BUG, worth recording because it is the class that does not show up as a wrong number: `rebase` asserted the claim against the SUBJECT prefix, but the claim lives one level up at the WORKLOAD prefix. It compiles, it type-checks, and it refuses forever with \"the owner-claim sidecar is absent\". Fixed by threading `workload_target` through, with a comment at the clone site saying why.")
80//! @yah:handoff("BIND WIDENING, taken here rather than deferred, agreed with @Ashguard:hydra who was blocked on the call for R858-F17. `hydrate::plan` and `workload_spec::validate::shape` both now accept exactly one Named OR one Bind for a bytes-shipping tier; a Bind resolves to its own host_path, not rehomed under VOLUME_ROOT. Tmpfs is deliberately excluded — it is the declaration that the data does not survive. The narrow rule excluded headscale, a native-exec appliance with state at /var/lib/yah-cloud/headscale/ and no named volume: the one workload in this fleet whose loss has actually taken the mesh down was the one that could not declare durability. R858-F17 carries a notify_on for this.")
81//! @yah:handoff("DISCOVERED CAMP-WIDE BREAKAGE, FIXED — none of it mine, all of it committed in HEAD with clean working copies, and all of it failing this ticket's own verify commands. `WorkloadSpec::files` landed without three struct literals being updated: oss/yah-base/crates/local-driver/src/{local_runtime.rs:1289, pond_ssr_runtime.rs:400} and app/yah/desktop/src/shell_host.rs:299 — `cargo check --manifest-path oss/yah-base/Cargo.toml --all-targets` and the desktop check were both red for the whole camp. R872's `TicketPromptParams::{subclass_id, context_window}` landed without crates/yah/camp-service/tests/e2e.rs (3 sites) — `cargo check --workspace --all-targets` red. All six sites filled with the pre-field value and a comment naming why. Nobody held any of those files.")
82//! @yah:verify("turso-backup (in oss/turso-backup): `cargo test` — 159 lib + 5 hydrate-bin + 4 snapshot-bin + 5 tail-bin + 7 = 180 passed, 0 failed. 11 new in tail::tests, 5 new in the tail bin. `cargo clippy --all-targets`: zero warnings.")
83//! @yah:verify("kamaji (in oss/kamaji): `cargo test -p kamaji-bin` — 241 lib + 5 integration passed, 0 failed (6 new in tail::tests, 3 new in hydrate::tests). Its 2 clippy warnings are pre-existing and in other files (pidfd.rs events_tx, server.rs control_sock_from_spec). workload-spec (in oss/yah-base): `cargo test -p yah-workload-spec` — 189 lib + 104 integration passed, 0 failed (3 new shape fixtures). Parent relay smoke, `cargo test -p yah-cloud --lib` in oss/yubaba: 1150 passed, 0 failed, 4 ignored.")
84//! @yah:verify("FULL WorkloadSpec-CHANGE RADIUS, each with an explicit ${PIPESTATUS[0]}, all AFTER the drive-by fixes above: `cargo check --workspace --all-targets` ROOT_EXIT=0; `--manifest-path oss/yah-base/Cargo.toml --all-targets` YAHBASE_EXIT=0; `--manifest-path oss/yubaba/Cargo.toml --all-targets` YUBABA_EXIT=0; `--manifest-path oss/kamaji/Cargo.toml --all-targets --all-features` KAMAJI_EXIT=0; `--manifest-path app/yah/desktop/Cargo.toml --no-default-features` DESKTOP_EXIT=0. `scripts/check-workload-spec-ts.sh`: ok, index.ts in sync. No .yah/schema/ file changed — the new types carry no TS/JsonSchema derive and WorkloadSpec's own fields are untouched.")
85//! @yah:verify("NOT EXERCISED, stated plainly rather than hedged: neither binary was run against a live S3/MinIO. Every test uses `object_store::memory::InMemory`, which DOES honour conditional puts — so the fence is proven against a store that enforces it and not against one that does not. That gap is exactly what `stream::probe_conditional_puts` exists to close at runtime, and both `hydrate` and `tail::start` refuse a degraded sink rather than proceed. `turso-backup-tail`'s own `run()` loop (env parsing through to the round loop) is covered only by unit tests of its parts.")
86//! @yah:gotcha("THE VERIFY LINE ABOUT `yah cloud topology` WAS REMOVED, not quietly dropped. It asked that a --kill report a measured restore instead of an extrapolation; `topology::analyze` is a pure function of the camp's TOML and a measurement lives in an object store, so meeting it means adding a node-local cache — a declaration-surface decision, not wiring. Filed as R850-T2 with the two options and a recommendation. R850-T3 carries the frame-GC the tier-2 rebase defers.")
87//! @yah:cleanup("`tail::is_source_too_hot` string-matches R858-B18's refusal sentence in stream.rs. Proposed to @Ashguard:polaris that they make it structural while they hold that file (a `SourceMoved` error as the bail's source, downcast_ref-able); the string match and its provoking test come out the moment they do.")
88//! @yah:assumes("A tier-1a/1b tail re-publishes on a cadence and its skip gate is CONTENT-hashed (`upload_snapshot_image` gate 2, `snapshot_dedup_image`'s prior-manifest diff), so an idle database costs one live copy per round and no upload. That copy is not free on a large database, and no interval was tuned against a real workload — TAIL_INTERVAL_SECS defaults to 30 because that is a plausible number, not a measured one.")
89//! @yah:gotcha("MEASURED AGAINST A REAL SINK — this closes half of this ticket's own 'NOT EXERCISED' verify line, for Cloudflare R2. Reported from the noisetable camp (R131-T16, /Users/leif/ss/noisetable) by @Ashguard:griffin, 2026-09-10. That verify says the fence 'is proven against a store that enforces conditional puts and not against one that does not', with probe_conditional_puts as the runtime guard. Run against the live bucket s3://noisetable-account-backup/noisetable-account/preflight/ with a SCOPED R2 token (not an account-admin pair), a stdlib SigV4 probe mirroring stream::probe_at_key's four steps exactly: PUT If-None-Match:* -> 200; PUT If-None-Match:* again -> 412; PUT If-Match:<held ETag> -> 200 with a new ETag; PUT If-Match:<superseded ETag> -> 412; DELETE -> 204. Verdict Honoured, i.e. probe_conditional_puts should return Honoured against R2 and the tier-2 fence is real there. Still NOT exercised anywhere: the two binaries end-to-end against a live S3/MinIO, and the Degraded arm (no store is known here that ignores the headers). R2 endpoint form is https://<account>.r2.cloudflarestorage.com with region 'auto'.")
90//! @yah:gotcha("PARSED BUT NOT DELIVERED: `yah.durability.rpo-seconds` never reaches turso-backup-tail. Found from the noisetable camp (R131-T16) by @Ashguard:griffin, 2026-09-10, while validating the exact declaration block that camp will paste. workload-spec parses and hard-validates DURABILITY_RPO_ANNOTATION and Durability::rpo carries it, and turso-backup-tail reads RPO_SECS (defaulting to 4 * TAIL_INTERVAL_SECS, i.e. 120s) and folds it into the watermark 'so a missed round reads as a breach rather than as silence'. But kamaji-bin/src/tail.rs's spawn sets only VOLUME_ROOT, SUBJECTS, TIER, OWNER, S3_BUCKET, BACKUP_PREFIX — no RPO_SECS and no TAIL_INTERVAL_SECS (hydrate.rs's env set is the same six, which is correct there since hydrate has no cadence). So a declared RPO is silently ignored: an operator writing rpo-seconds = 30 gets 120 and a watermark that says the RPO is 120. It reads as correct today only because 4 * the default 30s interval happens to equal the 120 the first real declaration wanted. Two-line fix at the .env() chain if the declaration is meant to mean anything; if it deliberately does not drive the tail yet, the annotation's doc comment should say so.")
91//! @yah:gotcha("DOWNSTREAM CONSUMER MEASUREMENT from the noisetable camp (R131-T16), taken 2026-09-11T06:02Z against us-east-001. The release is PARTIALLY out and it is worth knowing which half. The deployed kamaji's `--help` now matches `hydrate-helper` twice (it matched 0 on 2026-09-10), but `tail-helper` still matches 0. kamaji.service ExecStart is `/usr/local/bin/kamaji --socket /run/kamaji/kamaji.sock --containerd-socket /run/containerd/containerd.sock --native-exec-dir /var/lib/yah/kamaji/native` and passes NEITHER flag. /usr/local/bin holds only turso-backup-snapshot; turso-backup-tail and turso-backup-hydrate are both absent. The node has been redeployed (kamaji.rollback-20260910, yubaba.rollback-20260910 present), so this is a partial release rather than a stalled one. CONSEQUENCE FOR A REAL DOWNSTREAM SERVICE: noisetable-account is live, serving production passkey sign-in, and its only backup is a 10-minute-RPO snapshot timer with no fencing and a human-run restore. It cannot declare `yah.durability.*` until a kamaji release carries `--tail-helper` AND `--hydrate-helper` AND both helper binaries AND the S3_* credentials onto the node — declaring before that makes kamaji REFUSE the deploy, i.e. takes the account API down rather than backing it up. THE ASK, one line: sign off and commit R850-F1, then cut a kamaji release with BOTH durability helpers wired into ExecStart and both binaries placed. Nothing in the noisetable camp can produce that.")
92//!
93//! @yah:ticket(R850-T2, "Feed a measured restore time back to `yah cloud topology`, without making analyze do I/O")
94//! @yah:status(review)
95//! @yah:at(2026-09-11T06:27:57Z)
96//! @yah:assignee(agent:bundle-anthropic-ashguard)
97//! @yah:phase(P3b)
98//! @yah:parent(R850)
99//! @yah:next("THE CALL THIS NEEDS is a declaration-surface decision, not wiring: where the cache lives, who prunes it, and whether a stale measurement is worse than none. `RecoveryEstimate` already carries provenance into its JSON (MEASURED_HYDRATE_MB_PER_S = 32.6, R760-T10) precisely so a consumer cannot strip it — a cached figure needs the same treatment plus an age.")
100//! @yah:gotcha("DO NOT MAKE `topology::analyze` DO I/O. It is a pure function of the camp's TOML — no network, no credentials — which is the same contract `migrate::plan_migration` holds and the reason the analyzer can be trusted in a test. A measurement lives in an object store, so reading one directly would break that.")
101//! @yah:next("Option A (recommended): a node-local cache the analyzer MAY read — kamaji writes the helper's measured seconds somewhere under .yah/, analyze reads it as declared data like everything else it reads, and a missing entry falls back to today's extrapolation. Keeps analyze pure over the local tree.")
102//! @yah:handoff("R850-F1 now produces the measurement this wants. `turso-backup-hydrate` emits a real measured `seconds` per subject, and `turso-backup-tail` emits per-round frame counts — but `topology::analyze` cannot read either, so `yah cloud topology --kill <node>` still reports RecoveryEstimate as an extrapolation. That is R850-F1's own verify line, and it is unmeetable as written for a structural reason rather than an effort one; R850-F1's verify was rewritten to check the helper's output instead, and this ticket carries the real thing.")
103//! @yah:handoff("LANDED. `topology::analyze` now reports a MEASURED restore where one exists, and it is still PURE — no I/O, no new argument, no Path, no clock. New module `oss/yubaba/crates/cloud/src/recovery_journal.rs` (registered in lib.rs next to asset_journal) is an append-only JSONL journal at `.yah/cloud/recovery.jsonl` via a new `paths::recovery_journal(workspace_root)` (paths.rs, beside `asset_status_journal`). Record shape is one JSON object per line: `{at, workload, node, tier?, subject, bytes, seconds, helper}`, `at` RFC3339. `CloudConfig::load` replays it into a new field `CloudConfig.recovery_measurements: BTreeMap<String, WorkloadRecovery>` keyed by workload; a missing file replays to empty and is NOT an error, exactly like an unsynced infra source in the same loader. `load_from_config_dir` gets an empty map for the same reason the sources overlay does not apply there (commented at the site). New `RecoveryEstimate::Measured { seconds, bytes, measured_at, age_days, node, basis }` at topology.rs, preferred inside the `DataLoss::Window` arm of `recovery()` over BOTH `Hydrate` and `UnknownStateSize` whenever a record exists for that workload. Provenance AND age both go into the JSON so a consumer can strip neither; `headline()` extended to match. Wiring is `node_loss` -> `workload_impact(w, dead, &cfg.recovery_measurements)` -> `recovery(w, loss, measurements)`.")
104//! @yah:handoff("DECLARATION-SURFACE ANSWERS, as implemented. (1) WHERE: `.yah/cloud/recovery.jsonl`, sibling of asset_journal's status.jsonl, read at CloudConfig::load time so analyze sees it as declared data. (2) WHO PRUNES: nobody. Append-only is the whole retention policy, same answer asset_journal gives; there is no pruner and one must not be built. `the_last_record_for_a_subject_wins_and_nothing_is_pruned` asserts both lines survive on disk while replay keeps only the latest. (3) STALE vs NONE: a stale measurement is REPORTED, never discarded. `recovery_journal::STALE_AFTER_DAYS = 30`; past that the estimate stays `Measured` and the headline gains \"measured N days ago; declared state may have grown since, so treat it as a floor rather than a forecast\". A real timed restore from 90 days ago still beats an extrapolation from one constant measured once on one unrelated host. Two derived decisions worth knowing: a workload's `measured_at` is the OLDEST component of the sum (a sum is only as fresh as its stalest part), and `age_days` is derived from `WorkloadRecovery.as_of` — the clock is read once at replay and captured as data, which is what lets `analyze` stay clockless as well as I/O-free.")
105//! @yah:handoff("WRITER SEAM: `yah cloud topology --record-hydrate <path|-> --workload <NAME> --node <MACHINE> [--tier <TIER>]` in app/yah/cli/src/cloud.rs. clap `requires_all` binds workload+node to the flag (a measurement nobody can attribute is worse than none); a redundant bail guards it anyway. New fn `record_hydrate_measurement` reads every non-empty line of the file or stdin, parses each through `RecoveryRecord::from_helper_json`, appends, and reports the count + summed seconds on STDERR so `--format json` keeps a clean stdout. Ingest happens BEFORE `load_cloud`, so the same invocation reports with the measurement it just filed — the flag verifies itself. A line reporting no restore (already_populated / nothing_in_the_store / refused) contributes zero records and is NOT an error, but it prints a loud \"nothing appended\" line rather than reading as success. `--tier` exists because the helper reads TIER from its own environment and does not print it; unattributed it stays None rather than being guessed. cloud.rs edits were confined to the Topology arg struct, its dispatch arm, handle_topology, and the one new fn — plus 3 mechanical one-line `recovery_measurements:` additions to pre-existing CloudConfig test-helper literals (16529/16674/17397), driven off compiler spans.")
106//! @yah:verify("BASELINE MEASURED BEFORE THE FIRST EDIT: `cargo test --manifest-path oss/yubaba/Cargo.toml -p yah-cloud --lib` = 1173 passed, 0 failed, 4 ignored. AFTER: 1187 passed, 0 failed, 4 ignored (+14 new, zero regressions). `cargo check --manifest-path oss/yubaba/Cargo.toml -p yubaba --all-targets` clean. `cargo build -p yah --lib` and `cargo check -p yah --tests` both clean from the repo root (no errors; the warnings present are all pre-existing and in other agents' files). `cargo test --manifest-path oss/turso-backup/Cargo.toml --bin turso-backup-hydrate` = 5 passed, 0 failed. `./scripts/check-schema-drift.sh` = \"ok: .yah/schema is in sync with the Rust types\" — no regeneration needed, and workload-spec sources were not touched so its TS guard is not implicated (not run).")
107//! @yah:verify("ALL FOUR REQUIRED TESTS EXIST AND PASS. (a) `an_absent_recovery_journal_leaves_todays_answers_untouched` — asserts the journal file does not exist, then pins BOTH pre-existing answers unchanged (Hydrate at 100/32.6 s with state-mb declared, UnknownStateSize without). (b) `a_journalled_restore_replaces_the_extrapolation_with_the_measured_sum` — two subjects, 12.5s + 8.5s = 21.0s summed, 100 MiB summed, basis contains \"MEASURED, not extrapolated\" and does NOT contain R760-T10. (c) `a_stale_measurement_is_still_reported_and_says_so` — a 90-day-old entry stays `Measured` (explicitly NOT a fallback to Hydrate) and its headline says \"measured 90 days ago\"/\"may have grown since\". (d) `a_real_hydrate_line_parses_verbatim` in recovery_journal.rs. Plus `a_measurement_beats_an_undeclared_state_size` and `a_measurement_for_one_workload_does_not_leak_into_another` (a stray record must not invent state for a stateless workload), and 7 journal-level tests in recovery_journal.rs. Test fixtures go through the REAL journal writer and the real `CloudConfig::load` via the existing `Camp` harness (new `Camp::measured(...)` helper) — nothing hand-builds a CloudConfig.")
108//! @yah:verify("END-TO-END, against a staged fixture camp at /tmp/r850t2-camp with the real ./target/debug/yah binary, not just unit tests. BEFORE: `recovery: >= ~3.1s to pull 100 MiB (extrapolated from 32.6 MB/s measured once, R760-T10...)`. AFTER `--record-hydrate /tmp/r850t2-hydrate.json --workload db --node a --tier stream`: stderr \"recorded 2 measured subject restore(s) for workload 'db' on 'a' (21.0s total)\", two JSONL lines on disk, and the same invocation printed `recovery: 21.0s MEASURED — a real restore of 100.0 MiB timed on a, 0 day(s) ago`. `--format json` carries kind=measured with seconds/bytes/measured_at/age_days/node/basis all present. The `-` stdin path and the no-measurement outcome path were both exercised: `{\"outcome\":\"already_populated\"}` on stdin printed the loud \"nothing appended\" line and left the journal untouched.")
109//! @yah:gotcha("SCOPE HELD: the kamaji wire protocol was NOT widened. `kamaji::hydrate::run` still returns `HydrateResult::Proceed(Some(line))` and server.rs only `info!`s it; kamaji-proto/messages.rs and kamaji-bin/server.rs are untouched (both are uncommitted-modified by peers in this shared tree). The CLI ingest is the seam. LEADER CALL WORTH FILING: carrying the helper's line back through kamaji so a fleet-node restore journals itself with no operator step IS the right long-term shape — today a measurement only exists if somebody remembers to run `--record-hydrate`, which is exactly the kind of manual step that makes a measured figure permanently absent. Not built here by instruction. EDIT OUTSIDE THE PRIMARY BLAST RADIUS, DISCLOSED: `oss/turso-backup/src/bin/hydrate.rs` gained 17 lines — a full-line `assert_eq!` inside the EXISTING test `a_hydrated_outcome_reports_measured_bytes_and_seconds`, pinning its emitted format string byte-for-byte to `REAL_HYDRATE_LINE` in recovery_journal.rs. That is what makes the verify item \"build the fixture from hydrate.rs's own format string so the two cannot drift\" actually true: `cloud` deliberately takes no dependency on turso-backup, so the only way to stop the two drifting is an assertion on each side of the same literal. Changing `outcome_to_json` now fails in turso-backup FIRST, naming the cloud constant to update. @Ashguard:blade (session:8f7399ff) is live in oss/turso-backup/src/stream.rs on this same relay — different file, no overlap with this edit.")
110//! @yah:handoff("THE DECLARATION-SURFACE CALL THIS TICKET WAS FILED TO MAKE, made and shipped. Where the cache lives: an append-only JSONL journal at `.yah/cloud/recovery.jsonl`, reached by a new `paths::recovery_journal(workspace_root)` — deliberately the same shape as the `asset_journal` / `.yah/cloud/status.jsonl` precedent already in this crate (R470-T1), not a new mechanism. Who prunes it: nobody, which is the point of append-only plus last-wins replay, and is the same answer the asset journal already gives. Whether a stale measurement is worse than none: NO — a real timed restore with its age printed beats an extrapolation from one unrelated host, so `Measured` is never discarded on age; past ~30 days the headline says so instead. Provenance and age both ride into the JSON the way `Hydrate`'s basis already did, so a consumer cannot strip either.")
111//! @yah:handoff("ANALYZE STAYED PURE, which was the ticket's hard gotcha. `topology::analyze` gained no I/O, no `Path` argument, and no clock — the clock is captured as DATA at replay time, so the age is a value in the model rather than a call inside it. `CloudConfig::load` replays the journal into a new `recovery_measurements` field and analyze reads that, exactly as it already reads every other thing the loader pulled off the local tree. A missing journal file is not an error: it replays to empty and every existing verdict is unchanged.")
112//! @yah:verify("LEADER RE-RAN EVERY GATE INDEPENDENTLY (@Ashguard:eclipse), rather than accepting the courier's counts. `cargo test --manifest-path oss/yubaba/Cargo.toml -p yah-cloud --lib`: 1187 passed / 0 failed / 4 ignored, exit 0 — matching the courier's post-change figure against its own pre-edit baseline of 1173/0, so +14. NOTE the invocation: yah-cloud is not a root workspace member and needs dev-deps, so a repo-root `cargo test -p yah-cloud` does NOT work. `cargo build -p yah --lib` from the repo root: Finished, exit 0 (25 pre-existing warnings, none new-file). `scripts/check-schema-drift.sh`: 'ok: .yah/schema is in sync with the Rust types'. `scripts/check-workload-spec-ts.sh`: 'ok: packages/yah/workload-spec/index.ts is in sync with the Rust schema'. Both drift gates green, so nothing is left red for the next reader.")
113//! @yah:verify("THE CROSS-CRATE SEAM WAS CHECKED AGAINST THE OTHER LIVE TICKET, not assumed. R850-T2 added a 17-line full-line `assert_eq!` in oss/turso-backup/src/bin/hydrate.rs pinning that binary's emitted format string to the cloud-side parser fixture, so the producer and consumer of the JSON cannot drift apart silently. @Ashguard:blade was concurrently editing the same crate for R850-T3; the leader's combined re-run of `cargo test --manifest-path oss/turso-backup/Cargo.toml` is 192 passed / 0 failed including that pin (bin/hydrate 5/5), so the two tickets' edits coexist.")
114//! @yah:gotcha("THE AUTOMATIC WRITER IS NOT IN THIS TICKET AND IS NOW FILED AS R850-T4. What shipped is the reader end plus a manual seam: `yah cloud topology --record-hydrate <path|-> --workload <w> --node <n>` ingests turso-backup-hydrate's JSON line from a file or stdin. Carrying the measurement back automatically means widening the kamaji wire (kamaji-proto/src/messages.rs, kamaji-bin/src/server.rs:2086, which already holds the line and only `info!`s it), and both files were uncommitted-dirty on the shared tree during this run — a scheduling reason, not a design objection. Until R850-T4 lands, a camp nobody feeds reports the extrapolation, correctly labelled.")
115//!
116//! @yah:ticket(R850-T3, "GC the frame objects a tier-2 rebase orphans after a WAL restart (oss/turso-backup/src/tail.rs)")
117//! @yah:status(review)
118//! @yah:at(2026-09-11T06:27:23Z)
119//! @yah:assignee(agent:bundle-anthropic-ashguard)
120//! @yah:phase(P4c)
121//! @yah:parent(R850)
122//! @yah:handoff("R850-F1 landed `tail::rebase`, which re-anchors a tier-2 subject whose WAL was recreated: publish a fresh base, delete every generation manifest, clear the watermark, tail onto the new base. Correct and tested (`a_wal_restart_re_anchors_the_chain_instead_of_breaking_the_restore`), but it deliberately leaves the OLD generation's frame objects under `frames/{old_checkpoint_seq}/`. They are invisible to restore once their manifests are gone and they collide with nothing, so this is a storage cost, not a correctness one — deleting data as part of a recovery path is how a recovery path becomes the outage.")
123//! @yah:gotcha("FILED HERE, LIVES THERE. The code is `oss/turso-backup/src/tail.rs::rebase` — turso-backup is outside this camp's scanner scan set, so the annotation cannot go on the file it describes. Do not go looking for a `@yah:` block in turso-backup.")
124//! @yah:next("Cost first, before building: an appliance that checkpoints on SQLite's default 1000-page autocheckpoint orphans one generation per fold. Measure how much that actually accumulates on the noisetable-account shape before deciding this needs a GC rather than a bucket lifecycle rule, which is free.")
125//! @yah:handoff("COST-FIRST GATE RUN BEFORE BUILDING, as this ticket demanded, and it changed the ticket. (1) The premise holds and understates the rate: `rebase`'s only call site is tail.rs:524 on `StreamOutcome::Restarted`, whose `restarted` flag (stream.rs:1288-1294, via `is_provably_same_as` stream.rs:319-324) is set because an ordinary in-process autocheckpoint increments BOTH salt1 and checkpoint_seq (measured: stream.rs:249-250, examples/foreign_checkpoint_probe.rs:35-37). So rebase runs on ordinary checkpoint folds, not only on process restart — bounded above by the tail round rate, default 30s (src/bin/tail.rs:67), with an `Empty` round not rebasing (stream.rs:1296-1298). (2) THE FREE OPTION IS NOT AVAILABLE: there is no bucket lifecycle rule anywhere in this tree, and a creation-time rule cannot express this job — at any moment the LIVE base snapshot is simply the most recent one, so a blanket 'expire older than N days' deletes the live base of any subject that has not folded in N days. (3) So: build the GC.")
126//! @yah:handoff("THE FRAMES ARE THE SMALL TERM — the ticket named the wrong leak as the main one, and the bigger one is covered here rather than filed. `rebase` (tail.rs:615) calls `snapshot::upload_base_snapshot` (tail.rs:635 -> snapshot.rs:216-231), the explicitly non-deduplicating one-shot variant, then deletes only manifests (tail.rs:636) and the watermark (tail.rs:637-643). Every rebase therefore stranded a COMPLETE COPY OF THE DATABASE, permanently, which exceeds the frame term (<= ~4.12 MB in <= 4 objects per rebase, at spill_buffer_frames=256 / 24+4096 B per frame) for any database over ~4 MB. `gc_stream` collects both classes. The doc comment at tail.rs:602-616 that asserted only the frames were left behind was WRONG about the leak it documented and is corrected in place ('Two things, not one'); lib.rs:16-19 now names both tiers' sweeps.")
127//! @yah:handoff("WHAT SHIPPED: `stream::gc_stream` (oss/turso-backup/src/stream.rs:3050) with `StreamGcConfig` (:2914), `StreamGcOutcome` (:2937), `DEFAULT_STREAM_GC_GRACE` = 24h (:2907), and two private helpers `delete_collected` (:3138, absent-is-success, the same posture rebase's own deletes hold) and `list_recursive` (:3156). Plus a `turso-backup-gc` bin (oss/turso-backup/src/bin/gc.rs, wired in Cargo.toml) emitting one JSON line, mirroring turso-backup-snapshot's shape. DRY RUN UNLESS `GC_APPLY` is an explicit 1/true/yes — GC_APPLY=0, a typo and an empty string all leave the dry run in place, because the failure mode of guessing wrong points at deleted objects. `GC_GRACE_SECS` overrides the grace. Shape follows the existing `dedup::gc_dedup` (dedup.rs:472-573) tier-1b sweep rather than minting a parallel vocabulary. `rebase` still deletes nothing — the ticket's standing judgment that a recovery path must not delete is intact; the sweep is explicitly invoked.")
128//! @yah:handoff("LIVENESS WAS DERIVED BY READING THE RESTORE PATH, NOT ASSUMED, and the two are pinned together. Restore makes two selections and the GC's live set is their union: `restore_stream_from_manifests` (stream.rs:2094-2101) takes the chain's `base_snapshot_key`, and with no generations `hydrate::restore_subject` (hydrate.rs:522-528) falls back to `snapshot::restore_latest`, which takes the lexically-greatest key (snapshot.rs:432-444). Frames are live iff a manifest names them, resolved through `BackupTarget::frame_objects_of` — the same function restore's replay and `WalPuller` use, so both epoch key shapes and both batch/legacy layouts come along by construction. `gc_liveness_is_the_complement_of_restore_selection` asserts both branches against the real selection functions so they cannot drift. DELIBERATE: every manifest's base key is treated as live, not `validate_generation_chain`'s single answer — a chain spanning a WAL restart does not validate at all, and refusing to guess keeps both bases rather than deleting the one the next rebase is about to adopt.")
129//! @yah:handoff("TWO JUDGMENT CALLS, both recorded at the code site. (1) GRACE WINDOW, default 24h (stream.rs:2907), documenting the three live windows it must cover: a serialized cold restore in flight, a tail between uploading batches and writing their manifest, and rebase's publish-before-delete gap where the NEW base is reachable from nothing — that third case is where grace is the only thing preventing a concurrent GC from deleting a base a recovery just published. (2) DISCOVERED GAP, CLOSED: `base_snapshots_skipped` (:2937). A `snapshots/` prefix with no chain over it is indistinguishable from a plain tier-1a sink, whose older snapshots are HISTORY, not garbage — unguarded, the sweep would have pruned all but the newest. It now skips the snapshot half entirely when there are no generation manifests and reports that it did. Frames still go: a `frames/` prefix under a chainless sink is unreachable by construction.")
130//! @yah:verify("BASELINE MEASURED BEFORE THE FIRST EDIT, then re-measured, and INDEPENDENTLY RE-RUN BY THE LEADER (@Ashguard:eclipse) rather than taken on the courier's word. `cargo test --manifest-path oss/turso-backup/Cargo.toml` (turso-backup is its own workspace, excluded from the yah root workspace — a root-level `-p` invocation does not work): baseline 180 passed / 0 failed; after 192 passed / 0 failed. Leader's independent re-run of the combined tree, after R850-T2's own edit to oss/turso-backup/src/bin/hydrate.rs landed: 167 + 4 + 5 + 4 + 5 + 7 = 192 passed / 0 failed, exit 0. `cargo clippy --all-targets -- --deny=warnings` clean; `cargo doc --no-deps` warning locations byte-identical to before the edits.")
131//! @yah:verify("TWELVE NEW TESTS — all four the ticket named, plus four more that came out of the liveness derivation. `gc_collects_a_superseded_base_snapshot_and_keeps_the_current_one` · `gc_collects_orphaned_frame_prefixes_and_keeps_the_live_generation` (both epoch layouts) · `gc_spares_everything_inside_the_grace_window` · `gc_dry_run_reports_without_deleting` (and asserts the wet sweep executes the dry proposal key-for-key) · `gc_liveness_is_the_complement_of_restore_selection` · `gc_keeps_both_bases_while_a_rebase_is_half_landed` · `gc_without_a_chain_keeps_snapshots_but_still_collects_frames` · `gc_on_an_empty_prefix_collects_nothing`, plus 4 in bin/gc.rs. NO LIVE MinIO NEEDED: these use `object_store::memory::InMemory`, the same route the existing `dedup` GC tests take — no new infrastructure was stood up.")
132//! @yah:gotcha("THIS TICKET'S ANNOTATION WAS WRITTEN BY THE LEADER, NOT THE IMPLEMENTER, on purpose. R850-T3's annotation lives in oss/yubaba/crates/cloud/src/topology.rs (turso-backup is outside the board scanner's scan set), and @Ashguard:polaris held that file dirty with ~9 uncommitted writes for R850-T2 for the whole of this ticket's run. @Ashguard:blade was steered off `board.update` mid-turn and reported its account back to the leader instead, which then wrote it here once polaris was out. Content is blade's; the keystrokes are the leader's.")
133
134use std::collections::BTreeMap;
135
136use chrono::{DateTime, Utc};
137use serde::Serialize;
138
139use workload_spec::sovereign::SovereignRole;
140use workload_spec::{
141 Durability, DurabilityTier, LifecycleArchetype, RestartPolicy, VolumeSource, WorkloadSpec,
142};
143
144use crate::config::{CloudConfig, NodeAllocatable};
145use crate::migrate::{named_volume_path, VolumeDisposition};
146use crate::recovery_journal::{self, WorkloadRecovery};
147
148// ─── Measured constants the recovery estimate is built on ────────────────────
149
150/// Bulk object-store throughput, MB/s.
151///
152/// **Measured**, not modelled: R760-T10 on 2026-08-29, a real bulk range GET at
153/// 32.6 MB/s on one host — the 100 MB / 3.3 s figure in
154/// `oss/roadcase/docs/COST.md` §8. It is one measurement on one host against
155/// one backend, which is the whole of what this camp knows about the number;
156/// see [`RecoveryEstimate`] for how that limitation is carried outward rather
157/// than smoothed over.
158pub const MEASURED_HYDRATE_MB_PER_S: f64 = 32.6;
159
160/// Per-GET round-trip time, milliseconds. Measured in the same R760-T10 run.
161///
162/// Load-bearing because a tier-2 cold start issues its GETs **serially** —
163/// `turso_backup::stream` awaits one generation manifest and then one frame
164/// object at a time — so round trips, not bandwidth, are what dominates a
165/// restore of a small database with a long WAL history.
166pub const MEASURED_GET_RTT_MS: f64 = 16.0;
167
168/// `turso_backup::stream::DEFAULT_RPO_TARGET`, in seconds.
169///
170/// Duplicated rather than imported: `cloud` has no `turso-backup` dependency
171/// and should not grow one to print a number into a report. There is therefore
172/// **no test pinning the two together** — if this looks stale, the authority is
173/// `DEFAULT_RPO_TARGET` in `oss/turso-backup/src/stream.rs`, and a report that
174/// says "≤ 120s (turso-backup default)" is only as true as this line.
175pub const DEFAULT_STREAM_RPO_SECONDS: u32 = 120;
176
177// ─── The model ───────────────────────────────────────────────────────────────
178
179/// The declared fleet, read as a graph, plus every verdict derivable from it.
180///
181/// This is the single traversal. The Mermaid render ([`Topology::to_mermaid`])
182/// and the capacity report are *projections* of these fields — a diagram
183/// generated by its own second walk of the config would be free to disagree
184/// with the analysis printed above it, which is worse than no diagram.
185#[derive(Debug, Clone, Serialize)]
186pub struct Topology {
187 pub machines: Vec<MachineNode>,
188 pub workloads: Vec<WorkloadNode>,
189 /// One entry per declared machine: what is lost when that machine is.
190 pub node_losses: Vec<NodeLoss>,
191 /// Machines whose declared `[allocatable]` does not cover what is projected
192 /// onto them. Empty is the good case.
193 pub oversubscribed: Vec<Oversubscription>,
194}
195
196/// A declared machine, plus what the admission seam projects onto it.
197#[derive(Debug, Clone, Serialize)]
198pub struct MachineNode {
199 pub name: String,
200 pub region: Option<String>,
201 pub sovereign_group: Option<String>,
202 pub sovereign_role: SovereignRole,
203 pub taints: Vec<String>,
204 pub allocatable: Option<NodeAllocatable>,
205 /// Workloads whose projected placement lands here, in declaration order.
206 pub placed: Vec<String>,
207 /// Sum of `memory_request_mb()` over [`Self::placed`].
208 pub committed_memory_mb: u32,
209 /// Sum of `resources.cpu_millis` over [`Self::placed`].
210 pub committed_cpu_millis: u32,
211}
212
213/// Where the admission seam would put a workload, and what else could take it.
214///
215/// # Projected, not observed
216///
217/// This is `CloudConfig::admit_workload_candidates` run against the declared
218/// inventory. It differs from reality in two known ways, both of them the
219/// reason a *planning* surface wants the projection rather than a probe:
220///
221/// - A workload deployed with `--where=node:<name>` carries that pin in its
222/// spec's annotations and the analyzer sees it, but a workload deployed
223/// before the file was last edited is running against an older spec.
224/// - Admission has no liveness input, so `chosen` may be a box that is off.
225/// [`Self::alternates`] is the field that matters for survivability anyway.
226#[derive(Debug, Clone, Serialize)]
227#[serde(tag = "kind", rename_all = "snake_case")]
228pub enum Placement {
229 /// At least one machine admits this workload. `chosen` is the head of the
230 /// candidate pool in declaration order — the same first-fit
231 /// `admit_workload` returns.
232 Admitted {
233 chosen: String,
234 /// Every *other* admitting machine. **Empty is the survivability
235 /// finding**: a workload with no alternates has nowhere to go even if
236 /// its archetype would allow a move.
237 alternates: Vec<String>,
238 /// True when the spec names a node via the R833-F8 placement
239 /// annotation. A pin is still checked against capacity and taints, so
240 /// a pinned workload can still be `Unschedulable`.
241 pinned: bool,
242 },
243 /// Nothing in the declared fleet admits it. Carries admission's own
244 /// refusal, which names the pool it searched.
245 Unschedulable { reason: String },
246}
247
248impl Placement {
249 /// The machine this workload is projected onto, if any.
250 pub fn machine(&self) -> Option<&str> {
251 match self {
252 Self::Admitted { chosen, .. } => Some(chosen.as_str()),
253 Self::Unschedulable { .. } => None,
254 }
255 }
256
257 /// Machines that could take this workload if [`Self::machine`] were lost.
258 pub fn alternates(&self) -> &[String] {
259 match self {
260 Self::Admitted { alternates, .. } => alternates,
261 Self::Unschedulable { .. } => &[],
262 }
263 }
264}
265
266/// A workload's `yah.durability.*` declaration as the analyzer sees it.
267///
268/// [`Self::Undeclared`] and a declared [`DurabilityTier::None`] are separate
269/// variants on purpose — see `WorkloadSpec::durability`. The first is the
270/// shape that loses data by omission; the second is a decision.
271#[derive(Debug, Clone, Serialize)]
272#[serde(tag = "kind", rename_all = "snake_case")]
273pub enum DurabilityView {
274 /// Nobody said. For a workload with a named volume this is the finding.
275 Undeclared,
276 Declared(Durability),
277 /// The declaration exists and cannot be read. Surfaced rather than treated
278 /// as `Undeclared`, because "the operator tried and got it wrong" and "the
279 /// operator never considered it" call for different conversations.
280 Malformed {
281 reason: String,
282 },
283}
284
285/// One workload, everything about it that bears on survival, and where it goes.
286#[derive(Debug, Clone, Serialize)]
287pub struct WorkloadNode {
288 pub name: String,
289 pub tier: String,
290 pub mesh_identity: String,
291 pub archetype: LifecycleArchetype,
292 pub replicas: u32,
293 /// TOML-ish spelling of `restart_policy`, for report output.
294 pub restart_policy: String,
295 pub placement: Placement,
296 /// Reused wholesale from [`crate::migrate`]: the named/bind/tmpfs split is
297 /// the same classification a move needs, and minting a second vocabulary
298 /// for it would let the two answers drift.
299 pub volumes: Vec<VolumeDisposition>,
300 pub durability: DurabilityView,
301 /// Memory **request** (`memory_request_mb()`), not the cgroup ceiling.
302 pub memory_request_mb: u32,
303 pub cpu_millis: u32,
304 /// Public hostnames this workload fronts, if any.
305 pub public_hostnames: Vec<String>,
306 /// Mesh identities that must be `Ready` before this one starts.
307 pub depends_on: Vec<String>,
308}
309
310impl WorkloadNode {
311 /// Durable mounts — the ones whose bytes a node loss puts at risk.
312 /// Tmpfs is excluded by [`VolumeDisposition::is_durable`].
313 pub fn durable_volumes(&self) -> Vec<String> {
314 self.volumes
315 .iter()
316 .filter(|v| v.is_durable())
317 .map(|v| v.label())
318 .collect()
319 }
320
321 /// Whether any mount is a yubaba-managed named volume. This is the exact
322 /// shape whose only copy lives at `/var/lib/yah/kamaji/volumes/<name>` on
323 /// one box.
324 pub fn has_named_volume(&self) -> bool {
325 self.volumes
326 .iter()
327 .any(|v| matches!(v, VolumeDisposition::Copy { .. }))
328 }
329}
330
331// ─── Node loss ───────────────────────────────────────────────────────────────
332
333/// Everything that follows from losing one machine outright.
334#[derive(Debug, Clone, Serialize)]
335pub struct NodeLoss {
336 pub machine: String,
337 pub impacts: Vec<WorkloadImpact>,
338 /// Public hostnames served only from this machine.
339 pub public_endpoints_lost: Vec<String>,
340 pub quorum: QuorumEffect,
341}
342
343impl NodeLoss {
344 /// Impacts where bytes are gone for good. The headline of any report.
345 pub fn total_losses(&self) -> impl Iterator<Item = &WorkloadImpact> {
346 self.impacts
347 .iter()
348 .filter(|i| matches!(i.data_loss, DataLoss::Total { .. }))
349 }
350
351 /// Impacts that need a person. The second headline: an outage nobody is
352 /// paged for is an outage that lasts until someone notices.
353 pub fn needs_operator(&self) -> impl Iterator<Item = &WorkloadImpact> {
354 self.impacts.iter().filter(|i| i.outcome.is_manual())
355 }
356}
357
358/// What a node loss does to one workload projected onto it.
359#[derive(Debug, Clone, Serialize)]
360pub struct WorkloadImpact {
361 pub workload: String,
362 pub archetype: LifecycleArchetype,
363 pub outcome: Outcome,
364 pub data_loss: DataLoss,
365 pub recovery: RecoveryEstimate,
366}
367
368/// Whether the workload comes back, and who brings it back.
369#[derive(Debug, Clone, PartialEq, Serialize)]
370#[serde(tag = "kind", rename_all = "snake_case")]
371pub enum Outcome {
372 /// Fungible, and somewhere else in the declared fleet admits it. Note that
373 /// this says the *scheduler* could place it — it does not claim anything
374 /// automatically triggers that placement today.
375 Reschedulable { candidates: Vec<String> },
376
377 /// Fungible, but nothing else admits it. Down until the node is back or
378 /// the fleet grows. `reason` is the constraint that excludes everyone else
379 /// — a taint, the capacity floor, an arch mismatch.
380 NowhereToGo { reason: String },
381
382 /// [`LifecycleArchetype::Appliance`]: **pinned and non-drainable.**
383 ///
384 /// This is the variant the driving question lands on. `drain_workloads`
385 /// (`oss/yubaba/crates/yubaba/src/lib.rs`, R572-F4) skips appliances
386 /// outright, so there is no automatic move at any capacity — and the
387 /// operator verb that does move one, `yah cloud migrate`, *plans a
388 /// stop → copy → start* and expects the volume to already exist at the
389 /// destination. Against a node that is gone at the hardware level there is
390 /// nothing to copy from, which is why this variant does not promise
391 /// `migrate` will help.
392 PinnedAppliance {
393 /// Where a migrate could target, if anywhere admits it.
394 migrate_target: Option<String>,
395 /// True when the source volume is only reachable from the dead node,
396 /// i.e. `yah cloud migrate` has no source to copy from.
397 source_unreachable: bool,
398 },
399
400 /// `restart_policy = Never`: a run, not a service. Losing the node loses
401 /// the run; the answer is to run it again, not to fail it over.
402 RunLost,
403
404 /// `replicas > 1` and at least one other machine admits the workload, so
405 /// the survivors keep serving while the lost replica is replaced.
406 DegradedButServing { surviving_replicas: u32 },
407}
408
409impl Outcome {
410 /// Whether a human has to do something before this workload serves again.
411 pub fn is_manual(&self) -> bool {
412 matches!(
413 self,
414 Self::PinnedAppliance { .. } | Self::NowhereToGo { .. } | Self::RunLost
415 )
416 }
417
418 /// One line for a text report.
419 pub fn headline(&self) -> String {
420 match self {
421 Self::Reschedulable { candidates } => {
422 format!("automatic — schedulable onto {}", candidates.join(", "))
423 }
424 Self::NowhereToGo { reason } => {
425 format!("STAYS DOWN — nothing else admits it: {reason}")
426 }
427 Self::PinnedAppliance {
428 migrate_target,
429 source_unreachable,
430 } => {
431 let target = migrate_target.as_deref().unwrap_or("(nothing admits it)");
432 if *source_unreachable {
433 format!(
434 "OPERATOR — pinned appliance, never drained or rescheduled. \
435 `yah cloud migrate` would target {target}, but it plans a \
436 stop → copy → start and the copy has no source once the node \
437 is gone"
438 )
439 } else {
440 format!("OPERATOR — pinned appliance; `yah cloud migrate` to {target}")
441 }
442 }
443 Self::RunLost => "run lost — re-run it; nothing fails a job over".to_string(),
444 Self::DegradedButServing { surviving_replicas } => {
445 format!("degraded — {surviving_replicas} replica(s) still serving")
446 }
447 }
448 }
449}
450
451/// How much of the workload's state is gone, and how far back the copy is.
452#[derive(Debug, Clone, PartialEq, Serialize)]
453#[serde(tag = "kind", rename_all = "snake_case")]
454pub enum DataLoss {
455 /// Nothing durable is mounted.
456 None,
457 /// Only tmpfs. Discarded on stop by definition, so the node dying costs
458 /// nothing that a restart would not have.
459 EphemeralOnly,
460 /// Durable state exists and there is **no second copy anywhere**. The
461 /// bytes are gone with the node.
462 Total {
463 volumes: Vec<String>,
464 /// Why there is no copy: no declaration at all, or `tier = "none"`.
465 because: String,
466 },
467 /// A copy exists in an object store; the loss is the gap between the last
468 /// write and the last thing that reached the store.
469 Window {
470 tier: DurabilityTier,
471 store: String,
472 /// `None` for snapshot/dedup tiers, whose recovery point is set by
473 /// whatever schedules the snapshot and is therefore not in the spec.
474 rpo_seconds: Option<u32>,
475 /// Present when [`Self::rpo_seconds`] is the turso-backup default
476 /// rather than a declared value.
477 rpo_is_default: bool,
478 },
479 /// Bind mounts only. The camp did not create the host path and cannot know
480 /// whether the bytes exist elsewhere — the same refusal-to-guess
481 /// `migrate::preconditions` makes for the same mount kind.
482 Unknown { volumes: Vec<String> },
483}
484
485impl DataLoss {
486 /// One line for a text report.
487 pub fn headline(&self) -> String {
488 match self {
489 Self::None => "none — no durable state declared".to_string(),
490 Self::EphemeralOnly => "none — tmpfs only, discarded on stop anyway".to_string(),
491 Self::Total { volumes, because } => format!(
492 "TOTAL — {} has no second copy anywhere ({because})",
493 volumes.join(", ")
494 ),
495 Self::Window {
496 tier,
497 store,
498 rpo_seconds,
499 rpo_is_default,
500 } => match rpo_seconds {
501 Some(s) if *rpo_is_default => {
502 format!("≤ {s}s (turso-backup default, not declared) — tier {tier} → {store}")
503 }
504 Some(s) => format!("≤ {s}s (declared) — tier {tier} → {store}"),
505 None => format!(
506 "unbounded by the spec — tier {tier} → {store}; a snapshot tier's \
507 recovery point is set by whatever schedules it"
508 ),
509 },
510 Self::Unknown { volumes } => format!(
511 "UNKNOWN — {} are operator-managed bind mounts; the camp cannot say \
512 whether the bytes exist anywhere else",
513 volumes.join(", ")
514 ),
515 }
516 }
517}
518
519/// How long it takes to get the state back, and on what basis that is claimed.
520///
521/// [`Hydrate`](Self::Hydrate) is an **extrapolation from two measured
522/// constants** ([`MEASURED_HYDRATE_MB_PER_S`], [`MEASURED_GET_RTT_MS`]) applied
523/// to a **declared** state size — one measurement, on one host, against one
524/// backend, stretched over a number an operator typed. That is stated on the
525/// type rather than in a footnote because a recovery-time number without its
526/// provenance is the single easiest thing in a planning report to mistake for a
527/// measurement.
528///
529/// [`Measured`](Self::Measured) (R850-T2) is the exception and the thing to
530/// prefer: a restore that was actually timed, by
531/// `turso-backup-hydrate`, replayed out of `.yah/cloud/recovery.jsonl` by
532/// [`CloudConfig::load`] into [`CloudConfig::recovery_measurements`]. It carries
533/// its own provenance *and its age*, for the same reason — and it is never
534/// discarded for being old. A real restore from six weeks ago is a better
535/// answer than an extrapolation from an unrelated host; it is reported with
536/// "measured N days ago" attached so the reader can discount it themselves.
537#[derive(Debug, Clone, PartialEq, Serialize)]
538#[serde(tag = "kind", rename_all = "snake_case")]
539pub enum RecoveryEstimate {
540 /// Nothing to hydrate — the workload carries no durable state.
541 Immediate,
542 /// There is no copy to recover from. Recovery is not a duration.
543 NotRecoverable,
544 /// A copy exists but the spec does not say how big the state is, so the
545 /// transfer cannot be estimated. Names the annotation that would fix it.
546 UnknownStateSize { hint: &'static str },
547 /// Bulk transfer of a declared state size at the measured throughput.
548 Hydrate {
549 state_mb: u32,
550 seconds: f64,
551 /// Verbatim provenance, carried into JSON so a consumer cannot strip it.
552 basis: String,
553 },
554 /// A real restore of this workload, timed on a real host. Summed over the
555 /// workload's subjects, because the helper restores them one at a time and
556 /// a workload's recovery is all of them.
557 Measured {
558 /// Summed measured wall-clock seconds.
559 seconds: f64,
560 /// Summed bytes on disk after the restore. Not a declaration — this is
561 /// what actually landed.
562 bytes: u64,
563 /// When the oldest component of the sum was measured. A sum is only as
564 /// fresh as its stalest part.
565 measured_at: DateTime<Utc>,
566 /// Whole days from `measured_at` to when the journal was replayed.
567 /// Carried into JSON beside `basis` so a consumer can strip neither the
568 /// provenance nor the age.
569 age_days: i64,
570 /// The machine the most recent measurement was taken on. A restore time
571 /// is a property of a host as much as of a database.
572 node: String,
573 /// Verbatim provenance, carried into JSON so a consumer cannot strip it.
574 basis: String,
575 },
576}
577
578impl RecoveryEstimate {
579 /// One line for a text report.
580 pub fn headline(&self) -> String {
581 match self {
582 Self::Immediate => "immediate — stateless".to_string(),
583 Self::NotRecoverable => "n/a — nothing to recover from".to_string(),
584 Self::UnknownStateSize { hint } => {
585 format!("unknown — declare {hint} to get an estimate")
586 }
587 Self::Hydrate {
588 state_mb, seconds, ..
589 } => format!(
590 "≥ ~{seconds:.1}s to pull {state_mb} MiB (extrapolated from \
591 {MEASURED_HYDRATE_MB_PER_S} MB/s measured once, R760-T10; no restore was \
592 timed here, and WAL replay is on top)"
593 ),
594 Self::Measured {
595 seconds,
596 bytes,
597 age_days,
598 node,
599 ..
600 } => {
601 let mib = *bytes as f64 / (1024.0 * 1024.0);
602 let base = format!(
603 "{seconds:.1}s MEASURED — a real restore of {mib:.1} MiB timed on {node}, \
604 {age_days} day(s) ago"
605 );
606 if *age_days > recovery_journal::STALE_AFTER_DAYS {
607 format!(
608 "{base} — measured {age_days} days ago; declared state may have grown \
609 since, so treat it as a floor rather than a forecast"
610 )
611 } else {
612 base
613 }
614 }
615 }
616 }
617}
618
619/// What losing a machine does to its sovereign group's raft quorum.
620#[derive(Debug, Clone, PartialEq, Serialize)]
621#[serde(tag = "kind", rename_all = "snake_case")]
622pub enum QuorumEffect {
623 /// The machine declares no `sovereign_group`, so it votes in nothing.
624 NotInAGroup,
625 /// `sovereign_role = "non-voter"` — in the group's blast radius, holds no
626 /// seat. Its absence from `/raft/status` is correct, not drift.
627 NonVoter { group: String },
628 /// A voter is lost and the survivors still make a majority.
629 QuorumHolds {
630 group: String,
631 voters_before: usize,
632 voters_after: usize,
633 majority_needed: usize,
634 },
635 /// A voter is lost and the survivors do not. The group's raft stops
636 /// accepting writes, which includes cluster secrets — so workloads there
637 /// fail to resolve secrets even if their own containers are untouched.
638 QuorumLost {
639 group: String,
640 voters_before: usize,
641 voters_after: usize,
642 majority_needed: usize,
643 },
644}
645
646impl QuorumEffect {
647 /// One line for a text report.
648 pub fn headline(&self) -> String {
649 match self {
650 Self::NotInAGroup => "no sovereign group — votes in nothing".to_string(),
651 Self::NonVoter { group } => {
652 format!("non-voter in '{group}' — no quorum seat to lose")
653 }
654 Self::QuorumHolds {
655 group,
656 voters_after,
657 majority_needed,
658 ..
659 } => format!(
660 "'{group}' quorum holds — {voters_after} voter(s) left, {majority_needed} needed"
661 ),
662 Self::QuorumLost {
663 group,
664 voters_after,
665 majority_needed,
666 ..
667 } => format!(
668 "'{group}' LOSES QUORUM — {voters_after} voter(s) left, {majority_needed} \
669 needed; the group's raft stops accepting writes, and cluster secrets are \
670 read from the local raft replica, so workloads there fail to resolve \
671 secrets even where their containers are untouched"
672 ),
673 }
674 }
675}
676
677/// A machine whose declared capacity does not cover what is projected onto it.
678///
679/// Admission checks each workload against the node's `[allocatable]`
680/// *individually* (`RequiredSpec::matches`, R572-F5) and never subtracts what
681/// is already committed — so N workloads that each fit can all be admitted onto
682/// a node that cannot hold their sum. This is the arithmetic nothing in the
683/// tree does today.
684#[derive(Debug, Clone, Serialize)]
685pub struct Oversubscription {
686 pub machine: String,
687 pub allocatable: NodeAllocatable,
688 pub committed_memory_mb: u32,
689 pub committed_cpu_millis: u32,
690 pub workloads: Vec<String>,
691 pub memory_over: bool,
692 pub cpu_over: bool,
693}
694
695// ─── The traversal ───────────────────────────────────────────────────────────
696
697/// Walk the declared graph once and answer every question derivable from it.
698///
699/// Pure: same TOML in, same [`Topology`] out, no network. See the module header
700/// for what is deliberately outside the model.
701pub fn analyze(cfg: &CloudConfig) -> Topology {
702 let workloads: Vec<WorkloadNode> = cfg
703 .workloads
704 .iter()
705 .map(|w| workload_node(cfg, &w.spec))
706 .collect();
707
708 let mut machines: Vec<MachineNode> = cfg
709 .machines
710 .iter()
711 .map(|m| MachineNode {
712 name: m.name.clone(),
713 region: m.region.clone(),
714 sovereign_group: m.sovereign_group.clone(),
715 sovereign_role: m.sovereign_role.unwrap_or_default(),
716 taints: m.taints.clone(),
717 allocatable: m.allocatable.clone(),
718 placed: Vec::new(),
719 committed_memory_mb: 0,
720 committed_cpu_millis: 0,
721 })
722 .collect();
723
724 // Fold each workload's projected placement back onto its machine. Replicas
725 // multiply the commitment: `replicas = 3` asks the node for three copies of
726 // the request, and admission — which checks one workload against one node —
727 // never sees that multiplication.
728 let by_name: BTreeMap<String, usize> = machines
729 .iter()
730 .enumerate()
731 .map(|(i, m)| (m.name.clone(), i))
732 .collect();
733 for w in &workloads {
734 let Some(machine) = w.placement.machine() else {
735 continue;
736 };
737 let Some(&i) = by_name.get(machine) else {
738 continue;
739 };
740 let copies = w.replicas.max(1);
741 machines[i].placed.push(w.name.clone());
742 machines[i].committed_memory_mb = machines[i]
743 .committed_memory_mb
744 .saturating_add(w.memory_request_mb.saturating_mul(copies));
745 machines[i].committed_cpu_millis = machines[i]
746 .committed_cpu_millis
747 .saturating_add(w.cpu_millis.saturating_mul(copies));
748 }
749
750 let oversubscribed = machines.iter().filter_map(oversubscription).collect();
751 let node_losses = machines
752 .iter()
753 .map(|m| node_loss(cfg, &machines, &workloads, &m.name))
754 .collect();
755
756 Topology {
757 machines,
758 workloads,
759 node_losses,
760 oversubscribed,
761 }
762}
763
764fn workload_node(cfg: &CloudConfig, spec: &WorkloadSpec) -> WorkloadNode {
765 // One call into the admission seam — the same one `yah cloud apply` and
766 // `yah cloud migrate` use. Forking a second selector here would let the
767 // analyzer report a placement the fleet would never make.
768 let placement = match cfg.admit_workload_candidates(spec) {
769 Ok(candidates) => {
770 let mut names = candidates.iter().map(|m| m.name.clone());
771 let chosen = names
772 .next()
773 .expect("admit_workload_candidates never returns empty");
774 Placement::Admitted {
775 chosen,
776 alternates: names.collect(),
777 pinned: crate::config::node_selector_node(spec).is_some(),
778 }
779 }
780 Err(e) => Placement::Unschedulable {
781 reason: e.to_string(),
782 },
783 };
784
785 let volumes = spec
786 .volumes
787 .iter()
788 .map(|v| match &v.source {
789 VolumeSource::Named { name } => VolumeDisposition::Copy {
790 name: name.clone(),
791 host_path: named_volume_path(name),
792 mounted_at: v.target.clone(),
793 },
794 VolumeSource::Bind { host_path } => VolumeDisposition::Precondition {
795 host_path: host_path.clone(),
796 mounted_at: v.target.clone(),
797 },
798 VolumeSource::Tmpfs { size_mb } => VolumeDisposition::Discard {
799 mounted_at: v.target.clone(),
800 size_mb: *size_mb,
801 },
802 })
803 .collect();
804
805 let durability = match spec.durability() {
806 Ok(Some(d)) => DurabilityView::Declared(d.clone()),
807 Ok(None) => DurabilityView::Undeclared,
808 Err(e) => DurabilityView::Malformed {
809 reason: e.to_string(),
810 },
811 };
812
813 WorkloadNode {
814 name: spec.name.clone(),
815 tier: spec.tier.0.clone(),
816 mesh_identity: spec.fq_mesh_identity(),
817 archetype: spec.effective_archetype(),
818 replicas: spec.replicas,
819 restart_policy: restart_policy_label(&spec.restart_policy),
820 placement,
821 volumes,
822 durability,
823 memory_request_mb: spec.memory_request_mb(),
824 cpu_millis: spec.resources.cpu_millis,
825 public_hostnames: spec
826 .expose
827 .public
828 .iter()
829 .map(|p| p.hostname.clone())
830 .collect(),
831 depends_on: spec.depends_on.iter().map(|d| d.0.clone()).collect(),
832 }
833}
834
835fn restart_policy_label(p: &RestartPolicy) -> String {
836 match p {
837 RestartPolicy::Always => "always".to_string(),
838 RestartPolicy::OnFailure { max_attempts, .. } => {
839 format!("on-failure (max {max_attempts})")
840 }
841 RestartPolicy::Never => "never".to_string(),
842 }
843}
844
845fn oversubscription(m: &MachineNode) -> Option<Oversubscription> {
846 // No `[allocatable]` block is "unconstrained", exactly as admission reads
847 // it — not "zero capacity". Reporting an unbounded node as oversubscribed
848 // would flag every machine that has not been measured yet.
849 let alloc = m.allocatable.as_ref()?;
850 let memory_over = m.committed_memory_mb > alloc.memory_mb;
851 let cpu_over = alloc.cpu_millis > 0 && m.committed_cpu_millis > alloc.cpu_millis;
852 if !memory_over && !cpu_over {
853 return None;
854 }
855 Some(Oversubscription {
856 machine: m.name.clone(),
857 allocatable: alloc.clone(),
858 committed_memory_mb: m.committed_memory_mb,
859 committed_cpu_millis: m.committed_cpu_millis,
860 workloads: m.placed.clone(),
861 memory_over,
862 cpu_over,
863 })
864}
865
866fn node_loss(
867 cfg: &CloudConfig,
868 machines: &[MachineNode],
869 workloads: &[WorkloadNode],
870 dead: &str,
871) -> NodeLoss {
872 let impacts: Vec<WorkloadImpact> = workloads
873 .iter()
874 .filter(|w| w.placement.machine() == Some(dead))
875 .map(|w| workload_impact(w, dead, &cfg.recovery_measurements))
876 .collect();
877
878 let public_endpoints_lost = impacts
879 .iter()
880 .filter_map(|i| workloads.iter().find(|w| w.name == i.workload))
881 .flat_map(|w| w.public_hostnames.iter().cloned())
882 .collect();
883
884 NodeLoss {
885 machine: dead.to_string(),
886 impacts,
887 public_endpoints_lost,
888 quorum: quorum_effect(cfg, machines, dead),
889 }
890}
891
892fn workload_impact(
893 w: &WorkloadNode,
894 dead: &str,
895 measurements: &BTreeMap<String, WorkloadRecovery>,
896) -> WorkloadImpact {
897 let alternates: Vec<String> = w
898 .placement
899 .alternates()
900 .iter()
901 .filter(|m| m.as_str() != dead)
902 .cloned()
903 .collect();
904
905 let data_loss = data_loss(w);
906 let outcome = outcome(w, &alternates, dead);
907 let recovery = recovery(w, &data_loss, measurements);
908
909 WorkloadImpact {
910 workload: w.name.clone(),
911 archetype: w.archetype,
912 outcome,
913 data_loss,
914 recovery,
915 }
916}
917
918/// The core verdict. Archetype decides it, because archetype is what yubaba
919/// itself branches on — `drain_workloads` skips appliances (R572-F4), and
920/// `migrate` orders its steps by the same split.
921fn outcome(w: &WorkloadNode, alternates: &[String], dead: &str) -> Outcome {
922 match w.archetype {
923 LifecycleArchetype::Appliance => Outcome::PinnedAppliance {
924 migrate_target: alternates.first().cloned(),
925 // "Hardware-level kill" is the question being asked, so the source
926 // side of migrate's stop → copy → start has nothing to read from.
927 // A workload whose only durable mount is a named volume on the dead
928 // box is the exact shape with no source; one with no durable state
929 // has nothing to copy and so is not blocked on this.
930 source_unreachable: w.has_named_volume(),
931 },
932 LifecycleArchetype::Job => Outcome::RunLost,
933 LifecycleArchetype::Server => {
934 if alternates.is_empty() {
935 return Outcome::NowhereToGo {
936 reason: format!(
937 "{dead} is the only machine in the declared fleet that admits \
938 {} (archetype {}, {} MiB request, {} millicores)",
939 w.name,
940 w.archetype.taint_key(),
941 w.memory_request_mb,
942 w.cpu_millis,
943 ),
944 };
945 }
946 if w.replicas > 1 {
947 Outcome::DegradedButServing {
948 surviving_replicas: w.replicas - 1,
949 }
950 } else {
951 Outcome::Reschedulable {
952 candidates: alternates.to_vec(),
953 }
954 }
955 }
956 }
957}
958
959fn data_loss(w: &WorkloadNode) -> DataLoss {
960 let durable = w.durable_volumes();
961 if durable.is_empty() {
962 return if w.volumes.is_empty() {
963 DataLoss::None
964 } else {
965 DataLoss::EphemeralOnly
966 };
967 }
968
969 // A bind mount is operator-managed; the camp did not create the host path
970 // and has no basis for a claim about it either way. Only say "total" about
971 // volumes this camp is responsible for.
972 if !w.has_named_volume() {
973 return DataLoss::Unknown { volumes: durable };
974 }
975
976 match &w.durability {
977 DurabilityView::Undeclared => DataLoss::Total {
978 volumes: durable,
979 because: "no durability.tier declared, so the yubaba-managed named volume \
980 at /var/lib/yah/kamaji/volumes/ is the only copy"
981 .to_string(),
982 },
983 DurabilityView::Malformed { reason } => DataLoss::Total {
984 volumes: durable,
985 because: format!("the durability declaration cannot be read: {reason}"),
986 },
987 DurabilityView::Declared(d) => match d.tier {
988 DurabilityTier::None => DataLoss::Total {
989 volumes: durable,
990 because: "durability.tier = \"none\" — deliberately no second copy".to_string(),
991 },
992 DurabilityTier::Snapshot | DurabilityTier::Dedup => DataLoss::Window {
993 tier: d.tier,
994 store: d.store.clone().unwrap_or_default(),
995 rpo_seconds: None,
996 rpo_is_default: false,
997 },
998 DurabilityTier::Stream => DataLoss::Window {
999 tier: d.tier,
1000 store: d.store.clone().unwrap_or_default(),
1001 rpo_seconds: Some(d.rpo_seconds.unwrap_or(DEFAULT_STREAM_RPO_SECONDS)),
1002 rpo_is_default: d.rpo_seconds.is_none(),
1003 },
1004 },
1005 }
1006}
1007
1008/// R850-T2: `measurements` is [`CloudConfig::recovery_measurements`], already
1009/// replayed off the local tree by [`CloudConfig::load`]. It reaches here as
1010/// declared data, not as a path — `analyze` does no I/O, holds no clock, and
1011/// keeps the same contract [`crate::migrate::plan_migration`] holds.
1012fn recovery(
1013 w: &WorkloadNode,
1014 loss: &DataLoss,
1015 measurements: &BTreeMap<String, WorkloadRecovery>,
1016) -> RecoveryEstimate {
1017 match loss {
1018 DataLoss::None | DataLoss::EphemeralOnly => RecoveryEstimate::Immediate,
1019 DataLoss::Total { .. } | DataLoss::Unknown { .. } => RecoveryEstimate::NotRecoverable,
1020 DataLoss::Window { .. } => {
1021 // A timed restore beats both the extrapolation and the "we can't
1022 // say" — a measurement answers the question `state_mb` was only ever
1023 // a proxy for. Age does not disqualify it; see `RecoveryEstimate`.
1024 if let Some(m) = measurements.get(&w.name) {
1025 return measured_recovery(m);
1026 }
1027 let state_mb = match &w.durability {
1028 DurabilityView::Declared(Durability {
1029 state_mb: Some(mb), ..
1030 }) => *mb,
1031 _ => {
1032 return RecoveryEstimate::UnknownStateSize {
1033 hint: "durability.state_mb",
1034 }
1035 }
1036 };
1037 RecoveryEstimate::Hydrate {
1038 state_mb,
1039 seconds: f64::from(state_mb) / MEASURED_HYDRATE_MB_PER_S,
1040 // The *floor*, and it says so. A tier-2 restore also replays
1041 // WAL frames, and `turso_backup::stream` fetches those one at a
1042 // time — roadcase measured `2n + 1` serialized GETs for `n`
1043 // generations, which at this RTT reaches the same order as the
1044 // bulk transfer itself. `n` is not declared anywhere, so it is
1045 // named rather than guessed at.
1046 basis: format!(
1047 "bulk transfer at {MEASURED_HYDRATE_MB_PER_S} MB/s, measured R760-T10 \
1048 2026-08-29 on one host against one backend. A FLOOR: a tier-2 restore \
1049 adds 2n+1 serialized GETs at ~{MEASURED_GET_RTT_MS} ms each for n \
1050 generations, and n is not declared anywhere"
1051 ),
1052 }
1053 }
1054 }
1055}
1056
1057/// Turn one workload's replayed measurements into the estimate, provenance and
1058/// age included. Split out so the `basis` string lives next to the type that
1059/// justifies it rather than inside a `match` arm.
1060fn measured_recovery(m: &WorkloadRecovery) -> RecoveryEstimate {
1061 let measured_at = m.measured_at();
1062 RecoveryEstimate::Measured {
1063 seconds: m.seconds(),
1064 bytes: m.bytes(),
1065 measured_at,
1066 age_days: m.age_days(),
1067 node: m.node.clone(),
1068 basis: format!(
1069 "MEASURED, not extrapolated: {} subject(s) summed from a real restore recorded by \
1070 {} on {} at {}, replayed from `.yah/cloud/recovery.jsonl`. No \
1071 {MEASURED_HYDRATE_MB_PER_S} MB/s constant and no declared state-mb is involved. \
1072 Oldest component is {} day(s) old — the figure describes the state as it was then, \
1073 not as it is declared now.",
1074 m.subject_count(),
1075 m.helper(),
1076 m.node,
1077 measured_at.to_rfc3339(),
1078 m.age_days(),
1079 ),
1080 }
1081}
1082
1083fn quorum_effect(cfg: &CloudConfig, machines: &[MachineNode], dead: &str) -> QuorumEffect {
1084 let Some(m) = machines.iter().find(|m| m.name == dead) else {
1085 return QuorumEffect::NotInAGroup;
1086 };
1087 let Some(group) = m.sovereign_group.clone() else {
1088 return QuorumEffect::NotInAGroup;
1089 };
1090 if !m.sovereign_role.is_voter() {
1091 return QuorumEffect::NonVoter { group };
1092 }
1093
1094 let voters_before = cfg
1095 .machines_in_group(&group)
1096 .into_iter()
1097 .filter(|m| m.sovereign_role.unwrap_or_default().is_voter())
1098 .count();
1099 let voters_after = voters_before.saturating_sub(1);
1100 // Raft majority is over the *configured* membership, which the loss of a
1101 // box does not shrink — a dead voter still counts in the denominator until
1102 // someone removes it from the configuration.
1103 let majority_needed = voters_before / 2 + 1;
1104
1105 if voters_after >= majority_needed {
1106 QuorumEffect::QuorumHolds {
1107 group,
1108 voters_before,
1109 voters_after,
1110 majority_needed,
1111 }
1112 } else {
1113 QuorumEffect::QuorumLost {
1114 group,
1115 voters_before,
1116 voters_after,
1117 majority_needed,
1118 }
1119 }
1120}
1121
1122// ─── P2: renders ─────────────────────────────────────────────────────────────
1123
1124impl Topology {
1125 /// Render the model as a Mermaid `flowchart`.
1126 ///
1127 /// **A projection, not a second traversal.** Every node and edge below is
1128 /// read off fields [`analyze`] already computed, so the picture cannot
1129 /// disagree with the verdicts printed beside it. A diagram generated by its
1130 /// own walk of the config would be free to drift into decoration, which is
1131 /// worse than no diagram — it is the failure mode this method's shape
1132 /// exists to make impossible.
1133 ///
1134 /// The edge worth the whole render is `-.->|hydrate|`: the backup path is
1135 /// the one relationship in this graph that is invisible in the TOML, has no
1136 /// runtime today, and is exactly what decides whether a node loss is an
1137 /// incident or a restore.
1138 pub fn to_mermaid(&self) -> String {
1139 let mut out = String::from("flowchart TB\n");
1140
1141 // Machines, grouped by sovereign group. The grouping is the blast
1142 // radius (W305), so it is what a reader should see first.
1143 let mut groups: BTreeMap<Option<&str>, Vec<&MachineNode>> = BTreeMap::new();
1144 for m in &self.machines {
1145 groups
1146 .entry(m.sovereign_group.as_deref())
1147 .or_default()
1148 .push(m);
1149 }
1150 for (group, members) in &groups {
1151 let label = group.unwrap_or("ungrouped");
1152 out.push_str(&format!(
1153 " subgraph grp_{}[\"{label}\"]\n",
1154 sanitize(label)
1155 ));
1156 for m in members {
1157 // `sovereign_role` defaults to Voter, so a box that declares no
1158 // group would otherwise render as "voter" — a seat in a quorum
1159 // it is not in. Match what `QuorumEffect::NotInAGroup` says.
1160 let role = match (group, m.sovereign_role.is_voter()) {
1161 (None, _) => "no group",
1162 (Some(_), true) => "voter",
1163 (Some(_), false) => "non-voter",
1164 };
1165 let cap = match &m.allocatable {
1166 Some(a) => format!(
1167 "<br/>{}/{} MiB · {}/{} mCPU",
1168 m.committed_memory_mb, a.memory_mb, m.committed_cpu_millis, a.cpu_millis
1169 ),
1170 None => "<br/>no [allocatable] declared".to_string(),
1171 };
1172 out.push_str(&format!(
1173 " {}[\"{}<br/><i>{role}</i>{cap}\"]\n",
1174 node_id("m", &m.name),
1175 m.name
1176 ));
1177 }
1178 out.push_str(" end\n");
1179 }
1180
1181 // Workloads, their mounts, and their public front doors.
1182 for w in &self.workloads {
1183 let wid = node_id("w", &w.name);
1184 out.push_str(&format!(
1185 " {wid}(\"{}<br/><i>{}</i> · replicas {}\")\n",
1186 w.name,
1187 w.archetype.taint_key(),
1188 w.replicas
1189 ));
1190
1191 match &w.placement {
1192 Placement::Admitted { chosen, pinned, .. } => {
1193 let verb = if *pinned { "pinned" } else { "placed" };
1194 out.push_str(&format!(" {} -->|{verb}| {wid}\n", node_id("m", chosen)));
1195 }
1196 Placement::Unschedulable { .. } => {
1197 out.push_str(&format!(
1198 " unschedulable{{{{no node admits it}}}} --> {wid}\n"
1199 ));
1200 }
1201 }
1202
1203 for v in &w.volumes {
1204 let vid = node_id("v", &format!("{}-{}", w.name, v.label()));
1205 let (shape, edge) = match v {
1206 VolumeDisposition::Copy { name, .. } => {
1207 (format!("{vid}[(\"named: {name}\")]"), "mount")
1208 }
1209 VolumeDisposition::Precondition { host_path, .. } => (
1210 format!("{vid}[(\"bind: {}\")]", host_path.display()),
1211 "mount",
1212 ),
1213 VolumeDisposition::Discard { size_mb, .. } => {
1214 (format!("{vid}[(\"tmpfs {size_mb} MiB\")]"), "ephemeral")
1215 }
1216 };
1217 out.push_str(&format!(" {shape}\n {wid} -->|{edge}| {vid}\n"));
1218
1219 // The invisible edge. Only durable mounts can have one, and
1220 // only a declared tier draws it.
1221 if !v.is_durable() {
1222 continue;
1223 }
1224 if let DurabilityView::Declared(d) = &w.durability {
1225 if let Some(store) = &d.store {
1226 let sid = node_id("s", store);
1227 out.push_str(&format!(" {sid}[[\"{store}\"]]\n"));
1228 out.push_str(&format!(
1229 " {vid} -.->|backup: {}| {sid}\n {sid} -.->|hydrate| {vid}\n",
1230 d.tier
1231 ));
1232 }
1233 }
1234 }
1235
1236 for host in &w.public_hostnames {
1237 let hid = node_id("p", host);
1238 out.push_str(&format!(" {hid}>\"{host}\"]\n {hid} ==>|public| {wid}\n"));
1239 }
1240
1241 for dep in &w.depends_on {
1242 if let Some(target) = self.workloads.iter().find(|o| {
1243 o.mesh_identity == *dep || o.mesh_identity.ends_with(&format!("/{dep}"))
1244 }) {
1245 out.push_str(&format!(
1246 " {wid} -.->|mesh admit| {}\n",
1247 node_id("w", &target.name)
1248 ));
1249 }
1250 }
1251 }
1252
1253 out
1254 }
1255
1256 /// Render the survivability answer for one machine as plain text.
1257 ///
1258 /// Returns `None` when no machine by that name is declared — the caller
1259 /// owns the wording of that refusal, since it has the declared list.
1260 pub fn render_node_loss(&self, machine: &str) -> Option<String> {
1261 let loss = self.node_losses.iter().find(|l| l.machine == machine)?;
1262 let mut out = format!("If {machine} is lost at the hardware level:\n\n");
1263 out.push_str(&format!(" quorum: {}\n", loss.quorum.headline()));
1264 if loss.public_endpoints_lost.is_empty() {
1265 out.push_str(" public endpoints lost: none\n");
1266 } else {
1267 out.push_str(&format!(
1268 " public endpoints lost: {}\n",
1269 loss.public_endpoints_lost.join(", ")
1270 ));
1271 }
1272
1273 if loss.impacts.is_empty() {
1274 out.push_str("\n No declared workload is projected onto this machine.\n");
1275 return Some(out);
1276 }
1277
1278 for i in &loss.impacts {
1279 out.push_str(&format!(
1280 "\n {} ({})\n",
1281 i.workload,
1282 i.archetype.taint_key()
1283 ));
1284 out.push_str(&format!(" what happens: {}\n", i.outcome.headline()));
1285 out.push_str(&format!(" data loss: {}\n", i.data_loss.headline()));
1286 out.push_str(&format!(" recovery: {}\n", i.recovery.headline()));
1287 }
1288 Some(out)
1289 }
1290
1291 /// Render the capacity arithmetic — the whole fleet, oversubscription
1292 /// called out rather than left to the reader to spot.
1293 pub fn render_capacity(&self) -> String {
1294 let mut out = String::from("Declared capacity vs projected commitment:\n\n");
1295 for m in &self.machines {
1296 let placed = if m.placed.is_empty() {
1297 "(nothing)".to_string()
1298 } else {
1299 m.placed.join(", ")
1300 };
1301 match &m.allocatable {
1302 Some(a) => out.push_str(&format!(
1303 " {:<16} {:>6}/{:<6} MiB {:>6}/{:<6} mCPU {placed}\n",
1304 m.name,
1305 m.committed_memory_mb,
1306 a.memory_mb,
1307 m.committed_cpu_millis,
1308 a.cpu_millis
1309 )),
1310 None => out.push_str(&format!(
1311 " {:<16} {:>6}/{:<6} MiB {:>6}/{:<6} mCPU {placed}\n",
1312 m.name, m.committed_memory_mb, "?", m.committed_cpu_millis, "?"
1313 )),
1314 }
1315 }
1316 if self.oversubscribed.is_empty() {
1317 out.push_str("\nNo machine is oversubscribed against its declared [allocatable].\n");
1318 } else {
1319 out.push_str(
1320 "\nOVERSUBSCRIBED — admission checks each workload against a node \
1321 individually and never subtracts what is already committed (R572-F5), \
1322 so these all admitted and cannot all run:\n",
1323 );
1324 for o in &self.oversubscribed {
1325 let mut axes = Vec::new();
1326 if o.memory_over {
1327 axes.push(format!(
1328 "memory {} MiB > {} MiB",
1329 o.committed_memory_mb, o.allocatable.memory_mb
1330 ));
1331 }
1332 if o.cpu_over {
1333 axes.push(format!(
1334 "cpu {} > {} millicores",
1335 o.committed_cpu_millis, o.allocatable.cpu_millis
1336 ));
1337 }
1338 out.push_str(&format!(
1339 " {}: {} — {}\n",
1340 o.machine,
1341 axes.join(", "),
1342 o.workloads.join(", ")
1343 ));
1344 }
1345 }
1346 out
1347 }
1348
1349 /// Render every machine's loss, plus the capacity view — the default
1350 /// whole-fleet report.
1351 pub fn render(&self) -> String {
1352 let mut out = String::new();
1353 for m in &self.machines {
1354 if let Some(section) = self.render_node_loss(&m.name) {
1355 out.push_str(§ion);
1356 out.push('\n');
1357 }
1358 }
1359 out.push_str(&self.render_capacity());
1360 out
1361 }
1362}
1363
1364/// Mermaid node ids must be identifier-ish; machine and workload names are DNS
1365/// labels and hostnames, and buckets are URLs. One prefix per kind keeps two
1366/// different things that sanitize to the same string apart.
1367fn node_id(prefix: &str, name: &str) -> String {
1368 format!("{prefix}_{}", sanitize(name))
1369}
1370
1371fn sanitize(name: &str) -> String {
1372 name.chars()
1373 .map(|c| if c.is_ascii_alphanumeric() { c } else { '_' })
1374 .collect()
1375}
1376
1377#[cfg(test)]
1378mod tests {
1379 use super::*;
1380 use std::path::Path;
1381 use tempfile::{tempdir, TempDir};
1382
1383 /// Fixtures go through `CloudConfig::load`, never a struct literal, for the
1384 /// reason `migrate`'s own fixture records: a hand-built spec can express a
1385 /// shape no operator could write, and `load_workloads` shape-validates —
1386 /// which since R850-P4 includes the durability declaration. A test that
1387 /// skips the loader would pass on a spec the fleet would reject.
1388 struct Camp {
1389 dir: TempDir,
1390 }
1391
1392 impl Camp {
1393 fn new() -> Self {
1394 Self {
1395 dir: tempdir().unwrap(),
1396 }
1397 }
1398
1399 fn root(&self) -> &Path {
1400 self.dir.path()
1401 }
1402
1403 fn machine(self, name: &str, extra: &str) -> Self {
1404 let dir = self.root().join(".yah/infra/machines");
1405 std::fs::create_dir_all(&dir).unwrap();
1406 std::fs::write(
1407 dir.join(format!("{name}.toml")),
1408 format!(
1409 "name = \"{name}\"\nprovider = \"static\"\nmesh_tags = []\n\
1410 {extra}\n\
1411 [connect]\naddress = \"10.0.0.1\"\nssh = \"root@{name}\"\n\
1412 identity_file = \"~/.ssh/yah\"\n\
1413 yubaba = \"http://{name}:7443\"\n"
1414 ),
1415 )
1416 .unwrap();
1417 self
1418 }
1419
1420 /// A machine with a capacity budget, which is what makes the
1421 /// oversubscription and NowhereToGo cases expressible.
1422 fn sized(self, name: &str, memory_mb: u32, cpu_millis: u32, extra: &str) -> Self {
1423 self.machine(
1424 name,
1425 &format!(
1426 "{extra}\n[allocatable]\nmemory_mb = {memory_mb}\ncpu_millis = {cpu_millis}"
1427 ),
1428 )
1429 }
1430
1431 fn workload(self, name: &str, extra: &str) -> Self {
1432 self.workload_sized(name, 128, 100, extra)
1433 }
1434
1435 fn workload_sized(self, name: &str, memory_mb: u32, cpu_millis: u32, extra: &str) -> Self {
1436 let dir = self.root().join(".yah/infra/workloads");
1437 std::fs::create_dir_all(&dir).unwrap();
1438 std::fs::write(
1439 dir.join(format!("{name}.toml")),
1440 format!(
1441 "schema_version = 1\nname = \"{name}\"\ntier = \"infra\"\n\
1442 restart_policy = \"always\"\n\
1443 {extra}\n\
1444 [image]\nregistry = \"cr.yah.dev\"\nrepository = \"{name}\"\n\
1445 tag = \"v1\"\ndigest = \"sha256:abc\"\n\
1446 [resources]\nmemory_mb = {memory_mb}\ncpu_millis = {cpu_millis}\n\
1447 [stop_policy]\nsignal = 15\ngrace_period = 10000\n\
1448 [expose.mesh]\nidentity = \"{name}\"\nports = [8080]\nallow_from = []\n"
1449 ),
1450 )
1451 .unwrap();
1452 self
1453 }
1454
1455 /// Append a measured subject restore to `.yah/cloud/recovery.jsonl`
1456 /// (R850-T2), dated `days_ago` before now. Goes through the real
1457 /// journal writer so `analyze` sees exactly what `yah cloud topology
1458 /// --record-hydrate` would have left behind.
1459 fn measured(
1460 self,
1461 workload: &str,
1462 subject: &str,
1463 bytes: u64,
1464 seconds: f64,
1465 days_ago: i64,
1466 ) -> Self {
1467 crate::recovery_journal::RecoveryJournal::at_workspace(self.root())
1468 .append(&[crate::recovery_journal::RecoveryRecord {
1469 at: Utc::now() - chrono::Duration::days(days_ago),
1470 workload: workload.to_string(),
1471 node: "a".to_string(),
1472 tier: Some("stream".to_string()),
1473 subject: subject.to_string(),
1474 bytes,
1475 seconds,
1476 helper: crate::recovery_journal::HELPER_TURSO_BACKUP_HYDRATE.to_string(),
1477 }])
1478 .unwrap();
1479 self
1480 }
1481
1482 fn analyze(&self) -> Topology {
1483 super::analyze(&CloudConfig::load(self.root()).expect("fixture camp must load"))
1484 }
1485 }
1486
1487 /// The `[annotations]` block for a stream-tier workload, with `state-mb`
1488 /// declared unless `state_mb` is `None` — the two shapes that decide
1489 /// between `Hydrate` and `UnknownStateSize` when nothing was measured.
1490 fn stream_durability(state_mb: Option<u32>) -> String {
1491 let mut s = "[durability]\n\
1492 tier = \"stream\"\n\
1493 engine = \"turso\"\n\
1494 store = \"s3://backups/db\"\n\
1495 subjects = [\"accounts.db\"]\n"
1496 .to_string();
1497 if let Some(mb) = state_mb {
1498 s.push_str(&format!("state_mb = {mb}\n"));
1499 }
1500 s
1501 }
1502
1503 const NAMED_VOLUME: &str = "[[volumes]]\nsource = { named = { name = \"accounts\" } }\n\
1504 target = \"/var/lib/app\"\nread_only = false\n";
1505
1506 fn impact<'a>(topo: &'a Topology, machine: &str, workload: &str) -> &'a WorkloadImpact {
1507 topo.node_losses
1508 .iter()
1509 .find(|l| l.machine == machine)
1510 .unwrap_or_else(|| panic!("no node_loss for {machine}"))
1511 .impacts
1512 .iter()
1513 .find(|i| i.workload == workload)
1514 .unwrap_or_else(|| panic!("{workload} is not projected onto {machine}"))
1515 }
1516
1517 // ── The driving question ─────────────────────────────────────────────────
1518
1519 /// R850's acceptance test, and the reason the relay exists.
1520 ///
1521 /// The noisetable-account shape verbatim: one replica, `archetype =
1522 /// "appliance"`, a yubaba-managed named volume, no durability tier
1523 /// declared. Hardware-kill the node. The analyzer has to say *out loud*
1524 /// that nothing reschedules, that the bytes are gone, and that even the
1525 /// operator verb has no source to copy from — because before this module
1526 /// existed, learning any of those three meant reading yubaba's source.
1527 #[test]
1528 fn hardware_killing_the_node_under_a_singleton_appliance_loses_everything() {
1529 let topo = Camp::new()
1530 .sized("us-west-001", 12288, 6000, "")
1531 .sized("us-west-003", 16384, 16000, "")
1532 .workload(
1533 "noisetable-account",
1534 &format!("replicas = 1\narchetype = \"appliance\"\n{NAMED_VOLUME}"),
1535 )
1536 .analyze();
1537
1538 let i = impact(&topo, "us-west-001", "noisetable-account");
1539
1540 // 1. Nothing reschedules — and the reason is the archetype, not a
1541 // shortage of nodes. us-west-003 admits it fine and is irrelevant.
1542 let Outcome::PinnedAppliance {
1543 migrate_target,
1544 source_unreachable,
1545 } = &i.outcome
1546 else {
1547 panic!("expected PinnedAppliance, got {:?}", i.outcome);
1548 };
1549 assert_eq!(migrate_target.as_deref(), Some("us-west-003"));
1550
1551 // 2. `yah cloud migrate` is named as the human step AND disclaimed:
1552 // it plans a copy, and a dead box has nothing to copy from.
1553 assert!(source_unreachable);
1554 let headline = i.outcome.headline();
1555 assert!(headline.contains("yah cloud migrate"), "{headline}");
1556 assert!(headline.contains("no source"), "{headline}");
1557
1558 // 3. Every account, passkey and session is gone.
1559 let DataLoss::Total { volumes, because } = &i.data_loss else {
1560 panic!("expected Total, got {:?}", i.data_loss);
1561 };
1562 assert_eq!(volumes, &["accounts".to_string()]);
1563 assert!(because.contains("durability.tier"), "{because}");
1564 assert_eq!(i.recovery, RecoveryEstimate::NotRecoverable);
1565
1566 // 4. And the operator-facing text says all three without a source dive.
1567 let text = topo.render_node_loss("us-west-001").unwrap();
1568 assert!(text.contains("OPERATOR"), "{text}");
1569 assert!(text.contains("TOTAL"), "{text}");
1570 assert!(text.contains("nothing to recover from"), "{text}");
1571 }
1572
1573 /// The same spec's verdict must not soften when the appliance is the only
1574 /// thing the fleet could hold — `migrate_target` goes to `None` and the
1575 /// outcome stays manual rather than degrading into "nowhere to go", which
1576 /// would read as a capacity problem instead of an archetype one.
1577 #[test]
1578 fn an_appliance_with_no_alternate_still_reads_as_pinned_not_as_capacity() {
1579 let topo = Camp::new()
1580 .sized("solo", 1024, 1000, "")
1581 .workload(
1582 "db",
1583 &format!("replicas = 1\narchetype = \"appliance\"\n{NAMED_VOLUME}"),
1584 )
1585 .analyze();
1586
1587 let i = impact(&topo, "solo", "db");
1588 let Outcome::PinnedAppliance { migrate_target, .. } = &i.outcome else {
1589 panic!("expected PinnedAppliance, got {:?}", i.outcome);
1590 };
1591 assert_eq!(*migrate_target, None);
1592 assert!(i.outcome.is_manual());
1593 }
1594
1595 // ── The fungible half ────────────────────────────────────────────────────
1596
1597 #[test]
1598 fn a_stateless_server_with_another_admitting_node_is_automatic() {
1599 let topo = Camp::new()
1600 .sized("a", 4096, 4000, "")
1601 .sized("b", 4096, 4000, "")
1602 .workload("api", "replicas = 1\narchetype = \"server\"")
1603 .analyze();
1604
1605 let i = impact(&topo, "a", "api");
1606 assert_eq!(
1607 i.outcome,
1608 Outcome::Reschedulable {
1609 candidates: vec!["b".to_string()]
1610 }
1611 );
1612 assert_eq!(i.data_loss, DataLoss::None);
1613 assert_eq!(i.recovery, RecoveryEstimate::Immediate);
1614 assert!(!i.outcome.is_manual());
1615 }
1616
1617 /// The finding that is invisible in the TOML: a workload can be perfectly
1618 /// stateless and restartable and still be a single point of failure,
1619 /// because only one declared box clears its floor.
1620 #[test]
1621 fn a_stateless_server_that_only_one_node_admits_stays_down() {
1622 let topo = Camp::new()
1623 .sized("big", 8192, 8000, "")
1624 .sized("small", 512, 8000, "")
1625 .workload_sized("hungry", 4096, 100, "replicas = 1\narchetype = \"server\"")
1626 .analyze();
1627
1628 let i = impact(&topo, "big", "hungry");
1629 let Outcome::NowhereToGo { reason } = &i.outcome else {
1630 panic!("expected NowhereToGo, got {:?}", i.outcome);
1631 };
1632 assert!(reason.contains("big"), "{reason}");
1633 assert!(i.outcome.is_manual());
1634 }
1635
1636 #[test]
1637 fn replicas_above_one_degrade_rather_than_fail() {
1638 let topo = Camp::new()
1639 .sized("a", 4096, 4000, "")
1640 .sized("b", 4096, 4000, "")
1641 .workload("api", "replicas = 3\narchetype = \"server\"")
1642 .analyze();
1643
1644 assert_eq!(
1645 impact(&topo, "a", "api").outcome,
1646 Outcome::DegradedButServing {
1647 surviving_replicas: 2
1648 }
1649 );
1650 }
1651
1652 /// A `no-appliance` taint is absolute — there is no toleration anywhere in
1653 /// the tree (W305 finding 2) — so a tainted box must never show up as an
1654 /// alternate an appliance could fail over to.
1655 #[test]
1656 fn a_no_appliance_taint_removes_a_node_from_the_alternates() {
1657 let topo = Camp::new()
1658 .sized("prod", 8192, 8000, "")
1659 .sized("builder", 16384, 16000, "taints = [\"no-appliance\"]")
1660 .workload(
1661 "db",
1662 &format!("replicas = 1\narchetype = \"appliance\"\n{NAMED_VOLUME}"),
1663 )
1664 .analyze();
1665
1666 let db = topo.workloads.iter().find(|w| w.name == "db").unwrap();
1667 assert_eq!(db.placement.machine(), Some("prod"));
1668 assert!(
1669 db.placement.alternates().is_empty(),
1670 "builder carries no-appliance and must not be offered: {:?}",
1671 db.placement
1672 );
1673 }
1674
1675 #[test]
1676 fn a_job_is_a_lost_run_not_a_failover() {
1677 let topo = Camp::new()
1678 .sized("a", 4096, 4000, "")
1679 .workload("build", "replicas = 1\narchetype = \"job\"")
1680 .analyze();
1681 assert_eq!(impact(&topo, "a", "build").outcome, Outcome::RunLost);
1682 }
1683
1684 // ── Durability ───────────────────────────────────────────────────────────
1685
1686 #[test]
1687 fn a_declared_stream_tier_turns_total_loss_into_a_bounded_window() {
1688 let topo = Camp::new()
1689 .sized("a", 4096, 4000, "")
1690 .workload(
1691 "db",
1692 &format!(
1693 "replicas = 1\narchetype = \"appliance\"\n{NAMED_VOLUME}\n\
1694 [durability]\n\
1695 tier = \"stream\"\n\
1696 engine = \"turso\"\n\
1697 store = \"s3://backups/db\"\n\
1698 subjects = [\"accounts.db\"]\n\
1699 rpo_seconds = 30\n\
1700 state_mb = 100\n"
1701 ),
1702 )
1703 .analyze();
1704
1705 let i = impact(&topo, "a", "db");
1706 assert_eq!(
1707 i.data_loss,
1708 DataLoss::Window {
1709 tier: DurabilityTier::Stream,
1710 store: "s3://backups/db".to_string(),
1711 rpo_seconds: Some(30),
1712 rpo_is_default: false,
1713 }
1714 );
1715
1716 // The recovery figure is roadcase's measured constant applied to the
1717 // declared size — 100 MiB at 32.6 MB/s ≈ 3.1 s, the same order as the
1718 // 3.3 s R760-T10 actually measured for 100 MB.
1719 let RecoveryEstimate::Hydrate {
1720 state_mb,
1721 seconds,
1722 basis,
1723 } = &i.recovery
1724 else {
1725 panic!("expected Hydrate, got {:?}", i.recovery);
1726 };
1727 assert_eq!(*state_mb, 100);
1728 assert!((*seconds - 3.067).abs() < 0.01, "{seconds}");
1729 // Provenance survives into the structured output, not just the text.
1730 assert!(basis.contains("R760-T10"), "{basis}");
1731 assert!(basis.contains("FLOOR"), "{basis}");
1732 }
1733
1734 // ── Measured recovery (R850-T2) ──────────────────────────────────────────
1735
1736 /// The baseline this feature must not disturb: a camp that has never timed
1737 /// a restore has no `.yah/cloud/recovery.jsonl`, that is not an error, and
1738 /// both pre-existing answers come back exactly as before.
1739 #[test]
1740 fn an_absent_recovery_journal_leaves_todays_answers_untouched() {
1741 for (state_mb, expected) in [
1742 (
1743 Some(100),
1744 RecoveryEstimate::Hydrate {
1745 state_mb: 100,
1746 seconds: 100.0 / MEASURED_HYDRATE_MB_PER_S,
1747 // Compared field-wise below; `basis` is long and pinned by
1748 // `a_declared_stream_tier_turns_total_loss_into_a_bounded_window`.
1749 basis: String::new(),
1750 },
1751 ),
1752 (
1753 None,
1754 RecoveryEstimate::UnknownStateSize {
1755 hint: "durability.state_mb",
1756 },
1757 ),
1758 ] {
1759 let camp = Camp::new().sized("a", 4096, 4000, "").workload(
1760 "db",
1761 &format!(
1762 "replicas = 1\narchetype = \"appliance\"\n{NAMED_VOLUME}\n{}",
1763 stream_durability(state_mb)
1764 ),
1765 );
1766 assert!(
1767 !crate::paths::recovery_journal(camp.root()).exists(),
1768 "fixture must have no journal"
1769 );
1770
1771 let topo = camp.analyze();
1772 let got = &impact(&topo, "a", "db").recovery;
1773 match (&expected, got) {
1774 (
1775 RecoveryEstimate::Hydrate {
1776 state_mb: want_mb,
1777 seconds: want_s,
1778 ..
1779 },
1780 RecoveryEstimate::Hydrate {
1781 state_mb, seconds, ..
1782 },
1783 ) => {
1784 assert_eq!(state_mb, want_mb);
1785 assert!((seconds - want_s).abs() < 1e-9, "{seconds}");
1786 }
1787 (a, b) => assert_eq!(a, b),
1788 }
1789 }
1790 }
1791
1792 /// The point of the ticket: one journal entry turns the extrapolation into
1793 /// a measurement, and the figure is the sum over the workload's subjects.
1794 #[test]
1795 fn a_journalled_restore_replaces_the_extrapolation_with_the_measured_sum() {
1796 let topo = Camp::new()
1797 .sized("a", 4096, 4000, "")
1798 .workload(
1799 "db",
1800 &format!(
1801 "replicas = 1\narchetype = \"appliance\"\n{NAMED_VOLUME}\n{}",
1802 stream_durability(Some(100))
1803 ),
1804 )
1805 // 100 MiB declared would extrapolate to ~3.07s; the real restore
1806 // took 21s across two subjects, and that is what must be reported.
1807 .measured("db", "accounts.db", 1024 * 1024 * 60, 12.5, 2)
1808 .measured("db", "ledger.db", 1024 * 1024 * 40, 8.5, 2)
1809 .analyze();
1810
1811 let i = impact(&topo, "a", "db");
1812 let RecoveryEstimate::Measured {
1813 seconds,
1814 bytes,
1815 age_days,
1816 node,
1817 basis,
1818 ..
1819 } = &i.recovery
1820 else {
1821 panic!("expected Measured, got {:?}", i.recovery);
1822 };
1823 assert!((*seconds - 21.0).abs() < 1e-9, "{seconds}");
1824 assert_eq!(*bytes, 1024 * 1024 * 100);
1825 assert_eq!(*age_days, 2);
1826 assert_eq!(node, "a");
1827 // Provenance survives into the structured output, not just the text —
1828 // and it says which of the two kinds of number this is.
1829 assert!(basis.contains("MEASURED, not extrapolated"), "{basis}");
1830 assert!(basis.contains("recovery.jsonl"), "{basis}");
1831 assert!(basis.contains("turso-backup-hydrate"), "{basis}");
1832 // The extrapolation's constant must not appear as if it were involved.
1833 assert!(!basis.contains("R760-T10"), "{basis}");
1834
1835 let headline = i.recovery.headline();
1836 assert!(headline.contains("MEASURED"), "{headline}");
1837 assert!(headline.contains("21.0s"), "{headline}");
1838 assert!(!headline.contains("extrapolated"), "{headline}");
1839 }
1840
1841 /// A measurement is never discarded for being old: an extrapolation from
1842 /// one unrelated host is not an improvement on a real restore. The age goes
1843 /// into the headline instead, so the reader discounts it themselves.
1844 #[test]
1845 fn a_stale_measurement_is_still_reported_and_says_so() {
1846 let topo = Camp::new()
1847 .sized("a", 4096, 4000, "")
1848 .workload(
1849 "db",
1850 &format!(
1851 "replicas = 1\narchetype = \"appliance\"\n{NAMED_VOLUME}\n{}",
1852 stream_durability(Some(100))
1853 ),
1854 )
1855 .measured("db", "accounts.db", 1024 * 1024 * 100, 41.0, 90)
1856 .analyze();
1857
1858 let i = impact(&topo, "a", "db");
1859 let RecoveryEstimate::Measured {
1860 seconds, age_days, ..
1861 } = &i.recovery
1862 else {
1863 panic!("stale must stay Measured, not fall back to Hydrate: {:?}", i.recovery);
1864 };
1865 assert!((*seconds - 41.0).abs() < 1e-9, "{seconds}");
1866 assert_eq!(*age_days, 90);
1867 assert!(*age_days > recovery_journal::STALE_AFTER_DAYS);
1868
1869 let headline = i.recovery.headline();
1870 assert!(headline.contains("measured 90 days ago"), "{headline}");
1871 assert!(headline.contains("may have grown since"), "{headline}");
1872 }
1873
1874 /// A measurement answers the question `state-mb` was only ever a proxy for,
1875 /// so it beats `UnknownStateSize` too — an undeclared size is no longer a
1876 /// reason to refuse an answer once a real restore has been timed.
1877 #[test]
1878 fn a_measurement_beats_an_undeclared_state_size() {
1879 let topo = Camp::new()
1880 .sized("a", 4096, 4000, "")
1881 .workload(
1882 "db",
1883 &format!(
1884 "replicas = 1\narchetype = \"appliance\"\n{NAMED_VOLUME}\n{}",
1885 stream_durability(None)
1886 ),
1887 )
1888 .measured("db", "accounts.db", 4096, 1.5, 0)
1889 .analyze();
1890
1891 assert!(matches!(
1892 impact(&topo, "a", "db").recovery,
1893 RecoveryEstimate::Measured { .. }
1894 ));
1895 }
1896
1897 /// Attribution is per workload. A measurement filed against one workload
1898 /// must not leak into another's estimate, and must not turn a stateless
1899 /// workload into a recoverable one.
1900 #[test]
1901 fn a_measurement_for_one_workload_does_not_leak_into_another() {
1902 let topo = Camp::new()
1903 .sized("a", 4096, 4000, "")
1904 .workload(
1905 "db",
1906 &format!(
1907 "replicas = 1\narchetype = \"appliance\"\n{NAMED_VOLUME}\n{}",
1908 stream_durability(Some(100))
1909 ),
1910 )
1911 .workload("web", "replicas = 1\narchetype = \"server\"")
1912 .measured("web", "accounts.db", 4096, 99.0, 1)
1913 .analyze();
1914
1915 // `db` has its own Window and no measurement of its own.
1916 assert!(matches!(
1917 impact(&topo, "a", "db").recovery,
1918 RecoveryEstimate::Hydrate { .. }
1919 ));
1920 // `web` is stateless; a stray measurement does not invent state for it.
1921 assert_eq!(
1922 impact(&topo, "a", "web").recovery,
1923 RecoveryEstimate::Immediate
1924 );
1925 }
1926
1927 /// An undeclared RPO on a stream tier is the turso-backup default, and the
1928 /// report has to say which of the two it is printing — a number the
1929 /// operator believes they chose is worse than no number.
1930 #[test]
1931 fn an_undeclared_rpo_is_labelled_as_the_default_not_as_a_choice() {
1932 let topo = Camp::new()
1933 .sized("a", 4096, 4000, "")
1934 .workload(
1935 "db",
1936 &format!(
1937 "replicas = 1\narchetype = \"appliance\"\n{NAMED_VOLUME}\n\
1938 [durability]\n\
1939 tier = \"stream\"\n\
1940 engine = \"turso\"\n\
1941 store = \"s3://backups/db\"\n\
1942 subjects = [\"accounts.db\"]\n"
1943 ),
1944 )
1945 .analyze();
1946
1947 let i = impact(&topo, "a", "db");
1948 let DataLoss::Window {
1949 rpo_seconds,
1950 rpo_is_default,
1951 ..
1952 } = &i.data_loss
1953 else {
1954 panic!("expected Window, got {:?}", i.data_loss);
1955 };
1956 assert_eq!(*rpo_seconds, Some(DEFAULT_STREAM_RPO_SECONDS));
1957 assert!(rpo_is_default);
1958 assert!(i.data_loss.headline().contains("not declared"));
1959
1960 // No declared state size ⇒ no invented duration.
1961 assert_eq!(
1962 i.recovery,
1963 RecoveryEstimate::UnknownStateSize {
1964 hint: "durability.state_mb"
1965 }
1966 );
1967 }
1968
1969 /// A snapshot tier has a copy but no recovery point the spec can state,
1970 /// and saying "≤ 120s" for it would be a fabrication.
1971 #[test]
1972 fn a_snapshot_tier_reports_an_unbounded_window_rather_than_borrowing_the_stream_rpo() {
1973 let topo = Camp::new()
1974 .sized("a", 4096, 4000, "")
1975 .workload(
1976 "db",
1977 &format!(
1978 "replicas = 1\narchetype = \"appliance\"\n{NAMED_VOLUME}\n\
1979 [durability]\n\
1980 tier = \"snapshot\"\n\
1981 engine = \"turso\"\n\
1982 store = \"s3://backups/db\"\n\
1983 subjects = [\"accounts.db\"]\n"
1984 ),
1985 )
1986 .analyze();
1987
1988 let i = impact(&topo, "a", "db");
1989 let DataLoss::Window { rpo_seconds, .. } = &i.data_loss else {
1990 panic!("expected Window, got {:?}", i.data_loss);
1991 };
1992 assert_eq!(*rpo_seconds, None);
1993 assert!(i.data_loss.headline().contains("unbounded by the spec"));
1994 }
1995
1996 /// `tier = "none"` still loses everything — but as a decision, and the
1997 /// report must not read the same as the case where nobody looked.
1998 #[test]
1999 fn a_deliberate_none_tier_reads_differently_from_an_undeclared_one() {
2000 let camp = Camp::new().sized("a", 4096, 4000, "").workload(
2001 "cache",
2002 &format!(
2003 "replicas = 1\narchetype = \"appliance\"\n{NAMED_VOLUME}\n\
2004 [durability]\ntier = \"none\"\n"
2005 ),
2006 );
2007 let topo = camp.analyze();
2008 let DataLoss::Total { because, .. } = &impact(&topo, "a", "cache").data_loss else {
2009 panic!("expected Total");
2010 };
2011 assert!(because.contains("deliberately"), "{because}");
2012 assert!(
2013 !because.contains("no durability.tier declared"),
2014 "{because}"
2015 );
2016 }
2017
2018 /// A bind mount is operator-managed. The camp did not create the host path
2019 /// and has no basis for saying the bytes are gone — the same refusal to
2020 /// guess `migrate::preconditions` makes about the same mount kind.
2021 #[test]
2022 fn a_bind_mount_is_unknown_rather_than_total() {
2023 let topo = Camp::new()
2024 .sized("a", 4096, 4000, "")
2025 .workload(
2026 "svc",
2027 "replicas = 1\narchetype = \"appliance\"\n\
2028 [[volumes]]\nsource = { bind = { host_path = \"/srv/data\" } }\n\
2029 target = \"/data\"\nread_only = false\n",
2030 )
2031 .analyze();
2032
2033 let i = impact(&topo, "a", "svc");
2034 assert!(matches!(i.data_loss, DataLoss::Unknown { .. }));
2035 assert!(i.data_loss.headline().contains("UNKNOWN"));
2036 }
2037
2038 #[test]
2039 fn tmpfs_only_is_not_a_data_loss() {
2040 let topo = Camp::new()
2041 .sized("a", 4096, 4000, "")
2042 .workload(
2043 "svc",
2044 "replicas = 1\narchetype = \"server\"\n\
2045 [[volumes]]\nsource = { tmpfs = { size_mb = 64 } }\n\
2046 target = \"/scratch\"\nread_only = false\n",
2047 )
2048 .analyze();
2049 assert_eq!(impact(&topo, "a", "svc").data_loss, DataLoss::EphemeralOnly);
2050 }
2051
2052 // ── Capacity (P3) ────────────────────────────────────────────────────────
2053
2054 /// The gap admission structurally cannot see: `RequiredSpec::matches`
2055 /// checks one workload against one node's `[allocatable]` and never
2056 /// subtracts what is already committed, so three workloads that each fit
2057 /// are each admitted onto a node that cannot hold their sum.
2058 #[test]
2059 fn workloads_that_each_fit_can_still_oversubscribe_the_node_they_all_land_on() {
2060 let topo = Camp::new()
2061 .sized("small", 1024, 4000, "")
2062 .workload_sized("a", 512, 100, "replicas = 1\narchetype = \"server\"")
2063 .workload_sized("b", 512, 100, "replicas = 1\narchetype = \"server\"")
2064 .workload_sized("c", 512, 100, "replicas = 1\narchetype = \"server\"")
2065 .analyze();
2066
2067 assert_eq!(topo.oversubscribed.len(), 1);
2068 let o = &topo.oversubscribed[0];
2069 assert_eq!(o.machine, "small");
2070 assert_eq!(o.committed_memory_mb, 1536);
2071 assert!(o.memory_over);
2072 assert!(!o.cpu_over);
2073 assert!(topo.render_capacity().contains("OVERSUBSCRIBED"));
2074 }
2075
2076 /// Replicas multiply the ask. Admission checks one copy against one node
2077 /// and never multiplies, so `replicas = 4` of a fitting workload is exactly
2078 /// the shape that admits cleanly and cannot run.
2079 #[test]
2080 fn replicas_multiply_the_commitment() {
2081 let topo = Camp::new()
2082 .sized("small", 1024, 4000, "")
2083 .workload_sized("api", 512, 100, "replicas = 4\narchetype = \"server\"")
2084 .analyze();
2085
2086 assert_eq!(topo.machines[0].committed_memory_mb, 2048);
2087 assert_eq!(topo.oversubscribed.len(), 1);
2088 }
2089
2090 /// A node with no `[allocatable]` is *unconstrained*, exactly as admission
2091 /// reads it — not a zero-capacity node. Flagging it would light up every
2092 /// machine nobody has measured yet and train the operator to ignore this.
2093 #[test]
2094 fn a_node_with_no_allocatable_block_is_never_reported_oversubscribed() {
2095 let topo = Camp::new()
2096 .machine("unmeasured", "")
2097 .workload_sized("api", 99999, 99999, "replicas = 1\narchetype = \"server\"")
2098 .analyze();
2099
2100 assert!(topo.oversubscribed.is_empty());
2101 assert!(topo.render_capacity().contains("unmeasured"));
2102 }
2103
2104 // ── Quorum ───────────────────────────────────────────────────────────────
2105
2106 #[test]
2107 fn losing_one_of_three_voters_keeps_quorum() {
2108 let topo = Camp::new()
2109 .machine(
2110 "a",
2111 "sovereign_group = \"prod\"\nsovereign_role = \"voter\"",
2112 )
2113 .machine(
2114 "b",
2115 "sovereign_group = \"prod\"\nsovereign_role = \"voter\"",
2116 )
2117 .machine(
2118 "c",
2119 "sovereign_group = \"prod\"\nsovereign_role = \"voter\"",
2120 )
2121 .analyze();
2122
2123 let q = &topo.node_losses[0].quorum;
2124 assert_eq!(
2125 *q,
2126 QuorumEffect::QuorumHolds {
2127 group: "prod".to_string(),
2128 voters_before: 3,
2129 voters_after: 2,
2130 majority_needed: 2,
2131 }
2132 );
2133 }
2134
2135 /// The second half of a node loss that is easy to miss: cluster secrets are
2136 /// read from the *local raft replica*, so a group that loses quorum fails
2137 /// workloads whose containers were never touched.
2138 #[test]
2139 fn losing_one_of_two_voters_loses_quorum_and_says_what_that_costs() {
2140 let topo = Camp::new()
2141 .machine(
2142 "a",
2143 "sovereign_group = \"prod\"\nsovereign_role = \"voter\"",
2144 )
2145 .machine(
2146 "b",
2147 "sovereign_group = \"prod\"\nsovereign_role = \"voter\"",
2148 )
2149 .analyze();
2150
2151 let q = &topo.node_losses[0].quorum;
2152 assert!(matches!(q, QuorumEffect::QuorumLost { .. }), "{q:?}");
2153 let headline = q.headline();
2154 assert!(headline.contains("LOSES QUORUM"), "{headline}");
2155 assert!(headline.contains("secrets"), "{headline}");
2156 }
2157
2158 /// A non-voter's absence from `/raft/status` is correct, not drift
2159 /// (`workload_spec::sovereign`), so losing one costs no seat.
2160 #[test]
2161 fn a_non_voter_has_no_seat_to_lose() {
2162 let topo = Camp::new()
2163 .machine(
2164 "a",
2165 "sovereign_group = \"prod\"\nsovereign_role = \"voter\"",
2166 )
2167 .machine(
2168 "w",
2169 "sovereign_group = \"prod\"\nsovereign_role = \"non-voter\"",
2170 )
2171 .analyze();
2172
2173 let q = &topo
2174 .node_losses
2175 .iter()
2176 .find(|l| l.machine == "w")
2177 .unwrap()
2178 .quorum;
2179 assert_eq!(
2180 *q,
2181 QuorumEffect::NonVoter {
2182 group: "prod".to_string()
2183 }
2184 );
2185 }
2186
2187 // ── P2: the render is a projection ───────────────────────────────────────
2188
2189 /// The scope trap this relay was warned about: a diagram produced by its
2190 /// own walk of the config can disagree with the analysis printed beside it.
2191 /// This pins that it cannot — the placement edge in the Mermaid output is
2192 /// the placement the model computed, and the hydrate edge exists exactly
2193 /// when the model found a durability tier.
2194 #[test]
2195 fn the_mermaid_render_agrees_with_the_model_it_projects() {
2196 let topo = Camp::new()
2197 .sized("us-west-001", 12288, 6000, "sovereign_group = \"prod\"")
2198 .sized("us-west-003", 16384, 16000, "taints = [\"no-appliance\"]")
2199 .workload(
2200 "backed-up",
2201 &format!(
2202 "replicas = 1\narchetype = \"appliance\"\n{NAMED_VOLUME}\n\
2203 [durability]\n\
2204 tier = \"stream\"\n\
2205 engine = \"turso\"\n\
2206 store = \"s3://backups/db\"\n\
2207 subjects = [\"accounts.db\"]\n"
2208 ),
2209 )
2210 .analyze();
2211
2212 let mermaid = topo.to_mermaid();
2213 let w = topo
2214 .workloads
2215 .iter()
2216 .find(|w| w.name == "backed-up")
2217 .unwrap();
2218
2219 // The edge names the machine the model chose, not a re-derived one.
2220 assert_eq!(w.placement.machine(), Some("us-west-001"));
2221 assert!(
2222 mermaid.contains("m_us_west_001 -->|placed| w_backed_up"),
2223 "{mermaid}"
2224 );
2225 // The blast radius is what a reader sees first.
2226 assert!(mermaid.contains("subgraph grp_prod[\"prod\"]"), "{mermaid}");
2227 // The edge that is invisible in the TOML.
2228 assert!(mermaid.contains("|hydrate|"), "{mermaid}");
2229 assert!(mermaid.contains("s3://backups/db"), "{mermaid}");
2230 // Committed-vs-allocatable rides the same numbers as render_capacity.
2231 assert!(mermaid.contains("128/12288 MiB"), "{mermaid}");
2232 }
2233
2234 /// The negative half: no declared tier means no hydrate edge. A diagram
2235 /// that drew one anyway would show a safety property that does not exist,
2236 /// which is the single worst thing this render could do.
2237 #[test]
2238 fn no_declared_tier_draws_no_hydrate_edge() {
2239 let topo = Camp::new()
2240 .sized("a", 4096, 4000, "")
2241 .workload(
2242 "db",
2243 &format!("replicas = 1\narchetype = \"appliance\"\n{NAMED_VOLUME}"),
2244 )
2245 .analyze();
2246
2247 let mermaid = topo.to_mermaid();
2248 assert!(mermaid.contains("v_db_accounts"), "{mermaid}");
2249 assert!(!mermaid.contains("hydrate"), "{mermaid}");
2250 }
2251
2252 #[test]
2253 fn a_public_hostname_is_an_edge_and_is_reported_lost_with_its_node() {
2254 let topo = Camp::new()
2255 .sized("a", 4096, 4000, "")
2256 .workload(
2257 "web",
2258 "replicas = 1\narchetype = \"server\"\n\
2259 [expose.public]\nhostname = \"app.example.com\"\nport = 8080\n\
2260 tls = \"cf_managed\"\n",
2261 )
2262 .analyze();
2263
2264 assert_eq!(
2265 topo.node_losses[0].public_endpoints_lost,
2266 vec!["app.example.com".to_string()]
2267 );
2268 assert!(topo.to_mermaid().contains("|public|"));
2269 assert!(topo
2270 .render_node_loss("a")
2271 .unwrap()
2272 .contains("public endpoints lost: app.example.com"));
2273 }
2274
2275 /// `sovereign_role` defaults to `Voter`, so a box that declares no group
2276 /// renders as one unless the label is suppressed — a quorum seat in a
2277 /// quorum that does not exist. Caught against the live camp, where
2278 /// us-west-002 and us-west-015 declare no group.
2279 #[test]
2280 fn an_ungrouped_machine_is_not_labelled_a_voter() {
2281 let mermaid = Camp::new()
2282 .sized("loner", 1024, 1000, "")
2283 .analyze()
2284 .to_mermaid();
2285 assert!(mermaid.contains("<i>no group</i>"), "{mermaid}");
2286 assert!(!mermaid.contains("<i>voter</i>"), "{mermaid}");
2287 }
2288
2289 #[test]
2290 fn an_undeclared_machine_has_no_node_loss_section() {
2291 let topo = Camp::new().sized("a", 4096, 4000, "").analyze();
2292 assert!(topo.render_node_loss("typo").is_none());
2293 assert!(topo.render_node_loss("a").is_some());
2294 }
2295}