Skip to main content

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}