cloud/recovery_journal.rs
1//! Append-only JSONL journal of **measured** restore times (R850-T2).
2//!
3//! Path: `.yah/cloud/recovery.jsonl`, via [`crate::paths::recovery_journal`].
4//! One record per restored *subject*; replay yields the last record per
5//! `(workload, subject)`, summed into one [`WorkloadRecovery`] per workload.
6//!
7//! Same shape as [`crate::asset_journal`] — append-only JSONL, one `write` per
8//! line so concurrent writers are safe, replayed to a map, a missing file
9//! replaying to empty rather than erroring. Two deliberate differences:
10//!
11//! - **Synchronous.** [`crate::config::CloudConfig::load`] is sync and this
12//! journal is replayed from inside it, so `std::fs` rather than `tokio::fs`.
13//! There is no in-process `subscribe()` either: nothing watches a
14//! measurement the way the desktop panel watches an asset transition.
15//! - **Nobody prunes it.** Append-only is the whole answer to retention: a
16//! stale measurement is *reported with its age*, never discarded, because a
17//! real timed restore from six weeks ago still beats an extrapolation from
18//! `MEASURED_HYDRATE_MB_PER_S` — one constant measured once on one host
19//! against one backend. See [`crate::topology::RecoveryEstimate`].
20//!
21//! ## The input is the helper's own output
22//!
23//! `turso-backup-hydrate` prints one JSON line describing what it restored
24//! (`oss/turso-backup/src/bin/hydrate.rs::outcome_to_json`).
25//! [`RecoveryRecord::from_helper_json`] parses *that* line, keyed off its field
26//! names verbatim — `restored[].{subject,source,bytes,seconds}` — rather than a
27//! re-typed parallel vocabulary that could drift from it silently.
28//!
29//! `cloud` takes no dependency on `turso-backup` (see the note on
30//! [`crate::topology::DEFAULT_STREAM_RPO_SECONDS`] for why), so the coupling is
31//! pinned from both sides by a byte-identical fixture: the literal in
32//! `a_real_hydrate_line_parses_verbatim` here, and the full-line equality
33//! assertion in that helper's own
34//! `a_hydrated_outcome_reports_measured_bytes_and_seconds`.
35//!
36//! @yah:ticket(R850-T4, "Carry the hydrate measurement back over the kamaji wire so the recovery journal fills without an operator")
37//! @yah:status(review)
38//! @yah:at(2026-09-11T07:31:38Z)
39//! @yah:assignee(agent:bundle-anthropic-ashguard)
40//! @yah:phase(P3c)
41//! @yah:parent(R850)
42//! @yah:next("R850-T2 shipped the journal and the reader end: `.yah/cloud/recovery.jsonl`, replayed by `CloudConfig::load` into `CloudConfig.recovery_measurements`, surfacing as `RecoveryEstimate::Measured` in `yah cloud topology`. What it did NOT ship is an automatic writer. Today the only way a measurement reaches the camp tree is a human running `yah cloud topology --record-hydrate <path|-> --workload <w> --node <n>` against a JSON line they copied out of a node's kamaji log. That works and is tested, but a planning surface nobody remembers to feed reports an extrapolation forever.")
43//! @yah:next("THE SEAM, and why it was deliberately not taken in R850-T2. `kamaji::hydrate::run` already HAS the measurement in hand — it returns `HydrateResult::Proceed(Some(line))` carrying turso-backup-hydrate's JSON verbatim — and oss/kamaji/crates/kamaji-bin/src/server.rs:2086 does nothing with it but `info!`. Carrying it back to yubaba means widening the deploy response in oss/kamaji/crates/kamaji-proto/src/messages.rs and threading it through server.rs, and BOTH of those files were uncommitted-dirty in the shared working tree on 2026-09-10 and plausibly a live peer's. That is a scheduling reason, not a design objection — re-check `camp.roster` and the working tree before starting.")
44//! @yah:next("WATCH THE WIRE SHAPE: kamaji-proto is positional postcard (R590-B3), the same constraint that pushed durability onto annotations rather than WorkloadSpec fields in R850-P4. Add the field in the way that crate already handles additive change; do not assume a struct field is free.")
45//! @yah:next("ATTRIBUTION IS THE REAL WORK, not the transport. A journal record needs {workload, node} to be useful, and the helper's own JSON carries neither — the CLI seam takes them as flags for exactly that reason. Whatever lands here must attach the deploying node's name and the workload's name at the point where both are known, which is yubaba's side of the deploy, not the helper's.")
46//! @yah:gotcha("DO NOT MAKE `topology::analyze` DO I/O — the same hard constraint R850-T2 held. analyze is a pure function of what `CloudConfig::load` read off the local tree, which is what makes it trustworthy in a test and is the contract `migrate::plan_migration` also holds. This ticket changes only how the journal FILLS; the read path is done and must not grow a network call.")
47//! @yah:handoff("THE MEASUREMENT NOW REACHES `.yah/cloud/recovery.jsonl` WITHOUT A HUMAN. Seven hops, all on disk: `kamaji::hydrate::run` -> new `KamajiToYubaba::DeployAck { request_id, id, hydrate: Option<String> }` (kamaji-proto/src/messages.rs) -> `DeployResult.hydrate` -> yubaba's deploy response (yubaba/src/lib.rs:4712) -> `cloud-client`'s `WorkloadDeployResponse.hydrate` over JSON -> `cloud::recovery_journal::record_deploy_measurement` -> the journal R850-T2 built. The manual `yah cloud topology --record-hydrate` seam stays and is untouched.")
48//! @yah:handoff("THE WIRE CHANGE IS A REAL BREAK AND IS BUMPED TO `ProtocolVersion::V10`, FOR TWO INDEPENDENT REASONS — the second one is the easy-to-miss one and is written into the version stanza. (1) An old peer cannot decode the new reply mid-deploy. (2) `AckKind::Deploy` was REMOVED rather than kept beside the new variant, and removing a variant RENUMBERS the remaining `AckKind` discriminants — so a V9 `Stop` ack would decode as `Probe`. The removal is deliberate under CLAUDE.md's pre-1.0 rule: two deploy-reply shapes where which one you got depended on which backend arm answered is exactly the shim that compounds, and deleting the variant made the compiler name all 20 call sites. The crate's own additive-change conventions were read first (appended VARIANTS ride `#[non_exhaustive]` unbumped; appended FIELDS are breaking — the V2/V4/V5/V6/V8/V9 stanzas; messages.rs:670 records the R746-B11 failure where an appended reply variant went unclassified, so the new one is classified in `reply_request_id` and added to `every_reply_variant_correlates_to_its_request`). The bump was verified to actually do something: server.rs:1373 refuses any `Hello` whose version != CURRENT with `UnsupportedVersion`, so skew is a named handshake refusal rather than a postcard error.")
49//! @yah:handoff("ATTRIBUTION WAS THE REAL WORK, AND IT LANDED SOMEWHERE OTHER THAN THE BRIEF SAID — deliberately, with the reasoning written into the code at both ends. The brief said append in `oss/yubaba/`. Node-side yubaba cannot: it has no name for itself that `.yah/infra/machines/*.toml` would recognise, and it is not running in the camp tree the journal lives in. The first scope holding BOTH facts and sitting inside the camp tree is the CLI, which dialed a machine BY NAME to deploy a workload BY NAME — so the journal write is `record_hydrate_from_deploy` at app/yah/cli/src/cloud.rs:3827. Yubaba still does its half: it puts the line on the deploy response, and the cloud-client JSON leg is genuinely self-describing so `serde(default)` is compatible there, unlike the postcard leg.")
50//! @yah:handoff("TWO CLI SITES, NOT ONE, AND THE CALL SITS AFTER THE STATUS MATCH. `handle_workload_deploy` and `handle_workload_rolling` — the rolling path matters MORE, because its destroy leaves the volume empty, so the redeploy IS the restore. The call is deliberately past the match arm: a real kamaji deploy answers \"deployed\" and only stub mode answers \"accepted\", so an arm-local call would have journaled exactly the deploys that ran nothing. The other four `deploy_workload` sites (passway, cloudflared, mesofact bundle, inner door) are unwired on purpose — hydrate-on-place only runs on kamaji's container-deploy path, so they structurally return `None`. A failed journal write WARNS and prints the line rather than erroring: the workload is already deployed by that point, and erroring would make an operator redeploy a live workload to fix bookkeeping.")
51//! @yah:handoff("DISCOVERED WORK, FIXED IN THIS PASS RATHER THAN FILED. kamaji-bin/src/server.rs:1766 — `deploy_workload`'s `if let Ack{..}` gate records the R852-B4 spec digest; left unconverted after the variant removal it would have kept COMPILING and silently stopped recording digests on every deploy, which is a reconciler regression no test would have caught. (server.rs:1649 is the graceful-upgrade path and correctly stays `Ack`.) Also `attach_hydrate` decorates only a `DeployAck`, so a restore in front of a deploy that then FAILED measures nothing recoverable. `topology::analyze` is untouched: no I/O, no clock, no `Path` — only the fill path changed, which was this ticket's hard constraint.")
52//! @yah:verify("LEADER RE-RAN EVERY GATE INDEPENDENTLY (@Ashguard:eclipse) rather than accepting the courier's counts, and all four agree. `cargo test --manifest-path oss/yubaba/Cargo.toml -p yah-cloud --lib`: 1190 passed / 0 failed / 4 ignored (courier's pre-edit baseline 1187/0/4, so +3). `cargo test --manifest-path oss/kamaji/Cargo.toml --workspace`: every suite ok, 0 failed (kamaji-bin lib 234, kamaji lib 119, proto 33, plus the rest; baseline exit 0). `cargo check --workspace --all-targets` from the repo root: exit 0, warnings all pre-existing and none in the changed files. `scripts/check-schema-drift.sh` and `scripts/check-workload-spec-ts.sh`: both 'ok, in sync'. NOTE the yah-cloud invocation — a repo-root `cargo test -p yah-cloud` does NOT work; it is not a root workspace member and needs dev-deps.")
53//! @yah:verify("THE END-TO-END SHAPE IS PINNED FROM BOTH ENDS, which matters because a re-encoded line would fail to parse only on the far side of a real restore. `sibling_wire_e2e` asserts the hydrate line survives the postcard wire BYTE-FOR-BYTE. `recovery_journal`'s new test drives `record_deploy_measurement` with `REAL_HYDRATE_LINE` — the same literal `oss/turso-backup/src/bin/hydrate.rs` asserts full-line equality against, the two-sided pin R850-T2 added — and replays it back to {workload, node, 12.25s, 1024 bytes}; that existing assertion is intact. Plus: a deploy that measured nothing writes no file at all, and a non-JSON reply errors naming both the workload and the node.")
54//! @yah:verify("ONE FLAKE, CHECKED RATHER THAN ASSUMED, AND IT IS NOT THIS TICKET'S. `tenant_passway::deploy_arms_the_declared_socket_and_stop_releases_it` failed intermittently mid-run under `--all-features`; it passes 3/3 in isolation, the failing assertion is a post-`Stop` `TcpListener::bind` on a `free_port()`-chosen port, and nothing in this diff touches custody teardown. It is a port-reuse race between parallel passway tests. The leader's own full kamaji workspace re-run was green. Separately, a W298 skew advisory on the leader's re-run named oss/kamaji/crates/kamaji/src/cgroup.rs and oss/yah-base/crates/workload-spec/src/lib.rs as modified mid-run — both peers' files, neither touched by this ticket.")
55//! @yah:gotcha("V10 IS ALREADY DEPLOYED TO A PRODUCTION RAFT VOTER, AHEAD OF THIS TICKET REACHING REVIEW. Reported by @Ashguard:coffee (session:91597c1e, R881-B7) on 2026-09-11 07:00 UTC: us-east-001 now runs kamaji AND yubaba 0.8.38-h5 built from this shared tree, so the V10 bump and the untracked recovery_journal.rs are live there. The node is healthy — /health reports version and kamaji_version both 0.8.38-h5, so the handshake agrees, and all five workloads are Running. TWO CONSEQUENCES. (1) V10 is DEPLOYED BUT UNRELEASED: rolling us-east-001 back to a published release (scripts/roll-node.sh) silently reverts it. (2) Shipping ONE HALF of {kamaji, yubaba} reproduces a real incident — it already happened an hour earlier: V10 kamaji met the node's released V9 yubaba, they misframed ('decode failed: frame too large: 542393671 > 1048576'), and yubaba fell back to its in-process containerd runtime WITH NO LOG LINE AT ALL; that runtime wires no netns, so two noisetable-account deploys came up with no address and no resolver while their service records still advertised 10.128.3.2. Filed as R881-B8. The V10 doc stanza predicted the skew precisely; what it could not predict is that the symptom presents as a NETWORKING bug rather than a protocol error. @Ashguard:coffee added a guard to scripts/hotship.sh refusing a one-sided ship when the tree's ProtocolVersion is ahead of the last `release v*` commit (--allow-proto-skew overrides, --no-restart downgrades to a warning).")
56
57use std::collections::BTreeMap;
58use std::path::{Path, PathBuf};
59
60use anyhow::{Context, Result};
61use chrono::{DateTime, Utc};
62use serde::{Deserialize, Serialize};
63
64/// The helper that produces the measurements this journal ingests today.
65/// Recorded on every line so a second source of restore timings is
66/// distinguishable in the journal without re-reading the code that wrote it.
67pub const HELPER_TURSO_BACKUP_HYDRATE: &str = "turso-backup-hydrate";
68
69/// A measurement older than this is *stale*: reported, with its age, and
70/// flagged in the headline. Not discarded — see the module docs.
71pub const STALE_AFTER_DAYS: i64 = 30;
72
73/// One subject's measured restore, appended to `.yah/cloud/recovery.jsonl`.
74///
75/// `bytes` and `seconds` are the helper's own measured figures, copied
76/// unmodified. `at` is when the measurement was **ingested** — the helper emits
77/// no timestamp of its own, so for a line piped straight out of a restore this
78/// is within seconds of the event, and for a replayed old log it is not. Use
79/// [`RecoveryRecord::from_helper_json_at`] when the real time is known.
80#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
81pub struct RecoveryRecord {
82 pub at: DateTime<Utc>,
83 pub workload: String,
84 pub node: String,
85 /// The durability tier the restore read from, when the operator said so.
86 /// The helper does not print it — it reads `TIER` from its environment and
87 /// emits per-subject `source` object keys instead — so this is `None`
88 /// unless attributed at ingest.
89 #[serde(default, skip_serializing_if = "Option::is_none")]
90 pub tier: Option<String>,
91 pub subject: String,
92 pub bytes: u64,
93 pub seconds: f64,
94 pub helper: String,
95}
96
97/// The helper's line, named exactly as it prints it. Private: the public shape
98/// is [`RecoveryRecord`], and nothing outside this module should have to know
99/// the helper's envelope.
100/// `outcome`, `epoch`, `subjects` and the envelope `bytes`/`seconds` are
101/// deliberately not read: the per-subject `restored` array is the measurement,
102/// and a field this module ignores is one the helper is free to change.
103/// Absent on every non-`hydrated` outcome (`already_populated`,
104/// `nothing_in_the_store`, `refused`) — those are legitimate lines that carry
105/// no measurement, and parse to zero records.
106#[derive(Debug, Deserialize)]
107struct HelperLine {
108 #[serde(default)]
109 restored: Vec<HelperSubject>,
110}
111
112#[derive(Debug, Deserialize)]
113struct HelperSubject {
114 subject: String,
115 bytes: u64,
116 seconds: f64,
117}
118
119impl RecoveryRecord {
120 /// Parse one line of `turso-backup-hydrate` output into one record per
121 /// restored subject, stamped with the current time.
122 ///
123 /// A well-formed line reporting no restore (`already_populated`, a refusal)
124 /// yields an empty vec — that is an answer, not an error. Only a line that
125 /// is not the helper's JSON at all is an `Err`.
126 pub fn from_helper_json(line: &str, workload: &str, node: &str) -> Result<Vec<Self>> {
127 Self::from_helper_json_at(line, workload, node, Utc::now())
128 }
129
130 /// [`from_helper_json`](Self::from_helper_json) with the ingest timestamp
131 /// supplied — for tests, and for ingesting a log whose real time is known.
132 pub fn from_helper_json_at(
133 line: &str,
134 workload: &str,
135 node: &str,
136 at: DateTime<Utc>,
137 ) -> Result<Vec<Self>> {
138 let parsed: HelperLine = serde_json::from_str(line.trim()).with_context(|| {
139 format!(
140 "parsing turso-backup-hydrate output as JSON (expected one object per line, \
141 with a `restored` array): {line:?}"
142 )
143 })?;
144 Ok(parsed
145 .restored
146 .into_iter()
147 .map(|s| Self {
148 at,
149 workload: workload.to_string(),
150 node: node.to_string(),
151 tier: None,
152 subject: s.subject,
153 bytes: s.bytes,
154 seconds: s.seconds,
155 helper: HELPER_TURSO_BACKUP_HYDRATE.to_string(),
156 })
157 .collect())
158 }
159
160 /// Attribute these records to a durability tier the helper did not print.
161 pub fn with_tier(mut self, tier: impl Into<String>) -> Self {
162 self.tier = Some(tier.into());
163 self
164 }
165}
166
167/// Every measured subject for one workload, plus the clock reading that dates
168/// them.
169///
170/// `as_of` is captured at replay so that [`age_days`](Self::age_days) is a pure
171/// function of already-read data: `topology::analyze` must not read a clock any
172/// more than it may read a file.
173#[derive(Debug, Clone, PartialEq)]
174pub struct WorkloadRecovery {
175 pub workload: String,
176 /// The node the most recent measurement was taken on. Measurements taken on
177 /// different machines are all kept; this names the latest, because a
178 /// restore time is a property of a host as much as of a database.
179 pub node: String,
180 /// Last record per subject, ordered by subject name.
181 pub subjects: BTreeMap<String, RecoveryRecord>,
182 /// Wall clock at replay time.
183 pub as_of: DateTime<Utc>,
184}
185
186impl WorkloadRecovery {
187 /// Total measured wall-clock seconds: the workload's restore is all of its
188 /// subjects, and the helper measures them one at a time.
189 pub fn seconds(&self) -> f64 {
190 self.subjects.values().map(|r| r.seconds).sum()
191 }
192
193 /// Total measured bytes on disk after the restore.
194 pub fn bytes(&self) -> u64 {
195 self.subjects.values().map(|r| r.bytes).sum()
196 }
197
198 /// When the summed figure was measured — the **oldest** component, because
199 /// a sum is only as fresh as its stalest part.
200 pub fn measured_at(&self) -> DateTime<Utc> {
201 self.subjects
202 .values()
203 .map(|r| r.at)
204 .min()
205 .unwrap_or(self.as_of)
206 }
207
208 /// Whole days between [`measured_at`](Self::measured_at) and `as_of`,
209 /// floored at zero (a clock that went backwards is not negative age).
210 pub fn age_days(&self) -> i64 {
211 (self.as_of - self.measured_at()).num_days().max(0)
212 }
213
214 /// Whether the headline must say the figure may no longer describe the
215 /// declared state.
216 pub fn is_stale(&self) -> bool {
217 self.age_days() > STALE_AFTER_DAYS
218 }
219
220 /// How many subjects the summed figure covers.
221 pub fn subject_count(&self) -> usize {
222 self.subjects.len()
223 }
224
225 /// The helper that produced the most recent record.
226 pub fn helper(&self) -> &str {
227 self.subjects
228 .values()
229 .max_by_key(|r| r.at)
230 .map(|r| r.helper.as_str())
231 .unwrap_or(HELPER_TURSO_BACKUP_HYDRATE)
232 }
233}
234
235/// Append-only JSONL journal at `.yah/cloud/recovery.jsonl`.
236pub struct RecoveryJournal {
237 path: PathBuf,
238}
239
240impl RecoveryJournal {
241 pub fn new(path: PathBuf) -> Self {
242 Self { path }
243 }
244
245 /// Journal rooted at `<workspace_root>/.yah/cloud/recovery.jsonl`.
246 pub fn at_workspace(workspace_root: &Path) -> Self {
247 Self::new(crate::paths::recovery_journal(workspace_root))
248 }
249
250 /// Path to the journal file on disk.
251 pub fn path(&self) -> &Path {
252 &self.path
253 }
254
255 /// Append records, one JSONL line each.
256 ///
257 /// Unlike [`crate::asset_journal::AssetStatusJournal::append`] this
258 /// propagates its error: that one is a best-effort side note during a
259 /// reconcile, whereas this is the entire point of the command the operator
260 /// ran, and a measurement silently not written is one nobody takes again.
261 pub fn append(&self, records: &[RecoveryRecord]) -> Result<()> {
262 use std::io::Write;
263
264 if records.is_empty() {
265 return Ok(());
266 }
267 if let Some(parent) = self.path.parent() {
268 std::fs::create_dir_all(parent)
269 .with_context(|| format!("creating {}", parent.display()))?;
270 }
271 let mut buf = String::new();
272 for r in records {
273 buf.push_str(&serde_json::to_string(r).context("serializing RecoveryRecord")?);
274 buf.push('\n');
275 }
276 let mut file = std::fs::OpenOptions::new()
277 .create(true)
278 .append(true)
279 .open(&self.path)
280 .with_context(|| format!("opening {}", self.path.display()))?;
281 file.write_all(buf.as_bytes())
282 .with_context(|| format!("writing to {}", self.path.display()))?;
283 Ok(())
284 }
285
286 /// Replay the journal into one entry per workload, last record winning per
287 /// `(workload, subject)`.
288 ///
289 /// Never errors: an absent journal — the state of every camp that has never
290 /// timed a restore — replays to an empty map, exactly like an unsynced
291 /// infra source overlays nothing in `CloudConfig::load`. An unparseable
292 /// line is skipped with a warning rather than failing the whole load.
293 pub fn replay(&self) -> BTreeMap<String, WorkloadRecovery> {
294 self.replay_as_of(Utc::now())
295 }
296
297 /// [`replay`](Self::replay) with the clock supplied, so ages are
298 /// deterministic in tests.
299 pub fn replay_as_of(&self, now: DateTime<Utc>) -> BTreeMap<String, WorkloadRecovery> {
300 let content = match std::fs::read_to_string(&self.path) {
301 Ok(s) => s,
302 Err(e) if e.kind() == std::io::ErrorKind::NotFound => return BTreeMap::new(),
303 Err(e) => {
304 tracing::debug!(
305 error = %e,
306 journal = %self.path.display(),
307 "recovery journal unreadable — returning empty map",
308 );
309 return BTreeMap::new();
310 }
311 };
312
313 let mut out: BTreeMap<String, WorkloadRecovery> = BTreeMap::new();
314 for (lineno, line) in content.lines().enumerate() {
315 let line = line.trim();
316 if line.is_empty() {
317 continue;
318 }
319 let record: RecoveryRecord = match serde_json::from_str(line) {
320 Ok(r) => r,
321 Err(e) => {
322 tracing::warn!(
323 line = lineno + 1,
324 journal = %self.path.display(),
325 error = %e,
326 "skipping unparseable recovery journal line",
327 );
328 continue;
329 }
330 };
331 let entry = out
332 .entry(record.workload.clone())
333 .or_insert_with(|| WorkloadRecovery {
334 workload: record.workload.clone(),
335 node: record.node.clone(),
336 subjects: BTreeMap::new(),
337 as_of: now,
338 });
339 // Last line wins per subject.
340 entry.subjects.insert(record.subject.clone(), record);
341 }
342 // The most recently measured surviving record names the node — done
343 // after the replay so a superseded line cannot leave its host behind.
344 for entry in out.values_mut() {
345 if let Some(latest) = entry.subjects.values().max_by_key(|r| r.at) {
346 entry.node = latest.node.clone();
347 }
348 }
349 out
350 }
351}
352
353/// Journal the measurement a deploy just came back with (R850-T4).
354///
355/// This is the automatic half of the seam `yah cloud topology --record-hydrate`
356/// is the manual half of: same parser, same journal, same record shape — the
357/// only difference is that the line arrives on the deploy reply instead of
358/// being copied out of a node's log by a human who remembered to look.
359///
360/// `hydrate` is whatever `WorkloadDeployResponse.hydrate` carried, so `None`
361/// (no restore happened) and a line that restored nothing both end as
362/// `Ok(0)` — an answer, not an error. The whole point is that a caller can
363/// wire this into every deploy without first knowing whether the workload is
364/// durable.
365///
366/// **Attribution happens here and can happen nowhere earlier.** The helper's
367/// JSON names neither the workload nor the node; kamaji knows only a
368/// `WorkloadId`; the node cannot name itself as the camp's machine files do.
369/// The caller that dialed `node` to deploy `workload` is the first place both
370/// facts are in one scope, which is why this takes them as arguments rather
371/// than digging them out of the line.
372///
373/// Errors propagate for the same reason [`RecoveryJournal::append`]'s do: a
374/// restore is a rare, unrepeatable event, and a measurement dropped silently
375/// is one nobody takes again.
376pub fn record_deploy_measurement(
377 workspace_root: &Path,
378 workload: &str,
379 node: &str,
380 hydrate: Option<&str>,
381) -> Result<usize> {
382 let Some(line) = hydrate else {
383 return Ok(0);
384 };
385 let records = RecoveryRecord::from_helper_json(line, workload, node)
386 .with_context(|| format!("recording the hydrate measurement {node} reported for {workload}"))?;
387 RecoveryJournal::at_workspace(workspace_root).append(&records)?;
388 Ok(records.len())
389}
390
391// ── Tests ──────────────────────────────────────────────────────────────────────
392
393#[cfg(test)]
394mod tests {
395 use super::*;
396 use chrono::TimeZone;
397 use tempfile::tempdir;
398
399 /// Byte-identical to what `outcome_to_json` prints for the outcome built in
400 /// `oss/turso-backup/src/bin/hydrate.rs`'s
401 /// `a_hydrated_outcome_reports_measured_bytes_and_seconds` — that test
402 /// asserts full-line equality against this same literal, so the two cannot
403 /// drift without one of them going red.
404 const REAL_HYDRATE_LINE: &str = concat!(
405 r#"{"outcome":"hydrated","epoch":4,"subjects":1,"bytes":1024,"seconds":12.500,"#,
406 r#""restored":[{"subject":"accounts.db","#,
407 r#""source":"wl/acct/accounts.db/snapshots/snapshot-000.db","#,
408 r#""bytes":1024,"seconds":12.250}]}"#,
409 );
410
411 fn at(days_ago: i64) -> DateTime<Utc> {
412 now() - chrono::Duration::days(days_ago)
413 }
414
415 fn now() -> DateTime<Utc> {
416 Utc.with_ymd_and_hms(2026, 9, 10, 12, 0, 0).unwrap()
417 }
418
419 fn record(workload: &str, subject: &str, bytes: u64, seconds: f64, days_ago: i64) -> RecoveryRecord {
420 RecoveryRecord {
421 at: at(days_ago),
422 workload: workload.to_string(),
423 node: "us-west-001".to_string(),
424 tier: None,
425 subject: subject.to_string(),
426 bytes,
427 seconds,
428 helper: HELPER_TURSO_BACKUP_HYDRATE.to_string(),
429 }
430 }
431
432 #[test]
433 fn a_real_hydrate_line_parses_verbatim() {
434 let records =
435 RecoveryRecord::from_helper_json_at(REAL_HYDRATE_LINE, "acct", "us-west-001", now())
436 .unwrap();
437 assert_eq!(records.len(), 1);
438 let r = &records[0];
439 assert_eq!(r.subject, "accounts.db");
440 assert_eq!(r.bytes, 1024);
441 assert_eq!(r.seconds, 12.25);
442 assert_eq!(r.workload, "acct");
443 assert_eq!(r.node, "us-west-001");
444 assert_eq!(r.helper, HELPER_TURSO_BACKUP_HYDRATE);
445 // The helper prints no tier; nothing invents one.
446 assert_eq!(r.tier, None);
447 }
448
449 /// R850-T4, the end-to-end shape: what a deploy reply carries becomes a
450 /// journal entry `yah cloud topology` reads back as a Measured estimate,
451 /// with the workload and the node attached.
452 ///
453 /// The fixture is the helper's REAL emitted line — the same constant
454 /// `oss/turso-backup/src/bin/hydrate.rs` asserts full-line equality against
455 /// — because the producer and the consumer of this line live in crates that
456 /// deliberately do not depend on each other. A hand-written approximation
457 /// would let `outcome_to_json` change shape and leave this green while
458 /// every real restore silently failed to journal.
459 ///
460 /// It is deliberately driven through `record_deploy_measurement` rather
461 /// than `from_helper_json` + `append`: the thing under test is the whole
462 /// automatic path, including that attribution is applied at ingest and that
463 /// replay keys on it.
464 #[test]
465 fn a_measurement_off_the_deploy_wire_lands_as_a_replayable_record() {
466 let dir = tempdir().unwrap();
467
468 // What kamaji put on the DeployAck, verbatim, as the CLI would hand it
469 // over having dialed us-west-001 by name to deploy `acct`.
470 let written =
471 record_deploy_measurement(dir.path(), "acct", "us-west-001", Some(REAL_HYDRATE_LINE))
472 .unwrap();
473 assert_eq!(written, 1, "the line restored one subject");
474
475 let replayed = RecoveryJournal::at_workspace(dir.path()).replay_as_of(now());
476 let acct = replayed
477 .get("acct")
478 .expect("the workload the deploy named must be replayable by that name");
479 assert_eq!(acct.node, "us-west-001", "the node must survive the round trip");
480 assert_eq!(acct.subject_count(), 1);
481 assert_eq!(acct.seconds(), 12.25, "the helper's own measured seconds");
482 assert_eq!(acct.bytes(), 1024);
483 assert_eq!(acct.helper(), HELPER_TURSO_BACKUP_HYDRATE);
484 assert!(!acct.is_stale(), "a measurement taken just now is not stale");
485 }
486
487 /// A deploy that restored nothing must not write a line, and must not be an
488 /// error either — this runs on EVERY deploy, and the overwhelming majority
489 /// of them place a workload that declares no durability tier at all.
490 #[test]
491 fn a_deploy_that_measured_nothing_writes_no_journal_at_all() {
492 let dir = tempdir().unwrap();
493 let journal = RecoveryJournal::at_workspace(dir.path());
494
495 assert_eq!(
496 record_deploy_measurement(dir.path(), "marketing", "us-west-001", None).unwrap(),
497 0
498 );
499 assert_eq!(
500 record_deploy_measurement(
501 dir.path(),
502 "marketing",
503 "us-west-001",
504 Some(r#"{"outcome":"already_populated"}"#),
505 )
506 .unwrap(),
507 0
508 );
509 assert!(
510 !journal.path().exists(),
511 "nothing was measured, so no journal file should have been created"
512 );
513 assert!(journal.replay_as_of(now()).is_empty());
514 }
515
516 /// A node that answered with something that is not the helper's JSON is an
517 /// error the caller sees, not a silently-dropped measurement. The deploy
518 /// itself already succeeded by then, so this cannot mean "fail the deploy"
519 /// — it means the operator is told the measurement did not land.
520 #[test]
521 fn a_reply_that_is_not_the_helpers_json_is_an_error_naming_the_workload() {
522 let dir = tempdir().unwrap();
523 let err = record_deploy_measurement(dir.path(), "acct", "us-west-001", Some("not json"))
524 .unwrap_err();
525 let rendered = format!("{err:#}");
526 assert!(
527 rendered.contains("acct") && rendered.contains("us-west-001"),
528 "the error must name what it failed to record: {rendered}"
529 );
530 }
531
532 #[test]
533 fn an_outcome_that_restored_nothing_is_zero_records_not_an_error() {
534 for line in [
535 r#"{"outcome":"already_populated"}"#,
536 r#"{"outcome":"nothing_in_the_store","epoch":4,"subjects":["a.db"]}"#,
537 r#"{"outcome":"refused","reason":"torn_volume","message":"whatever"}"#,
538 ] {
539 let records = RecoveryRecord::from_helper_json_at(line, "w", "n", now()).unwrap();
540 assert!(records.is_empty(), "{line}");
541 }
542 }
543
544 #[test]
545 fn a_line_that_is_not_the_helpers_json_is_an_error() {
546 let err = RecoveryRecord::from_helper_json_at("hydrating /var/lib/...", "w", "n", now())
547 .unwrap_err();
548 assert!(format!("{err:#}").contains("turso-backup-hydrate"), "{err:#}");
549 }
550
551 #[test]
552 fn an_absent_journal_replays_to_empty() {
553 let tmp = tempdir().unwrap();
554 let j = RecoveryJournal::at_workspace(tmp.path());
555 assert!(!j.path().exists());
556 assert!(j.replay().is_empty());
557 }
558
559 #[test]
560 fn a_workloads_figure_is_the_sum_of_its_subjects() {
561 let tmp = tempdir().unwrap();
562 let j = RecoveryJournal::at_workspace(tmp.path());
563 j.append(&[
564 record("acct", "accounts.db", 1024, 12.25, 1),
565 record("acct", "ledger.db", 2048, 7.75, 1),
566 record("other", "x.db", 16, 1.0, 1),
567 ])
568 .unwrap();
569
570 let map = j.replay_as_of(now());
571 assert_eq!(map.len(), 2);
572 let acct = &map["acct"];
573 assert_eq!(acct.subject_count(), 2);
574 assert_eq!(acct.bytes(), 3072);
575 assert!((acct.seconds() - 20.0).abs() < 1e-9, "{}", acct.seconds());
576 assert_eq!(acct.age_days(), 1);
577 assert!(!acct.is_stale());
578 assert_eq!(acct.node, "us-west-001");
579 }
580
581 #[test]
582 fn the_last_record_for_a_subject_wins_and_nothing_is_pruned() {
583 let tmp = tempdir().unwrap();
584 let j = RecoveryJournal::at_workspace(tmp.path());
585 j.append(&[record("acct", "accounts.db", 1024, 99.0, 40)])
586 .unwrap();
587 j.append(&[record("acct", "accounts.db", 4096, 3.5, 2)])
588 .unwrap();
589
590 // Both lines are still on disk — append-only, nobody prunes.
591 let raw = std::fs::read_to_string(j.path()).unwrap();
592 assert_eq!(raw.lines().count(), 2);
593
594 let acct = &j.replay_as_of(now())["acct"];
595 assert_eq!(acct.subject_count(), 1);
596 assert_eq!(acct.bytes(), 4096);
597 assert!((acct.seconds() - 3.5).abs() < 1e-9);
598 assert_eq!(acct.age_days(), 2);
599 }
600
601 #[test]
602 fn a_sum_is_only_as_fresh_as_its_stalest_subject() {
603 let tmp = tempdir().unwrap();
604 let j = RecoveryJournal::at_workspace(tmp.path());
605 j.append(&[
606 record("acct", "accounts.db", 1024, 12.0, 90),
607 record("acct", "ledger.db", 2048, 8.0, 1),
608 ])
609 .unwrap();
610
611 let acct = &j.replay_as_of(now())["acct"];
612 assert_eq!(acct.age_days(), 90);
613 assert!(acct.is_stale());
614 // Stale, and still reported in full — never discarded.
615 assert!((acct.seconds() - 20.0).abs() < 1e-9);
616 }
617
618 #[test]
619 fn an_unparseable_line_is_skipped_not_fatal() {
620 let tmp = tempdir().unwrap();
621 let j = RecoveryJournal::at_workspace(tmp.path());
622 j.append(&[record("acct", "accounts.db", 1024, 12.0, 1)])
623 .unwrap();
624 {
625 use std::io::Write;
626 let mut f = std::fs::OpenOptions::new()
627 .append(true)
628 .open(j.path())
629 .unwrap();
630 writeln!(f, "{{not json").unwrap();
631 }
632 let map = j.replay_as_of(now());
633 assert_eq!(map["acct"].subject_count(), 1);
634 }
635
636 #[test]
637 fn a_record_round_trips_through_the_journal_line() {
638 let r = record("acct", "accounts.db", 1024, 12.25, 3).with_tier("stream");
639 let line = serde_json::to_string(&r).unwrap();
640 assert!(line.contains("\"tier\":\"stream\""), "{line}");
641 assert_eq!(serde_json::from_str::<RecoveryRecord>(&line).unwrap(), r);
642 }
643}