Skip to main content

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:at(2026-09-05T09:34:40Z)
47//! @yah:assignee(agent:bundle-anthropic-ashguard)
48//! @yah:phase(P4b)
49//! @yah:parent(R850)
50//! @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.")
51//! @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.")
52//! @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.")
53//! @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.")
54//! @yah:verify("`yah cloud topology --kill <node>` for a workload declaring a real tier must stop reporting `RecoveryEstimate` as an extrapolation and start reporting a measured restore. The analyzer is already the acceptance surface — oss/yubaba/crates/cloud/src/topology.rs, test `hardware_killing_the_node_under_a_singleton_appliance_loses_everything`.")
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
74use std::collections::BTreeMap;
75
76use serde::Serialize;
77
78use workload_spec::sovereign::SovereignRole;
79use workload_spec::{
80    Durability, DurabilityTier, LifecycleArchetype, RestartPolicy, VolumeSource, WorkloadSpec,
81};
82
83use crate::config::{CloudConfig, NodeAllocatable};
84use crate::migrate::{named_volume_path, VolumeDisposition};
85
86// ─── Measured constants the recovery estimate is built on ────────────────────
87
88/// Bulk object-store throughput, MB/s.
89///
90/// **Measured**, not modelled: R760-T10 on 2026-08-29, a real bulk range GET at
91/// 32.6 MB/s on one host — the 100 MB / 3.3 s figure in
92/// `oss/roadcase/docs/COST.md` §8. It is one measurement on one host against
93/// one backend, which is the whole of what this camp knows about the number;
94/// see [`RecoveryEstimate`] for how that limitation is carried outward rather
95/// than smoothed over.
96pub const MEASURED_HYDRATE_MB_PER_S: f64 = 32.6;
97
98/// Per-GET round-trip time, milliseconds. Measured in the same R760-T10 run.
99///
100/// Load-bearing because a tier-2 cold start issues its GETs **serially** —
101/// `turso_backup::stream` awaits one generation manifest and then one frame
102/// object at a time — so round trips, not bandwidth, are what dominates a
103/// restore of a small database with a long WAL history.
104pub const MEASURED_GET_RTT_MS: f64 = 16.0;
105
106/// `turso_backup::stream::DEFAULT_RPO_TARGET`, in seconds.
107///
108/// Duplicated rather than imported: `cloud` has no `turso-backup` dependency
109/// and should not grow one to print a number into a report. There is therefore
110/// **no test pinning the two together** — if this looks stale, the authority is
111/// `DEFAULT_RPO_TARGET` in `oss/turso-backup/src/stream.rs`, and a report that
112/// says "≤ 120s (turso-backup default)" is only as true as this line.
113pub const DEFAULT_STREAM_RPO_SECONDS: u32 = 120;
114
115// ─── The model ───────────────────────────────────────────────────────────────
116
117/// The declared fleet, read as a graph, plus every verdict derivable from it.
118///
119/// This is the single traversal. The Mermaid render ([`Topology::to_mermaid`])
120/// and the capacity report are *projections* of these fields — a diagram
121/// generated by its own second walk of the config would be free to disagree
122/// with the analysis printed above it, which is worse than no diagram.
123#[derive(Debug, Clone, Serialize)]
124pub struct Topology {
125    pub machines: Vec<MachineNode>,
126    pub workloads: Vec<WorkloadNode>,
127    /// One entry per declared machine: what is lost when that machine is.
128    pub node_losses: Vec<NodeLoss>,
129    /// Machines whose declared `[allocatable]` does not cover what is projected
130    /// onto them. Empty is the good case.
131    pub oversubscribed: Vec<Oversubscription>,
132}
133
134/// A declared machine, plus what the admission seam projects onto it.
135#[derive(Debug, Clone, Serialize)]
136pub struct MachineNode {
137    pub name: String,
138    pub region: Option<String>,
139    pub sovereign_group: Option<String>,
140    pub sovereign_role: SovereignRole,
141    pub taints: Vec<String>,
142    pub allocatable: Option<NodeAllocatable>,
143    /// Workloads whose projected placement lands here, in declaration order.
144    pub placed: Vec<String>,
145    /// Sum of `memory_request_mb()` over [`Self::placed`].
146    pub committed_memory_mb: u32,
147    /// Sum of `resources.cpu_millis` over [`Self::placed`].
148    pub committed_cpu_millis: u32,
149}
150
151/// Where the admission seam would put a workload, and what else could take it.
152///
153/// # Projected, not observed
154///
155/// This is `CloudConfig::admit_workload_candidates` run against the declared
156/// inventory. It differs from reality in two known ways, both of them the
157/// reason a *planning* surface wants the projection rather than a probe:
158///
159/// - A workload deployed with `--where=node:<name>` carries that pin in its
160///   spec's annotations and the analyzer sees it, but a workload deployed
161///   before the file was last edited is running against an older spec.
162/// - Admission has no liveness input, so `chosen` may be a box that is off.
163///   [`Self::alternates`] is the field that matters for survivability anyway.
164#[derive(Debug, Clone, Serialize)]
165#[serde(tag = "kind", rename_all = "snake_case")]
166pub enum Placement {
167    /// At least one machine admits this workload. `chosen` is the head of the
168    /// candidate pool in declaration order — the same first-fit
169    /// `admit_workload` returns.
170    Admitted {
171        chosen: String,
172        /// Every *other* admitting machine. **Empty is the survivability
173        /// finding**: a workload with no alternates has nowhere to go even if
174        /// its archetype would allow a move.
175        alternates: Vec<String>,
176        /// True when the spec names a node via the R833-F8 placement
177        /// annotation. A pin is still checked against capacity and taints, so
178        /// a pinned workload can still be `Unschedulable`.
179        pinned: bool,
180    },
181    /// Nothing in the declared fleet admits it. Carries admission's own
182    /// refusal, which names the pool it searched.
183    Unschedulable { reason: String },
184}
185
186impl Placement {
187    /// The machine this workload is projected onto, if any.
188    pub fn machine(&self) -> Option<&str> {
189        match self {
190            Self::Admitted { chosen, .. } => Some(chosen.as_str()),
191            Self::Unschedulable { .. } => None,
192        }
193    }
194
195    /// Machines that could take this workload if [`Self::machine`] were lost.
196    pub fn alternates(&self) -> &[String] {
197        match self {
198            Self::Admitted { alternates, .. } => alternates,
199            Self::Unschedulable { .. } => &[],
200        }
201    }
202}
203
204/// A workload's `yah.durability.*` declaration as the analyzer sees it.
205///
206/// [`Self::Undeclared`] and a declared [`DurabilityTier::None`] are separate
207/// variants on purpose — see `WorkloadSpec::durability`. The first is the
208/// shape that loses data by omission; the second is a decision.
209#[derive(Debug, Clone, Serialize)]
210#[serde(tag = "kind", rename_all = "snake_case")]
211pub enum DurabilityView {
212    /// Nobody said. For a workload with a named volume this is the finding.
213    Undeclared,
214    Declared(Durability),
215    /// The declaration exists and cannot be read. Surfaced rather than treated
216    /// as `Undeclared`, because "the operator tried and got it wrong" and "the
217    /// operator never considered it" call for different conversations.
218    Malformed {
219        reason: String,
220    },
221}
222
223/// One workload, everything about it that bears on survival, and where it goes.
224#[derive(Debug, Clone, Serialize)]
225pub struct WorkloadNode {
226    pub name: String,
227    pub tier: String,
228    pub mesh_identity: String,
229    pub archetype: LifecycleArchetype,
230    pub replicas: u32,
231    /// TOML-ish spelling of `restart_policy`, for report output.
232    pub restart_policy: String,
233    pub placement: Placement,
234    /// Reused wholesale from [`crate::migrate`]: the named/bind/tmpfs split is
235    /// the same classification a move needs, and minting a second vocabulary
236    /// for it would let the two answers drift.
237    pub volumes: Vec<VolumeDisposition>,
238    pub durability: DurabilityView,
239    /// Memory **request** (`memory_request_mb()`), not the cgroup ceiling.
240    pub memory_request_mb: u32,
241    pub cpu_millis: u32,
242    /// Public hostnames this workload fronts, if any.
243    pub public_hostnames: Vec<String>,
244    /// Mesh identities that must be `Ready` before this one starts.
245    pub depends_on: Vec<String>,
246}
247
248impl WorkloadNode {
249    /// Durable mounts — the ones whose bytes a node loss puts at risk.
250    /// Tmpfs is excluded by [`VolumeDisposition::is_durable`].
251    pub fn durable_volumes(&self) -> Vec<String> {
252        self.volumes
253            .iter()
254            .filter(|v| v.is_durable())
255            .map(|v| v.label())
256            .collect()
257    }
258
259    /// Whether any mount is a yubaba-managed named volume. This is the exact
260    /// shape whose only copy lives at `/var/lib/yah/kamaji/volumes/<name>` on
261    /// one box.
262    pub fn has_named_volume(&self) -> bool {
263        self.volumes
264            .iter()
265            .any(|v| matches!(v, VolumeDisposition::Copy { .. }))
266    }
267}
268
269// ─── Node loss ───────────────────────────────────────────────────────────────
270
271/// Everything that follows from losing one machine outright.
272#[derive(Debug, Clone, Serialize)]
273pub struct NodeLoss {
274    pub machine: String,
275    pub impacts: Vec<WorkloadImpact>,
276    /// Public hostnames served only from this machine.
277    pub public_endpoints_lost: Vec<String>,
278    pub quorum: QuorumEffect,
279}
280
281impl NodeLoss {
282    /// Impacts where bytes are gone for good. The headline of any report.
283    pub fn total_losses(&self) -> impl Iterator<Item = &WorkloadImpact> {
284        self.impacts
285            .iter()
286            .filter(|i| matches!(i.data_loss, DataLoss::Total { .. }))
287    }
288
289    /// Impacts that need a person. The second headline: an outage nobody is
290    /// paged for is an outage that lasts until someone notices.
291    pub fn needs_operator(&self) -> impl Iterator<Item = &WorkloadImpact> {
292        self.impacts.iter().filter(|i| i.outcome.is_manual())
293    }
294}
295
296/// What a node loss does to one workload projected onto it.
297#[derive(Debug, Clone, Serialize)]
298pub struct WorkloadImpact {
299    pub workload: String,
300    pub archetype: LifecycleArchetype,
301    pub outcome: Outcome,
302    pub data_loss: DataLoss,
303    pub recovery: RecoveryEstimate,
304}
305
306/// Whether the workload comes back, and who brings it back.
307#[derive(Debug, Clone, PartialEq, Serialize)]
308#[serde(tag = "kind", rename_all = "snake_case")]
309pub enum Outcome {
310    /// Fungible, and somewhere else in the declared fleet admits it. Note that
311    /// this says the *scheduler* could place it — it does not claim anything
312    /// automatically triggers that placement today.
313    Reschedulable { candidates: Vec<String> },
314
315    /// Fungible, but nothing else admits it. Down until the node is back or
316    /// the fleet grows. `reason` is the constraint that excludes everyone else
317    /// — a taint, the capacity floor, an arch mismatch.
318    NowhereToGo { reason: String },
319
320    /// [`LifecycleArchetype::Appliance`]: **pinned and non-drainable.**
321    ///
322    /// This is the variant the driving question lands on. `drain_workloads`
323    /// (`oss/yubaba/crates/yubaba/src/lib.rs`, R572-F4) skips appliances
324    /// outright, so there is no automatic move at any capacity — and the
325    /// operator verb that does move one, `yah cloud migrate`, *plans a
326    /// stop → copy → start* and expects the volume to already exist at the
327    /// destination. Against a node that is gone at the hardware level there is
328    /// nothing to copy from, which is why this variant does not promise
329    /// `migrate` will help.
330    PinnedAppliance {
331        /// Where a migrate could target, if anywhere admits it.
332        migrate_target: Option<String>,
333        /// True when the source volume is only reachable from the dead node,
334        /// i.e. `yah cloud migrate` has no source to copy from.
335        source_unreachable: bool,
336    },
337
338    /// `restart_policy = Never`: a run, not a service. Losing the node loses
339    /// the run; the answer is to run it again, not to fail it over.
340    RunLost,
341
342    /// `replicas > 1` and at least one other machine admits the workload, so
343    /// the survivors keep serving while the lost replica is replaced.
344    DegradedButServing { surviving_replicas: u32 },
345}
346
347impl Outcome {
348    /// Whether a human has to do something before this workload serves again.
349    pub fn is_manual(&self) -> bool {
350        matches!(
351            self,
352            Self::PinnedAppliance { .. } | Self::NowhereToGo { .. } | Self::RunLost
353        )
354    }
355
356    /// One line for a text report.
357    pub fn headline(&self) -> String {
358        match self {
359            Self::Reschedulable { candidates } => {
360                format!("automatic — schedulable onto {}", candidates.join(", "))
361            }
362            Self::NowhereToGo { reason } => {
363                format!("STAYS DOWN — nothing else admits it: {reason}")
364            }
365            Self::PinnedAppliance {
366                migrate_target,
367                source_unreachable,
368            } => {
369                let target = migrate_target.as_deref().unwrap_or("(nothing admits it)");
370                if *source_unreachable {
371                    format!(
372                        "OPERATOR — pinned appliance, never drained or rescheduled. \
373                         `yah cloud migrate` would target {target}, but it plans a \
374                         stop → copy → start and the copy has no source once the node \
375                         is gone"
376                    )
377                } else {
378                    format!("OPERATOR — pinned appliance; `yah cloud migrate` to {target}")
379                }
380            }
381            Self::RunLost => "run lost — re-run it; nothing fails a job over".to_string(),
382            Self::DegradedButServing { surviving_replicas } => {
383                format!("degraded — {surviving_replicas} replica(s) still serving")
384            }
385        }
386    }
387}
388
389/// How much of the workload's state is gone, and how far back the copy is.
390#[derive(Debug, Clone, PartialEq, Serialize)]
391#[serde(tag = "kind", rename_all = "snake_case")]
392pub enum DataLoss {
393    /// Nothing durable is mounted.
394    None,
395    /// Only tmpfs. Discarded on stop by definition, so the node dying costs
396    /// nothing that a restart would not have.
397    EphemeralOnly,
398    /// Durable state exists and there is **no second copy anywhere**. The
399    /// bytes are gone with the node.
400    Total {
401        volumes: Vec<String>,
402        /// Why there is no copy: no declaration at all, or `tier = "none"`.
403        because: String,
404    },
405    /// A copy exists in an object store; the loss is the gap between the last
406    /// write and the last thing that reached the store.
407    Window {
408        tier: DurabilityTier,
409        store: String,
410        /// `None` for snapshot/dedup tiers, whose recovery point is set by
411        /// whatever schedules the snapshot and is therefore not in the spec.
412        rpo_seconds: Option<u32>,
413        /// Present when [`Self::rpo_seconds`] is the turso-backup default
414        /// rather than a declared value.
415        rpo_is_default: bool,
416    },
417    /// Bind mounts only. The camp did not create the host path and cannot know
418    /// whether the bytes exist elsewhere — the same refusal-to-guess
419    /// `migrate::preconditions` makes for the same mount kind.
420    Unknown { volumes: Vec<String> },
421}
422
423impl DataLoss {
424    /// One line for a text report.
425    pub fn headline(&self) -> String {
426        match self {
427            Self::None => "none — no durable state declared".to_string(),
428            Self::EphemeralOnly => "none — tmpfs only, discarded on stop anyway".to_string(),
429            Self::Total { volumes, because } => format!(
430                "TOTAL — {} has no second copy anywhere ({because})",
431                volumes.join(", ")
432            ),
433            Self::Window {
434                tier,
435                store,
436                rpo_seconds,
437                rpo_is_default,
438            } => match rpo_seconds {
439                Some(s) if *rpo_is_default => {
440                    format!("≤ {s}s (turso-backup default, not declared) — tier {tier} → {store}")
441                }
442                Some(s) => format!("≤ {s}s (declared) — tier {tier} → {store}"),
443                None => format!(
444                    "unbounded by the spec — tier {tier} → {store}; a snapshot tier's \
445                     recovery point is set by whatever schedules it"
446                ),
447            },
448            Self::Unknown { volumes } => format!(
449                "UNKNOWN — {} are operator-managed bind mounts; the camp cannot say \
450                 whether the bytes exist anywhere else",
451                volumes.join(", ")
452            ),
453        }
454    }
455}
456
457/// How long it takes to get the state back, and on what basis that is claimed.
458///
459/// Every non-trivial variant here is an **extrapolation from two measured
460/// constants** ([`MEASURED_HYDRATE_MB_PER_S`], [`MEASURED_GET_RTT_MS`]) applied
461/// to a **declared** state size. Nothing in this module has ever timed a real
462/// restore. That is stated on the type rather than in a footnote because a
463/// recovery-time number without its provenance is the single easiest thing in
464/// a planning report to mistake for a measurement.
465#[derive(Debug, Clone, PartialEq, Serialize)]
466#[serde(tag = "kind", rename_all = "snake_case")]
467pub enum RecoveryEstimate {
468    /// Nothing to hydrate — the workload carries no durable state.
469    Immediate,
470    /// There is no copy to recover from. Recovery is not a duration.
471    NotRecoverable,
472    /// A copy exists but the spec does not say how big the state is, so the
473    /// transfer cannot be estimated. Names the annotation that would fix it.
474    UnknownStateSize { hint: &'static str },
475    /// Bulk transfer of a declared state size at the measured throughput.
476    Hydrate {
477        state_mb: u32,
478        seconds: f64,
479        /// Verbatim provenance, carried into JSON so a consumer cannot strip it.
480        basis: String,
481    },
482}
483
484impl RecoveryEstimate {
485    /// One line for a text report.
486    pub fn headline(&self) -> String {
487        match self {
488            Self::Immediate => "immediate — stateless".to_string(),
489            Self::NotRecoverable => "n/a — nothing to recover from".to_string(),
490            Self::UnknownStateSize { hint } => {
491                format!("unknown — declare {hint} to get an estimate")
492            }
493            Self::Hydrate {
494                state_mb, seconds, ..
495            } => format!(
496                "≥ ~{seconds:.1}s to pull {state_mb} MiB (extrapolated from \
497                 {MEASURED_HYDRATE_MB_PER_S} MB/s measured once, R760-T10; no restore was \
498                 timed here, and WAL replay is on top)"
499            ),
500        }
501    }
502}
503
504/// What losing a machine does to its sovereign group's raft quorum.
505#[derive(Debug, Clone, PartialEq, Serialize)]
506#[serde(tag = "kind", rename_all = "snake_case")]
507pub enum QuorumEffect {
508    /// The machine declares no `sovereign_group`, so it votes in nothing.
509    NotInAGroup,
510    /// `sovereign_role = "non-voter"` — in the group's blast radius, holds no
511    /// seat. Its absence from `/raft/status` is correct, not drift.
512    NonVoter { group: String },
513    /// A voter is lost and the survivors still make a majority.
514    QuorumHolds {
515        group: String,
516        voters_before: usize,
517        voters_after: usize,
518        majority_needed: usize,
519    },
520    /// A voter is lost and the survivors do not. The group's raft stops
521    /// accepting writes, which includes cluster secrets — so workloads there
522    /// fail to resolve secrets even if their own containers are untouched.
523    QuorumLost {
524        group: String,
525        voters_before: usize,
526        voters_after: usize,
527        majority_needed: usize,
528    },
529}
530
531impl QuorumEffect {
532    /// One line for a text report.
533    pub fn headline(&self) -> String {
534        match self {
535            Self::NotInAGroup => "no sovereign group — votes in nothing".to_string(),
536            Self::NonVoter { group } => {
537                format!("non-voter in '{group}' — no quorum seat to lose")
538            }
539            Self::QuorumHolds {
540                group,
541                voters_after,
542                majority_needed,
543                ..
544            } => format!(
545                "'{group}' quorum holds — {voters_after} voter(s) left, {majority_needed} needed"
546            ),
547            Self::QuorumLost {
548                group,
549                voters_after,
550                majority_needed,
551                ..
552            } => format!(
553                "'{group}' LOSES QUORUM — {voters_after} voter(s) left, {majority_needed} \
554                 needed; the group's raft stops accepting writes, and cluster secrets are \
555                 read from the local raft replica, so workloads there fail to resolve \
556                 secrets even where their containers are untouched"
557            ),
558        }
559    }
560}
561
562/// A machine whose declared capacity does not cover what is projected onto it.
563///
564/// Admission checks each workload against the node's `[allocatable]`
565/// *individually* (`RequiredSpec::matches`, R572-F5) and never subtracts what
566/// is already committed — so N workloads that each fit can all be admitted onto
567/// a node that cannot hold their sum. This is the arithmetic nothing in the
568/// tree does today.
569#[derive(Debug, Clone, Serialize)]
570pub struct Oversubscription {
571    pub machine: String,
572    pub allocatable: NodeAllocatable,
573    pub committed_memory_mb: u32,
574    pub committed_cpu_millis: u32,
575    pub workloads: Vec<String>,
576    pub memory_over: bool,
577    pub cpu_over: bool,
578}
579
580// ─── The traversal ───────────────────────────────────────────────────────────
581
582/// Walk the declared graph once and answer every question derivable from it.
583///
584/// Pure: same TOML in, same [`Topology`] out, no network. See the module header
585/// for what is deliberately outside the model.
586pub fn analyze(cfg: &CloudConfig) -> Topology {
587    let workloads: Vec<WorkloadNode> = cfg
588        .workloads
589        .iter()
590        .map(|w| workload_node(cfg, &w.spec))
591        .collect();
592
593    let mut machines: Vec<MachineNode> = cfg
594        .machines
595        .iter()
596        .map(|m| MachineNode {
597            name: m.name.clone(),
598            region: m.region.clone(),
599            sovereign_group: m.sovereign_group.clone(),
600            sovereign_role: m.sovereign_role.unwrap_or_default(),
601            taints: m.taints.clone(),
602            allocatable: m.allocatable.clone(),
603            placed: Vec::new(),
604            committed_memory_mb: 0,
605            committed_cpu_millis: 0,
606        })
607        .collect();
608
609    // Fold each workload's projected placement back onto its machine. Replicas
610    // multiply the commitment: `replicas = 3` asks the node for three copies of
611    // the request, and admission — which checks one workload against one node —
612    // never sees that multiplication.
613    let by_name: BTreeMap<String, usize> = machines
614        .iter()
615        .enumerate()
616        .map(|(i, m)| (m.name.clone(), i))
617        .collect();
618    for w in &workloads {
619        let Some(machine) = w.placement.machine() else {
620            continue;
621        };
622        let Some(&i) = by_name.get(machine) else {
623            continue;
624        };
625        let copies = w.replicas.max(1);
626        machines[i].placed.push(w.name.clone());
627        machines[i].committed_memory_mb = machines[i]
628            .committed_memory_mb
629            .saturating_add(w.memory_request_mb.saturating_mul(copies));
630        machines[i].committed_cpu_millis = machines[i]
631            .committed_cpu_millis
632            .saturating_add(w.cpu_millis.saturating_mul(copies));
633    }
634
635    let oversubscribed = machines.iter().filter_map(oversubscription).collect();
636    let node_losses = machines
637        .iter()
638        .map(|m| node_loss(cfg, &machines, &workloads, &m.name))
639        .collect();
640
641    Topology {
642        machines,
643        workloads,
644        node_losses,
645        oversubscribed,
646    }
647}
648
649fn workload_node(cfg: &CloudConfig, spec: &WorkloadSpec) -> WorkloadNode {
650    // One call into the admission seam — the same one `yah cloud apply` and
651    // `yah cloud migrate` use. Forking a second selector here would let the
652    // analyzer report a placement the fleet would never make.
653    let placement = match cfg.admit_workload_candidates(spec) {
654        Ok(candidates) => {
655            let mut names = candidates.iter().map(|m| m.name.clone());
656            let chosen = names
657                .next()
658                .expect("admit_workload_candidates never returns empty");
659            Placement::Admitted {
660                chosen,
661                alternates: names.collect(),
662                pinned: crate::config::node_selector_node(spec).is_some(),
663            }
664        }
665        Err(e) => Placement::Unschedulable {
666            reason: e.to_string(),
667        },
668    };
669
670    let volumes = spec
671        .volumes
672        .iter()
673        .map(|v| match &v.source {
674            VolumeSource::Named { name } => VolumeDisposition::Copy {
675                name: name.clone(),
676                host_path: named_volume_path(name),
677                mounted_at: v.target.clone(),
678            },
679            VolumeSource::Bind { host_path } => VolumeDisposition::Precondition {
680                host_path: host_path.clone(),
681                mounted_at: v.target.clone(),
682            },
683            VolumeSource::Tmpfs { size_mb } => VolumeDisposition::Discard {
684                mounted_at: v.target.clone(),
685                size_mb: *size_mb,
686            },
687        })
688        .collect();
689
690    let durability = match spec.durability() {
691        Ok(Some(d)) => DurabilityView::Declared(d),
692        Ok(None) => DurabilityView::Undeclared,
693        Err(e) => DurabilityView::Malformed {
694            reason: e.to_string(),
695        },
696    };
697
698    WorkloadNode {
699        name: spec.name.clone(),
700        tier: spec.tier.0.clone(),
701        mesh_identity: spec.fq_mesh_identity(),
702        archetype: spec.effective_archetype(),
703        replicas: spec.replicas,
704        restart_policy: restart_policy_label(&spec.restart_policy),
705        placement,
706        volumes,
707        durability,
708        memory_request_mb: spec.memory_request_mb(),
709        cpu_millis: spec.resources.cpu_millis,
710        public_hostnames: spec
711            .expose
712            .public
713            .iter()
714            .map(|p| p.hostname.clone())
715            .collect(),
716        depends_on: spec.depends_on.iter().map(|d| d.0.clone()).collect(),
717    }
718}
719
720fn restart_policy_label(p: &RestartPolicy) -> String {
721    match p {
722        RestartPolicy::Always => "always".to_string(),
723        RestartPolicy::OnFailure { max_attempts, .. } => {
724            format!("on-failure (max {max_attempts})")
725        }
726        RestartPolicy::Never => "never".to_string(),
727    }
728}
729
730fn oversubscription(m: &MachineNode) -> Option<Oversubscription> {
731    // No `[allocatable]` block is "unconstrained", exactly as admission reads
732    // it — not "zero capacity". Reporting an unbounded node as oversubscribed
733    // would flag every machine that has not been measured yet.
734    let alloc = m.allocatable.as_ref()?;
735    let memory_over = m.committed_memory_mb > alloc.memory_mb;
736    let cpu_over = alloc.cpu_millis > 0 && m.committed_cpu_millis > alloc.cpu_millis;
737    if !memory_over && !cpu_over {
738        return None;
739    }
740    Some(Oversubscription {
741        machine: m.name.clone(),
742        allocatable: alloc.clone(),
743        committed_memory_mb: m.committed_memory_mb,
744        committed_cpu_millis: m.committed_cpu_millis,
745        workloads: m.placed.clone(),
746        memory_over,
747        cpu_over,
748    })
749}
750
751fn node_loss(
752    cfg: &CloudConfig,
753    machines: &[MachineNode],
754    workloads: &[WorkloadNode],
755    dead: &str,
756) -> NodeLoss {
757    let impacts: Vec<WorkloadImpact> = workloads
758        .iter()
759        .filter(|w| w.placement.machine() == Some(dead))
760        .map(|w| workload_impact(w, dead))
761        .collect();
762
763    let public_endpoints_lost = impacts
764        .iter()
765        .filter_map(|i| workloads.iter().find(|w| w.name == i.workload))
766        .flat_map(|w| w.public_hostnames.iter().cloned())
767        .collect();
768
769    NodeLoss {
770        machine: dead.to_string(),
771        impacts,
772        public_endpoints_lost,
773        quorum: quorum_effect(cfg, machines, dead),
774    }
775}
776
777fn workload_impact(w: &WorkloadNode, dead: &str) -> WorkloadImpact {
778    let alternates: Vec<String> = w
779        .placement
780        .alternates()
781        .iter()
782        .filter(|m| m.as_str() != dead)
783        .cloned()
784        .collect();
785
786    let data_loss = data_loss(w);
787    let outcome = outcome(w, &alternates, dead);
788    let recovery = recovery(w, &data_loss);
789
790    WorkloadImpact {
791        workload: w.name.clone(),
792        archetype: w.archetype,
793        outcome,
794        data_loss,
795        recovery,
796    }
797}
798
799/// The core verdict. Archetype decides it, because archetype is what yubaba
800/// itself branches on — `drain_workloads` skips appliances (R572-F4), and
801/// `migrate` orders its steps by the same split.
802fn outcome(w: &WorkloadNode, alternates: &[String], dead: &str) -> Outcome {
803    match w.archetype {
804        LifecycleArchetype::Appliance => Outcome::PinnedAppliance {
805            migrate_target: alternates.first().cloned(),
806            // "Hardware-level kill" is the question being asked, so the source
807            // side of migrate's stop → copy → start has nothing to read from.
808            // A workload whose only durable mount is a named volume on the dead
809            // box is the exact shape with no source; one with no durable state
810            // has nothing to copy and so is not blocked on this.
811            source_unreachable: w.has_named_volume(),
812        },
813        LifecycleArchetype::Job => Outcome::RunLost,
814        LifecycleArchetype::Server => {
815            if alternates.is_empty() {
816                return Outcome::NowhereToGo {
817                    reason: format!(
818                        "{dead} is the only machine in the declared fleet that admits \
819                         {} (archetype {}, {} MiB request, {} millicores)",
820                        w.name,
821                        w.archetype.taint_key(),
822                        w.memory_request_mb,
823                        w.cpu_millis,
824                    ),
825                };
826            }
827            if w.replicas > 1 {
828                Outcome::DegradedButServing {
829                    surviving_replicas: w.replicas - 1,
830                }
831            } else {
832                Outcome::Reschedulable {
833                    candidates: alternates.to_vec(),
834                }
835            }
836        }
837    }
838}
839
840fn data_loss(w: &WorkloadNode) -> DataLoss {
841    let durable = w.durable_volumes();
842    if durable.is_empty() {
843        return if w.volumes.is_empty() {
844            DataLoss::None
845        } else {
846            DataLoss::EphemeralOnly
847        };
848    }
849
850    // A bind mount is operator-managed; the camp did not create the host path
851    // and has no basis for a claim about it either way. Only say "total" about
852    // volumes this camp is responsible for.
853    if !w.has_named_volume() {
854        return DataLoss::Unknown { volumes: durable };
855    }
856
857    match &w.durability {
858        DurabilityView::Undeclared => DataLoss::Total {
859            volumes: durable,
860            because: "no yah.durability.tier declared, so the yubaba-managed named volume \
861                      at /var/lib/yah/kamaji/volumes/ is the only copy"
862                .to_string(),
863        },
864        DurabilityView::Malformed { reason } => DataLoss::Total {
865            volumes: durable,
866            because: format!("the durability declaration cannot be read: {reason}"),
867        },
868        DurabilityView::Declared(d) => match d.tier {
869            DurabilityTier::None => DataLoss::Total {
870                volumes: durable,
871                because: "yah.durability.tier = \"none\" — deliberately no second copy".to_string(),
872            },
873            DurabilityTier::Snapshot | DurabilityTier::Dedup => DataLoss::Window {
874                tier: d.tier,
875                store: d.store.clone().unwrap_or_default(),
876                rpo_seconds: None,
877                rpo_is_default: false,
878            },
879            DurabilityTier::Stream => DataLoss::Window {
880                tier: d.tier,
881                store: d.store.clone().unwrap_or_default(),
882                rpo_seconds: Some(d.rpo_seconds.unwrap_or(DEFAULT_STREAM_RPO_SECONDS)),
883                rpo_is_default: d.rpo_seconds.is_none(),
884            },
885        },
886    }
887}
888
889fn recovery(w: &WorkloadNode, loss: &DataLoss) -> RecoveryEstimate {
890    match loss {
891        DataLoss::None | DataLoss::EphemeralOnly => RecoveryEstimate::Immediate,
892        DataLoss::Total { .. } | DataLoss::Unknown { .. } => RecoveryEstimate::NotRecoverable,
893        DataLoss::Window { .. } => {
894            let state_mb = match &w.durability {
895                DurabilityView::Declared(Durability {
896                    state_mb: Some(mb), ..
897                }) => *mb,
898                _ => {
899                    return RecoveryEstimate::UnknownStateSize {
900                        hint: workload_spec::DURABILITY_STATE_MB_ANNOTATION,
901                    }
902                }
903            };
904            RecoveryEstimate::Hydrate {
905                state_mb,
906                seconds: f64::from(state_mb) / MEASURED_HYDRATE_MB_PER_S,
907                // The *floor*, and it says so. A tier-2 restore also replays
908                // WAL frames, and `turso_backup::stream` fetches those one at a
909                // time — roadcase measured `2n + 1` serialized GETs for `n`
910                // generations, which at this RTT reaches the same order as the
911                // bulk transfer itself. `n` is not declared anywhere, so it is
912                // named rather than guessed at.
913                basis: format!(
914                    "bulk transfer at {MEASURED_HYDRATE_MB_PER_S} MB/s, measured R760-T10 \
915                     2026-08-29 on one host against one backend. A FLOOR: a tier-2 restore \
916                     adds 2n+1 serialized GETs at ~{MEASURED_GET_RTT_MS} ms each for n \
917                     generations, and n is not declared anywhere"
918                ),
919            }
920        }
921    }
922}
923
924fn quorum_effect(cfg: &CloudConfig, machines: &[MachineNode], dead: &str) -> QuorumEffect {
925    let Some(m) = machines.iter().find(|m| m.name == dead) else {
926        return QuorumEffect::NotInAGroup;
927    };
928    let Some(group) = m.sovereign_group.clone() else {
929        return QuorumEffect::NotInAGroup;
930    };
931    if !m.sovereign_role.is_voter() {
932        return QuorumEffect::NonVoter { group };
933    }
934
935    let voters_before = cfg
936        .machines_in_group(&group)
937        .into_iter()
938        .filter(|m| m.sovereign_role.unwrap_or_default().is_voter())
939        .count();
940    let voters_after = voters_before.saturating_sub(1);
941    // Raft majority is over the *configured* membership, which the loss of a
942    // box does not shrink — a dead voter still counts in the denominator until
943    // someone removes it from the configuration.
944    let majority_needed = voters_before / 2 + 1;
945
946    if voters_after >= majority_needed {
947        QuorumEffect::QuorumHolds {
948            group,
949            voters_before,
950            voters_after,
951            majority_needed,
952        }
953    } else {
954        QuorumEffect::QuorumLost {
955            group,
956            voters_before,
957            voters_after,
958            majority_needed,
959        }
960    }
961}
962
963// ─── P2: renders ─────────────────────────────────────────────────────────────
964
965impl Topology {
966    /// Render the model as a Mermaid `flowchart`.
967    ///
968    /// **A projection, not a second traversal.** Every node and edge below is
969    /// read off fields [`analyze`] already computed, so the picture cannot
970    /// disagree with the verdicts printed beside it. A diagram generated by its
971    /// own walk of the config would be free to drift into decoration, which is
972    /// worse than no diagram — it is the failure mode this method's shape
973    /// exists to make impossible.
974    ///
975    /// The edge worth the whole render is `-.->|hydrate|`: the backup path is
976    /// the one relationship in this graph that is invisible in the TOML, has no
977    /// runtime today, and is exactly what decides whether a node loss is an
978    /// incident or a restore.
979    pub fn to_mermaid(&self) -> String {
980        let mut out = String::from("flowchart TB\n");
981
982        // Machines, grouped by sovereign group. The grouping is the blast
983        // radius (W305), so it is what a reader should see first.
984        let mut groups: BTreeMap<Option<&str>, Vec<&MachineNode>> = BTreeMap::new();
985        for m in &self.machines {
986            groups
987                .entry(m.sovereign_group.as_deref())
988                .or_default()
989                .push(m);
990        }
991        for (group, members) in &groups {
992            let label = group.unwrap_or("ungrouped");
993            out.push_str(&format!(
994                "  subgraph grp_{}[\"{label}\"]\n",
995                sanitize(label)
996            ));
997            for m in members {
998                // `sovereign_role` defaults to Voter, so a box that declares no
999                // group would otherwise render as "voter" — a seat in a quorum
1000                // it is not in. Match what `QuorumEffect::NotInAGroup` says.
1001                let role = match (group, m.sovereign_role.is_voter()) {
1002                    (None, _) => "no group",
1003                    (Some(_), true) => "voter",
1004                    (Some(_), false) => "non-voter",
1005                };
1006                let cap = match &m.allocatable {
1007                    Some(a) => format!(
1008                        "<br/>{}/{} MiB · {}/{} mCPU",
1009                        m.committed_memory_mb, a.memory_mb, m.committed_cpu_millis, a.cpu_millis
1010                    ),
1011                    None => "<br/>no [allocatable] declared".to_string(),
1012                };
1013                out.push_str(&format!(
1014                    "    {}[\"{}<br/><i>{role}</i>{cap}\"]\n",
1015                    node_id("m", &m.name),
1016                    m.name
1017                ));
1018            }
1019            out.push_str("  end\n");
1020        }
1021
1022        // Workloads, their mounts, and their public front doors.
1023        for w in &self.workloads {
1024            let wid = node_id("w", &w.name);
1025            out.push_str(&format!(
1026                "  {wid}(\"{}<br/><i>{}</i> · replicas {}\")\n",
1027                w.name,
1028                w.archetype.taint_key(),
1029                w.replicas
1030            ));
1031
1032            match &w.placement {
1033                Placement::Admitted { chosen, pinned, .. } => {
1034                    let verb = if *pinned { "pinned" } else { "placed" };
1035                    out.push_str(&format!("  {} -->|{verb}| {wid}\n", node_id("m", chosen)));
1036                }
1037                Placement::Unschedulable { .. } => {
1038                    out.push_str(&format!(
1039                        "  unschedulable{{{{no node admits it}}}} --> {wid}\n"
1040                    ));
1041                }
1042            }
1043
1044            for v in &w.volumes {
1045                let vid = node_id("v", &format!("{}-{}", w.name, v.label()));
1046                let (shape, edge) = match v {
1047                    VolumeDisposition::Copy { name, .. } => {
1048                        (format!("{vid}[(\"named: {name}\")]"), "mount")
1049                    }
1050                    VolumeDisposition::Precondition { host_path, .. } => (
1051                        format!("{vid}[(\"bind: {}\")]", host_path.display()),
1052                        "mount",
1053                    ),
1054                    VolumeDisposition::Discard { size_mb, .. } => {
1055                        (format!("{vid}[(\"tmpfs {size_mb} MiB\")]"), "ephemeral")
1056                    }
1057                };
1058                out.push_str(&format!("  {shape}\n  {wid} -->|{edge}| {vid}\n"));
1059
1060                // The invisible edge. Only durable mounts can have one, and
1061                // only a declared tier draws it.
1062                if !v.is_durable() {
1063                    continue;
1064                }
1065                if let DurabilityView::Declared(d) = &w.durability {
1066                    if let Some(store) = &d.store {
1067                        let sid = node_id("s", store);
1068                        out.push_str(&format!("  {sid}[[\"{store}\"]]\n"));
1069                        out.push_str(&format!(
1070                            "  {vid} -.->|backup: {}| {sid}\n  {sid} -.->|hydrate| {vid}\n",
1071                            d.tier
1072                        ));
1073                    }
1074                }
1075            }
1076
1077            for host in &w.public_hostnames {
1078                let hid = node_id("p", host);
1079                out.push_str(&format!("  {hid}>\"{host}\"]\n  {hid} ==>|public| {wid}\n"));
1080            }
1081
1082            for dep in &w.depends_on {
1083                if let Some(target) = self.workloads.iter().find(|o| {
1084                    o.mesh_identity == *dep || o.mesh_identity.ends_with(&format!("/{dep}"))
1085                }) {
1086                    out.push_str(&format!(
1087                        "  {wid} -.->|mesh admit| {}\n",
1088                        node_id("w", &target.name)
1089                    ));
1090                }
1091            }
1092        }
1093
1094        out
1095    }
1096
1097    /// Render the survivability answer for one machine as plain text.
1098    ///
1099    /// Returns `None` when no machine by that name is declared — the caller
1100    /// owns the wording of that refusal, since it has the declared list.
1101    pub fn render_node_loss(&self, machine: &str) -> Option<String> {
1102        let loss = self.node_losses.iter().find(|l| l.machine == machine)?;
1103        let mut out = format!("If {machine} is lost at the hardware level:\n\n");
1104        out.push_str(&format!("  quorum: {}\n", loss.quorum.headline()));
1105        if loss.public_endpoints_lost.is_empty() {
1106            out.push_str("  public endpoints lost: none\n");
1107        } else {
1108            out.push_str(&format!(
1109                "  public endpoints lost: {}\n",
1110                loss.public_endpoints_lost.join(", ")
1111            ));
1112        }
1113
1114        if loss.impacts.is_empty() {
1115            out.push_str("\n  No declared workload is projected onto this machine.\n");
1116            return Some(out);
1117        }
1118
1119        for i in &loss.impacts {
1120            out.push_str(&format!(
1121                "\n  {} ({})\n",
1122                i.workload,
1123                i.archetype.taint_key()
1124            ));
1125            out.push_str(&format!("    what happens: {}\n", i.outcome.headline()));
1126            out.push_str(&format!("    data loss:    {}\n", i.data_loss.headline()));
1127            out.push_str(&format!("    recovery:     {}\n", i.recovery.headline()));
1128        }
1129        Some(out)
1130    }
1131
1132    /// Render the capacity arithmetic — the whole fleet, oversubscription
1133    /// called out rather than left to the reader to spot.
1134    pub fn render_capacity(&self) -> String {
1135        let mut out = String::from("Declared capacity vs projected commitment:\n\n");
1136        for m in &self.machines {
1137            let placed = if m.placed.is_empty() {
1138                "(nothing)".to_string()
1139            } else {
1140                m.placed.join(", ")
1141            };
1142            match &m.allocatable {
1143                Some(a) => out.push_str(&format!(
1144                    "  {:<16} {:>6}/{:<6} MiB   {:>6}/{:<6} mCPU   {placed}\n",
1145                    m.name,
1146                    m.committed_memory_mb,
1147                    a.memory_mb,
1148                    m.committed_cpu_millis,
1149                    a.cpu_millis
1150                )),
1151                None => out.push_str(&format!(
1152                    "  {:<16} {:>6}/{:<6} MiB   {:>6}/{:<6} mCPU   {placed}\n",
1153                    m.name, m.committed_memory_mb, "?", m.committed_cpu_millis, "?"
1154                )),
1155            }
1156        }
1157        if self.oversubscribed.is_empty() {
1158            out.push_str("\nNo machine is oversubscribed against its declared [allocatable].\n");
1159        } else {
1160            out.push_str(
1161                "\nOVERSUBSCRIBED — admission checks each workload against a node \
1162                 individually and never subtracts what is already committed (R572-F5), \
1163                 so these all admitted and cannot all run:\n",
1164            );
1165            for o in &self.oversubscribed {
1166                let mut axes = Vec::new();
1167                if o.memory_over {
1168                    axes.push(format!(
1169                        "memory {} MiB > {} MiB",
1170                        o.committed_memory_mb, o.allocatable.memory_mb
1171                    ));
1172                }
1173                if o.cpu_over {
1174                    axes.push(format!(
1175                        "cpu {} > {} millicores",
1176                        o.committed_cpu_millis, o.allocatable.cpu_millis
1177                    ));
1178                }
1179                out.push_str(&format!(
1180                    "  {}: {} — {}\n",
1181                    o.machine,
1182                    axes.join(", "),
1183                    o.workloads.join(", ")
1184                ));
1185            }
1186        }
1187        out
1188    }
1189
1190    /// Render every machine's loss, plus the capacity view — the default
1191    /// whole-fleet report.
1192    pub fn render(&self) -> String {
1193        let mut out = String::new();
1194        for m in &self.machines {
1195            if let Some(section) = self.render_node_loss(&m.name) {
1196                out.push_str(&section);
1197                out.push('\n');
1198            }
1199        }
1200        out.push_str(&self.render_capacity());
1201        out
1202    }
1203}
1204
1205/// Mermaid node ids must be identifier-ish; machine and workload names are DNS
1206/// labels and hostnames, and buckets are URLs. One prefix per kind keeps two
1207/// different things that sanitize to the same string apart.
1208fn node_id(prefix: &str, name: &str) -> String {
1209    format!("{prefix}_{}", sanitize(name))
1210}
1211
1212fn sanitize(name: &str) -> String {
1213    name.chars()
1214        .map(|c| if c.is_ascii_alphanumeric() { c } else { '_' })
1215        .collect()
1216}
1217
1218#[cfg(test)]
1219mod tests {
1220    use super::*;
1221    use std::path::Path;
1222    use tempfile::{tempdir, TempDir};
1223
1224    /// Fixtures go through `CloudConfig::load`, never a struct literal, for the
1225    /// reason `migrate`'s own fixture records: a hand-built spec can express a
1226    /// shape no operator could write, and `load_workloads` shape-validates —
1227    /// which since R850-P4 includes the durability declaration. A test that
1228    /// skips the loader would pass on a spec the fleet would reject.
1229    struct Camp {
1230        dir: TempDir,
1231    }
1232
1233    impl Camp {
1234        fn new() -> Self {
1235            Self {
1236                dir: tempdir().unwrap(),
1237            }
1238        }
1239
1240        fn root(&self) -> &Path {
1241            self.dir.path()
1242        }
1243
1244        fn machine(self, name: &str, extra: &str) -> Self {
1245            let dir = self.root().join(".yah/infra/machines");
1246            std::fs::create_dir_all(&dir).unwrap();
1247            std::fs::write(
1248                dir.join(format!("{name}.toml")),
1249                format!(
1250                    "name = \"{name}\"\nprovider = \"static\"\nmesh_tags = []\n\
1251                     {extra}\n\
1252                     [connect]\naddress = \"10.0.0.1\"\nssh = \"root@{name}\"\n\
1253                     identity_file = \"~/.ssh/yah\"\n\
1254                     yubaba = \"http://{name}:7443\"\n"
1255                ),
1256            )
1257            .unwrap();
1258            self
1259        }
1260
1261        /// A machine with a capacity budget, which is what makes the
1262        /// oversubscription and NowhereToGo cases expressible.
1263        fn sized(self, name: &str, memory_mb: u32, cpu_millis: u32, extra: &str) -> Self {
1264            self.machine(
1265                name,
1266                &format!(
1267                    "{extra}\n[allocatable]\nmemory_mb = {memory_mb}\ncpu_millis = {cpu_millis}"
1268                ),
1269            )
1270        }
1271
1272        fn workload(self, name: &str, extra: &str) -> Self {
1273            self.workload_sized(name, 128, 100, extra)
1274        }
1275
1276        fn workload_sized(self, name: &str, memory_mb: u32, cpu_millis: u32, extra: &str) -> Self {
1277            let dir = self.root().join(".yah/infra/workloads");
1278            std::fs::create_dir_all(&dir).unwrap();
1279            std::fs::write(
1280                dir.join(format!("{name}.toml")),
1281                format!(
1282                    "schema_version = 1\nname = \"{name}\"\ntier = \"infra\"\n\
1283                     restart_policy = \"always\"\n\
1284                     {extra}\n\
1285                     [image]\nregistry = \"cr.yah.dev\"\nrepository = \"{name}\"\n\
1286                     tag = \"v1\"\ndigest = \"sha256:abc\"\n\
1287                     [resources]\nmemory_mb = {memory_mb}\ncpu_millis = {cpu_millis}\n\
1288                     ephemeral_storage_mb = 64\n\
1289                     [stop_policy]\nsignal = 15\ngrace_period = 10000\n\
1290                     [expose.mesh]\nidentity = \"{name}\"\nports = [8080]\nallow_from = []\n"
1291                ),
1292            )
1293            .unwrap();
1294            self
1295        }
1296
1297        fn analyze(&self) -> Topology {
1298            super::analyze(&CloudConfig::load(self.root()).expect("fixture camp must load"))
1299        }
1300    }
1301
1302    const NAMED_VOLUME: &str = "[[volumes]]\nsource = { named = { name = \"accounts\" } }\n\
1303                                target = \"/var/lib/app\"\nread_only = false\n";
1304
1305    fn impact<'a>(topo: &'a Topology, machine: &str, workload: &str) -> &'a WorkloadImpact {
1306        topo.node_losses
1307            .iter()
1308            .find(|l| l.machine == machine)
1309            .unwrap_or_else(|| panic!("no node_loss for {machine}"))
1310            .impacts
1311            .iter()
1312            .find(|i| i.workload == workload)
1313            .unwrap_or_else(|| panic!("{workload} is not projected onto {machine}"))
1314    }
1315
1316    // ── The driving question ─────────────────────────────────────────────────
1317
1318    /// R850's acceptance test, and the reason the relay exists.
1319    ///
1320    /// The noisetable-account shape verbatim: one replica, `archetype =
1321    /// "appliance"`, a yubaba-managed named volume, no durability tier
1322    /// declared. Hardware-kill the node. The analyzer has to say *out loud*
1323    /// that nothing reschedules, that the bytes are gone, and that even the
1324    /// operator verb has no source to copy from — because before this module
1325    /// existed, learning any of those three meant reading yubaba's source.
1326    #[test]
1327    fn hardware_killing_the_node_under_a_singleton_appliance_loses_everything() {
1328        let topo = Camp::new()
1329            .sized("us-west-001", 12288, 6000, "")
1330            .sized("us-west-003", 16384, 16000, "")
1331            .workload(
1332                "noisetable-account",
1333                &format!("replicas = 1\narchetype = \"appliance\"\n{NAMED_VOLUME}"),
1334            )
1335            .analyze();
1336
1337        let i = impact(&topo, "us-west-001", "noisetable-account");
1338
1339        // 1. Nothing reschedules — and the reason is the archetype, not a
1340        //    shortage of nodes. us-west-003 admits it fine and is irrelevant.
1341        let Outcome::PinnedAppliance {
1342            migrate_target,
1343            source_unreachable,
1344        } = &i.outcome
1345        else {
1346            panic!("expected PinnedAppliance, got {:?}", i.outcome);
1347        };
1348        assert_eq!(migrate_target.as_deref(), Some("us-west-003"));
1349
1350        // 2. `yah cloud migrate` is named as the human step AND disclaimed:
1351        //    it plans a copy, and a dead box has nothing to copy from.
1352        assert!(source_unreachable);
1353        let headline = i.outcome.headline();
1354        assert!(headline.contains("yah cloud migrate"), "{headline}");
1355        assert!(headline.contains("no source"), "{headline}");
1356
1357        // 3. Every account, passkey and session is gone.
1358        let DataLoss::Total { volumes, because } = &i.data_loss else {
1359            panic!("expected Total, got {:?}", i.data_loss);
1360        };
1361        assert_eq!(volumes, &["accounts".to_string()]);
1362        assert!(because.contains("yah.durability.tier"), "{because}");
1363        assert_eq!(i.recovery, RecoveryEstimate::NotRecoverable);
1364
1365        // 4. And the operator-facing text says all three without a source dive.
1366        let text = topo.render_node_loss("us-west-001").unwrap();
1367        assert!(text.contains("OPERATOR"), "{text}");
1368        assert!(text.contains("TOTAL"), "{text}");
1369        assert!(text.contains("nothing to recover from"), "{text}");
1370    }
1371
1372    /// The same spec's verdict must not soften when the appliance is the only
1373    /// thing the fleet could hold — `migrate_target` goes to `None` and the
1374    /// outcome stays manual rather than degrading into "nowhere to go", which
1375    /// would read as a capacity problem instead of an archetype one.
1376    #[test]
1377    fn an_appliance_with_no_alternate_still_reads_as_pinned_not_as_capacity() {
1378        let topo = Camp::new()
1379            .sized("solo", 1024, 1000, "")
1380            .workload(
1381                "db",
1382                &format!("replicas = 1\narchetype = \"appliance\"\n{NAMED_VOLUME}"),
1383            )
1384            .analyze();
1385
1386        let i = impact(&topo, "solo", "db");
1387        let Outcome::PinnedAppliance { migrate_target, .. } = &i.outcome else {
1388            panic!("expected PinnedAppliance, got {:?}", i.outcome);
1389        };
1390        assert_eq!(*migrate_target, None);
1391        assert!(i.outcome.is_manual());
1392    }
1393
1394    // ── The fungible half ────────────────────────────────────────────────────
1395
1396    #[test]
1397    fn a_stateless_server_with_another_admitting_node_is_automatic() {
1398        let topo = Camp::new()
1399            .sized("a", 4096, 4000, "")
1400            .sized("b", 4096, 4000, "")
1401            .workload("api", "replicas = 1\narchetype = \"server\"")
1402            .analyze();
1403
1404        let i = impact(&topo, "a", "api");
1405        assert_eq!(
1406            i.outcome,
1407            Outcome::Reschedulable {
1408                candidates: vec!["b".to_string()]
1409            }
1410        );
1411        assert_eq!(i.data_loss, DataLoss::None);
1412        assert_eq!(i.recovery, RecoveryEstimate::Immediate);
1413        assert!(!i.outcome.is_manual());
1414    }
1415
1416    /// The finding that is invisible in the TOML: a workload can be perfectly
1417    /// stateless and restartable and still be a single point of failure,
1418    /// because only one declared box clears its floor.
1419    #[test]
1420    fn a_stateless_server_that_only_one_node_admits_stays_down() {
1421        let topo = Camp::new()
1422            .sized("big", 8192, 8000, "")
1423            .sized("small", 512, 8000, "")
1424            .workload_sized("hungry", 4096, 100, "replicas = 1\narchetype = \"server\"")
1425            .analyze();
1426
1427        let i = impact(&topo, "big", "hungry");
1428        let Outcome::NowhereToGo { reason } = &i.outcome else {
1429            panic!("expected NowhereToGo, got {:?}", i.outcome);
1430        };
1431        assert!(reason.contains("big"), "{reason}");
1432        assert!(i.outcome.is_manual());
1433    }
1434
1435    #[test]
1436    fn replicas_above_one_degrade_rather_than_fail() {
1437        let topo = Camp::new()
1438            .sized("a", 4096, 4000, "")
1439            .sized("b", 4096, 4000, "")
1440            .workload("api", "replicas = 3\narchetype = \"server\"")
1441            .analyze();
1442
1443        assert_eq!(
1444            impact(&topo, "a", "api").outcome,
1445            Outcome::DegradedButServing {
1446                surviving_replicas: 2
1447            }
1448        );
1449    }
1450
1451    /// A `no-appliance` taint is absolute — there is no toleration anywhere in
1452    /// the tree (W305 finding 2) — so a tainted box must never show up as an
1453    /// alternate an appliance could fail over to.
1454    #[test]
1455    fn a_no_appliance_taint_removes_a_node_from_the_alternates() {
1456        let topo = Camp::new()
1457            .sized("prod", 8192, 8000, "")
1458            .sized("builder", 16384, 16000, "taints = [\"no-appliance\"]")
1459            .workload(
1460                "db",
1461                &format!("replicas = 1\narchetype = \"appliance\"\n{NAMED_VOLUME}"),
1462            )
1463            .analyze();
1464
1465        let db = topo.workloads.iter().find(|w| w.name == "db").unwrap();
1466        assert_eq!(db.placement.machine(), Some("prod"));
1467        assert!(
1468            db.placement.alternates().is_empty(),
1469            "builder carries no-appliance and must not be offered: {:?}",
1470            db.placement
1471        );
1472    }
1473
1474    #[test]
1475    fn a_job_is_a_lost_run_not_a_failover() {
1476        let topo = Camp::new()
1477            .sized("a", 4096, 4000, "")
1478            .workload("build", "replicas = 1\narchetype = \"job\"")
1479            .analyze();
1480        assert_eq!(impact(&topo, "a", "build").outcome, Outcome::RunLost);
1481    }
1482
1483    // ── Durability ───────────────────────────────────────────────────────────
1484
1485    #[test]
1486    fn a_declared_stream_tier_turns_total_loss_into_a_bounded_window() {
1487        let topo = Camp::new()
1488            .sized("a", 4096, 4000, "")
1489            .workload(
1490                "db",
1491                &format!(
1492                    "replicas = 1\narchetype = \"appliance\"\n{NAMED_VOLUME}\n\
1493                     [annotations]\n\
1494                     \"yah.durability.tier\" = \"stream\"\n\
1495                     \"yah.durability.engine\" = \"turso\"\n\
1496                     \"yah.durability.store\" = \"s3://backups/db\"\n\
1497                     \"yah.durability.subjects\" = \"accounts.db\"\n\
1498                     \"yah.durability.rpo-seconds\" = \"30\"\n\
1499                     \"yah.durability.state-mb\" = \"100\"\n"
1500                ),
1501            )
1502            .analyze();
1503
1504        let i = impact(&topo, "a", "db");
1505        assert_eq!(
1506            i.data_loss,
1507            DataLoss::Window {
1508                tier: DurabilityTier::Stream,
1509                store: "s3://backups/db".to_string(),
1510                rpo_seconds: Some(30),
1511                rpo_is_default: false,
1512            }
1513        );
1514
1515        // The recovery figure is roadcase's measured constant applied to the
1516        // declared size — 100 MiB at 32.6 MB/s ≈ 3.1 s, the same order as the
1517        // 3.3 s R760-T10 actually measured for 100 MB.
1518        let RecoveryEstimate::Hydrate {
1519            state_mb,
1520            seconds,
1521            basis,
1522        } = &i.recovery
1523        else {
1524            panic!("expected Hydrate, got {:?}", i.recovery);
1525        };
1526        assert_eq!(*state_mb, 100);
1527        assert!((*seconds - 3.067).abs() < 0.01, "{seconds}");
1528        // Provenance survives into the structured output, not just the text.
1529        assert!(basis.contains("R760-T10"), "{basis}");
1530        assert!(basis.contains("FLOOR"), "{basis}");
1531    }
1532
1533    /// An undeclared RPO on a stream tier is the turso-backup default, and the
1534    /// report has to say which of the two it is printing — a number the
1535    /// operator believes they chose is worse than no number.
1536    #[test]
1537    fn an_undeclared_rpo_is_labelled_as_the_default_not_as_a_choice() {
1538        let topo = Camp::new()
1539            .sized("a", 4096, 4000, "")
1540            .workload(
1541                "db",
1542                &format!(
1543                    "replicas = 1\narchetype = \"appliance\"\n{NAMED_VOLUME}\n\
1544                     [annotations]\n\
1545                     \"yah.durability.tier\" = \"stream\"\n\
1546                     \"yah.durability.engine\" = \"turso\"\n\
1547                     \"yah.durability.store\" = \"s3://backups/db\"\n\
1548                     \"yah.durability.subjects\" = \"accounts.db\"\n"
1549                ),
1550            )
1551            .analyze();
1552
1553        let i = impact(&topo, "a", "db");
1554        let DataLoss::Window {
1555            rpo_seconds,
1556            rpo_is_default,
1557            ..
1558        } = &i.data_loss
1559        else {
1560            panic!("expected Window, got {:?}", i.data_loss);
1561        };
1562        assert_eq!(*rpo_seconds, Some(DEFAULT_STREAM_RPO_SECONDS));
1563        assert!(rpo_is_default);
1564        assert!(i.data_loss.headline().contains("not declared"));
1565
1566        // No declared state size ⇒ no invented duration.
1567        assert_eq!(
1568            i.recovery,
1569            RecoveryEstimate::UnknownStateSize {
1570                hint: "yah.durability.state-mb"
1571            }
1572        );
1573    }
1574
1575    /// A snapshot tier has a copy but no recovery point the spec can state,
1576    /// and saying "≤ 120s" for it would be a fabrication.
1577    #[test]
1578    fn a_snapshot_tier_reports_an_unbounded_window_rather_than_borrowing_the_stream_rpo() {
1579        let topo = Camp::new()
1580            .sized("a", 4096, 4000, "")
1581            .workload(
1582                "db",
1583                &format!(
1584                    "replicas = 1\narchetype = \"appliance\"\n{NAMED_VOLUME}\n\
1585                     [annotations]\n\
1586                     \"yah.durability.tier\" = \"snapshot\"\n\
1587                     \"yah.durability.engine\" = \"turso\"\n\
1588                     \"yah.durability.store\" = \"s3://backups/db\"\n\
1589                     \"yah.durability.subjects\" = \"accounts.db\"\n"
1590                ),
1591            )
1592            .analyze();
1593
1594        let i = impact(&topo, "a", "db");
1595        let DataLoss::Window { rpo_seconds, .. } = &i.data_loss else {
1596            panic!("expected Window, got {:?}", i.data_loss);
1597        };
1598        assert_eq!(*rpo_seconds, None);
1599        assert!(i.data_loss.headline().contains("unbounded by the spec"));
1600    }
1601
1602    /// `tier = "none"` still loses everything — but as a decision, and the
1603    /// report must not read the same as the case where nobody looked.
1604    #[test]
1605    fn a_deliberate_none_tier_reads_differently_from_an_undeclared_one() {
1606        let camp = Camp::new().sized("a", 4096, 4000, "").workload(
1607            "cache",
1608            &format!(
1609                "replicas = 1\narchetype = \"appliance\"\n{NAMED_VOLUME}\n\
1610                 [annotations]\n\"yah.durability.tier\" = \"none\"\n"
1611            ),
1612        );
1613        let topo = camp.analyze();
1614        let DataLoss::Total { because, .. } = &impact(&topo, "a", "cache").data_loss else {
1615            panic!("expected Total");
1616        };
1617        assert!(because.contains("deliberately"), "{because}");
1618        assert!(
1619            !because.contains("no yah.durability.tier declared"),
1620            "{because}"
1621        );
1622    }
1623
1624    /// A bind mount is operator-managed. The camp did not create the host path
1625    /// and has no basis for saying the bytes are gone — the same refusal to
1626    /// guess `migrate::preconditions` makes about the same mount kind.
1627    #[test]
1628    fn a_bind_mount_is_unknown_rather_than_total() {
1629        let topo = Camp::new()
1630            .sized("a", 4096, 4000, "")
1631            .workload(
1632                "svc",
1633                "replicas = 1\narchetype = \"appliance\"\n\
1634                 [[volumes]]\nsource = { bind = { host_path = \"/srv/data\" } }\n\
1635                 target = \"/data\"\nread_only = false\n",
1636            )
1637            .analyze();
1638
1639        let i = impact(&topo, "a", "svc");
1640        assert!(matches!(i.data_loss, DataLoss::Unknown { .. }));
1641        assert!(i.data_loss.headline().contains("UNKNOWN"));
1642    }
1643
1644    #[test]
1645    fn tmpfs_only_is_not_a_data_loss() {
1646        let topo = Camp::new()
1647            .sized("a", 4096, 4000, "")
1648            .workload(
1649                "svc",
1650                "replicas = 1\narchetype = \"server\"\n\
1651                 [[volumes]]\nsource = { tmpfs = { size_mb = 64 } }\n\
1652                 target = \"/scratch\"\nread_only = false\n",
1653            )
1654            .analyze();
1655        assert_eq!(impact(&topo, "a", "svc").data_loss, DataLoss::EphemeralOnly);
1656    }
1657
1658    // ── Capacity (P3) ────────────────────────────────────────────────────────
1659
1660    /// The gap admission structurally cannot see: `RequiredSpec::matches`
1661    /// checks one workload against one node's `[allocatable]` and never
1662    /// subtracts what is already committed, so three workloads that each fit
1663    /// are each admitted onto a node that cannot hold their sum.
1664    #[test]
1665    fn workloads_that_each_fit_can_still_oversubscribe_the_node_they_all_land_on() {
1666        let topo = Camp::new()
1667            .sized("small", 1024, 4000, "")
1668            .workload_sized("a", 512, 100, "replicas = 1\narchetype = \"server\"")
1669            .workload_sized("b", 512, 100, "replicas = 1\narchetype = \"server\"")
1670            .workload_sized("c", 512, 100, "replicas = 1\narchetype = \"server\"")
1671            .analyze();
1672
1673        assert_eq!(topo.oversubscribed.len(), 1);
1674        let o = &topo.oversubscribed[0];
1675        assert_eq!(o.machine, "small");
1676        assert_eq!(o.committed_memory_mb, 1536);
1677        assert!(o.memory_over);
1678        assert!(!o.cpu_over);
1679        assert!(topo.render_capacity().contains("OVERSUBSCRIBED"));
1680    }
1681
1682    /// Replicas multiply the ask. Admission checks one copy against one node
1683    /// and never multiplies, so `replicas = 4` of a fitting workload is exactly
1684    /// the shape that admits cleanly and cannot run.
1685    #[test]
1686    fn replicas_multiply_the_commitment() {
1687        let topo = Camp::new()
1688            .sized("small", 1024, 4000, "")
1689            .workload_sized("api", 512, 100, "replicas = 4\narchetype = \"server\"")
1690            .analyze();
1691
1692        assert_eq!(topo.machines[0].committed_memory_mb, 2048);
1693        assert_eq!(topo.oversubscribed.len(), 1);
1694    }
1695
1696    /// A node with no `[allocatable]` is *unconstrained*, exactly as admission
1697    /// reads it — not a zero-capacity node. Flagging it would light up every
1698    /// machine nobody has measured yet and train the operator to ignore this.
1699    #[test]
1700    fn a_node_with_no_allocatable_block_is_never_reported_oversubscribed() {
1701        let topo = Camp::new()
1702            .machine("unmeasured", "")
1703            .workload_sized("api", 99999, 99999, "replicas = 1\narchetype = \"server\"")
1704            .analyze();
1705
1706        assert!(topo.oversubscribed.is_empty());
1707        assert!(topo.render_capacity().contains("unmeasured"));
1708    }
1709
1710    // ── Quorum ───────────────────────────────────────────────────────────────
1711
1712    #[test]
1713    fn losing_one_of_three_voters_keeps_quorum() {
1714        let topo = Camp::new()
1715            .machine(
1716                "a",
1717                "sovereign_group = \"prod\"\nsovereign_role = \"voter\"",
1718            )
1719            .machine(
1720                "b",
1721                "sovereign_group = \"prod\"\nsovereign_role = \"voter\"",
1722            )
1723            .machine(
1724                "c",
1725                "sovereign_group = \"prod\"\nsovereign_role = \"voter\"",
1726            )
1727            .analyze();
1728
1729        let q = &topo.node_losses[0].quorum;
1730        assert_eq!(
1731            *q,
1732            QuorumEffect::QuorumHolds {
1733                group: "prod".to_string(),
1734                voters_before: 3,
1735                voters_after: 2,
1736                majority_needed: 2,
1737            }
1738        );
1739    }
1740
1741    /// The second half of a node loss that is easy to miss: cluster secrets are
1742    /// read from the *local raft replica*, so a group that loses quorum fails
1743    /// workloads whose containers were never touched.
1744    #[test]
1745    fn losing_one_of_two_voters_loses_quorum_and_says_what_that_costs() {
1746        let topo = Camp::new()
1747            .machine(
1748                "a",
1749                "sovereign_group = \"prod\"\nsovereign_role = \"voter\"",
1750            )
1751            .machine(
1752                "b",
1753                "sovereign_group = \"prod\"\nsovereign_role = \"voter\"",
1754            )
1755            .analyze();
1756
1757        let q = &topo.node_losses[0].quorum;
1758        assert!(matches!(q, QuorumEffect::QuorumLost { .. }), "{q:?}");
1759        let headline = q.headline();
1760        assert!(headline.contains("LOSES QUORUM"), "{headline}");
1761        assert!(headline.contains("secrets"), "{headline}");
1762    }
1763
1764    /// A non-voter's absence from `/raft/status` is correct, not drift
1765    /// (`workload_spec::sovereign`), so losing one costs no seat.
1766    #[test]
1767    fn a_non_voter_has_no_seat_to_lose() {
1768        let topo = Camp::new()
1769            .machine(
1770                "a",
1771                "sovereign_group = \"prod\"\nsovereign_role = \"voter\"",
1772            )
1773            .machine(
1774                "w",
1775                "sovereign_group = \"prod\"\nsovereign_role = \"non-voter\"",
1776            )
1777            .analyze();
1778
1779        let q = &topo
1780            .node_losses
1781            .iter()
1782            .find(|l| l.machine == "w")
1783            .unwrap()
1784            .quorum;
1785        assert_eq!(
1786            *q,
1787            QuorumEffect::NonVoter {
1788                group: "prod".to_string()
1789            }
1790        );
1791    }
1792
1793    // ── P2: the render is a projection ───────────────────────────────────────
1794
1795    /// The scope trap this relay was warned about: a diagram produced by its
1796    /// own walk of the config can disagree with the analysis printed beside it.
1797    /// This pins that it cannot — the placement edge in the Mermaid output is
1798    /// the placement the model computed, and the hydrate edge exists exactly
1799    /// when the model found a durability tier.
1800    #[test]
1801    fn the_mermaid_render_agrees_with_the_model_it_projects() {
1802        let topo = Camp::new()
1803            .sized("us-west-001", 12288, 6000, "sovereign_group = \"prod\"")
1804            .sized("us-west-003", 16384, 16000, "taints = [\"no-appliance\"]")
1805            .workload(
1806                "backed-up",
1807                &format!(
1808                    "replicas = 1\narchetype = \"appliance\"\n{NAMED_VOLUME}\n\
1809                     [annotations]\n\
1810                     \"yah.durability.tier\" = \"stream\"\n\
1811                     \"yah.durability.engine\" = \"turso\"\n\
1812                     \"yah.durability.store\" = \"s3://backups/db\"\n\
1813                     \"yah.durability.subjects\" = \"accounts.db\"\n"
1814                ),
1815            )
1816            .analyze();
1817
1818        let mermaid = topo.to_mermaid();
1819        let w = topo
1820            .workloads
1821            .iter()
1822            .find(|w| w.name == "backed-up")
1823            .unwrap();
1824
1825        // The edge names the machine the model chose, not a re-derived one.
1826        assert_eq!(w.placement.machine(), Some("us-west-001"));
1827        assert!(
1828            mermaid.contains("m_us_west_001 -->|placed| w_backed_up"),
1829            "{mermaid}"
1830        );
1831        // The blast radius is what a reader sees first.
1832        assert!(mermaid.contains("subgraph grp_prod[\"prod\"]"), "{mermaid}");
1833        // The edge that is invisible in the TOML.
1834        assert!(mermaid.contains("|hydrate|"), "{mermaid}");
1835        assert!(mermaid.contains("s3://backups/db"), "{mermaid}");
1836        // Committed-vs-allocatable rides the same numbers as render_capacity.
1837        assert!(mermaid.contains("128/12288 MiB"), "{mermaid}");
1838    }
1839
1840    /// The negative half: no declared tier means no hydrate edge. A diagram
1841    /// that drew one anyway would show a safety property that does not exist,
1842    /// which is the single worst thing this render could do.
1843    #[test]
1844    fn no_declared_tier_draws_no_hydrate_edge() {
1845        let topo = Camp::new()
1846            .sized("a", 4096, 4000, "")
1847            .workload(
1848                "db",
1849                &format!("replicas = 1\narchetype = \"appliance\"\n{NAMED_VOLUME}"),
1850            )
1851            .analyze();
1852
1853        let mermaid = topo.to_mermaid();
1854        assert!(mermaid.contains("v_db_accounts"), "{mermaid}");
1855        assert!(!mermaid.contains("hydrate"), "{mermaid}");
1856    }
1857
1858    #[test]
1859    fn a_public_hostname_is_an_edge_and_is_reported_lost_with_its_node() {
1860        let topo = Camp::new()
1861            .sized("a", 4096, 4000, "")
1862            .workload(
1863                "web",
1864                "replicas = 1\narchetype = \"server\"\n\
1865                 [expose.public]\nhostname = \"app.example.com\"\nport = 8080\n\
1866                 tls = \"cf_managed\"\n",
1867            )
1868            .analyze();
1869
1870        assert_eq!(
1871            topo.node_losses[0].public_endpoints_lost,
1872            vec!["app.example.com".to_string()]
1873        );
1874        assert!(topo.to_mermaid().contains("|public|"));
1875        assert!(topo
1876            .render_node_loss("a")
1877            .unwrap()
1878            .contains("public endpoints lost: app.example.com"));
1879    }
1880
1881    /// `sovereign_role` defaults to `Voter`, so a box that declares no group
1882    /// renders as one unless the label is suppressed — a quorum seat in a
1883    /// quorum that does not exist. Caught against the live camp, where
1884    /// us-west-002 and us-west-015 declare no group.
1885    #[test]
1886    fn an_ungrouped_machine_is_not_labelled_a_voter() {
1887        let mermaid = Camp::new()
1888            .sized("loner", 1024, 1000, "")
1889            .analyze()
1890            .to_mermaid();
1891        assert!(mermaid.contains("<i>no group</i>"), "{mermaid}");
1892        assert!(!mermaid.contains("<i>voter</i>"), "{mermaid}");
1893    }
1894
1895    #[test]
1896    fn an_undeclared_machine_has_no_node_loss_section() {
1897        let topo = Camp::new().sized("a", 4096, 4000, "").analyze();
1898        assert!(topo.render_node_loss("typo").is_none());
1899        assert!(topo.render_node_loss("a").is_some());
1900    }
1901}