Skip to main content

cloud/
status.rs

1//! Drift detection between declared `.yah/cloud/` config and live cloud state.
2//!
3//! Phase 1 scope (R040-T5): server existence + machine type + bucket existence.
4//! Hostkey-fingerprint verification stays stubbed for now; the yubaba-side
5//! `/health` probe is wired via the [`AgentProbe`] trait — callers that want
6//! an agent-reachability check pass an implementation, otherwise no agent
7//! findings are emitted.
8//!
9//! @yah:ticket(R470-T1, "Status journal + replay (.yah/cloud/status.jsonl)")
10//! @yah:assignee(agent:claude)
11//! @yah:at(2026-06-06T21:02:37Z)
12//! @yah:status(review)
13//! @yah:phase(P1)
14//! @yah:parent(R470)
15//! @yah:next("Generalize the drift-only shape in crates/yah/cloud/src/status.rs to an append-only JSONL journal at .yah/cloud/status.jsonl. One record per reconciler decision: {at, asset, from, to, bytes?, blake3?}. Replay-on-load yields last_state_per_asset.")
16//! @yah:next("Wire the reconciler's StaticAssetSyncReport emissions through this journal so every apply produces transition records.")
17//! @yah:next("Implement tail-F-style subscribe for the --watch path: a Stream<Item = StatusEvent> that emits each new record as it's appended.")
18//! @yah:verify("cargo test -p yah-cloud status")
19//! @yah:verify("echo > .yah/cloud/status.jsonl && yah cloud apply --env cloud --service yah-desktop && test -s .yah/cloud/status.jsonl")
20//! @arch:see(.yah/docs/working/W193-asset-dependency-status-surface.md)
21//! @yah:handoff("Delivered: new crates/yah/cloud/src/asset_journal.rs with AssetState (8-state enum, kebab-case serde), AssetStatusEvent, and AssetStatusJournal (append/replay/subscribe). Path helper asset_status_journal() added to paths.rs → .yah/cloud/status.jsonl. Journal registered in lib.rs with pub re-exports. Wired into StaticAssetReconciler.up() → sync_to_r2/sync_to_minio → sync_assets: emits DriftBucket on hash_mismatch, Published (with correct from-state: PlaceholderOutput/PlaceholderFetch/PinnedNotPublished) on successful upload. In-process broadcast::Sender for subscribe(). 9 new asset_journal tests pass; 36 static_asset tests pass; cloud crate checks clean. Pre-existing cloud_init::tests::render_warden_channel... failure unchanged.")
22//! @yah:verify("cargo test -p cloud asset_journal  # 9 pass")
23//! @yah:verify("cargo test -p cloud --lib reconciler::static_asset  # 36 pass")
24//! @yah:verify("cargo check -p cloud  # clean")
25
26use async_trait::async_trait;
27
28use crate::config::MachineConfig;
29use crate::provider::{Location, MachineProvider, ServerStatus, ServerSummary};
30
31/// Hook for probing a per-machine `yah-yubaba`. Implemented in the CLI on top
32/// of `cloud-client` so this crate can stay independent of `reqwest` /
33/// transport concerns.
34#[async_trait]
35pub trait AgentProbe: Send + Sync {
36    /// `Ok(())` if the yubaba answered healthy. `Err(reason)` becomes the
37    /// `AgentUnreachable.reason` string in the report.
38    async fn probe(&self, machine: &MachineConfig) -> std::result::Result<(), String>;
39}
40
41#[derive(Debug, Clone, PartialEq, Eq)]
42pub enum DriftFinding {
43    /// No Hetzner server with the declared name.
44    MissingServer,
45    /// Server exists but its `server_type` doesn't match declared.
46    WrongMachineType { declared: String, actual: String },
47    /// Server exists but isn't Running.
48    NotRunning(ServerStatus),
49    /// Declared bucket isn't present at the location's S3 endpoint.
50    MissingBucket { name: String },
51    /// Live agent's hostkey fingerprint doesn't match declared (A8).
52    HostkeyFingerprintMismatch { declared: String, actual: String },
53
54    // ── Soft findings: surfaced but don't trigger non-zero exit ─────────────
55    /// Couldn't probe the bucket (no S3 creds, transient API error).
56    BucketUnchecked { name: String, reason: String },
57    /// yah-yubaba didn't respond. Always emitted pre-A8 with a "lands with A8" reason.
58    AgentUnreachable { reason: String },
59    /// Hetzner API call failed for this machine; report is incomplete.
60    ProviderError(String),
61}
62
63impl DriftFinding {
64    /// Real divergence between declared and actual cloud state. Soft
65    /// findings (uncheckable conditions, missing creds) return false.
66    pub fn is_drift(&self) -> bool {
67        matches!(
68            self,
69            Self::MissingServer
70                | Self::WrongMachineType { .. }
71                | Self::NotRunning(_)
72                | Self::MissingBucket { .. }
73                | Self::HostkeyFingerprintMismatch { .. }
74        )
75    }
76}
77
78#[derive(Debug, Clone)]
79pub struct MachineReport {
80    pub name: String,
81    pub server: Option<ServerSummary>,
82    pub bucket_present: Option<bool>,
83    pub findings: Vec<DriftFinding>,
84}
85
86impl MachineReport {
87    pub fn has_drift(&self) -> bool {
88        self.findings.iter().any(DriftFinding::is_drift)
89    }
90
91    /// One-word summary for the top of a status block.
92    pub fn headline(&self) -> &'static str {
93        if self
94            .findings
95            .iter()
96            .any(|f| matches!(f, DriftFinding::MissingServer))
97        {
98            "not provisioned"
99        } else if self.has_drift() {
100            "drift"
101        } else if self.server.is_some() {
102            "in sync"
103        } else {
104            "unknown"
105        }
106    }
107}
108
109/// Run drift checks for a single declared machine.
110///
111/// - `provider = None` → cloud credentials aren't available; the report
112///   contains a single `ProviderError` finding describing why.
113/// - `agent = None`    → no yubaba probe attempted; no `AgentUnreachable`
114///   finding is emitted regardless of server state. Pass an `AgentProbe`
115///   when the caller wants reachability surfaced as drift.
116pub async fn collect_machine_report(
117    machine: &MachineConfig,
118    provider: Option<&dyn MachineProvider>,
119    agent: Option<&dyn AgentProbe>,
120) -> MachineReport {
121    let mut findings: Vec<DriftFinding> = Vec::new();
122    let mut server: Option<ServerSummary> = None;
123    let mut bucket_present: Option<bool> = None;
124
125    let location = match Location::try_from(machine.location()) {
126        Ok(l) => Some(l),
127        Err(e) => {
128            findings.push(DriftFinding::ProviderError(format!(
129                "unknown location '{}': {e}",
130                machine.location()
131            )));
132            None
133        }
134    };
135
136    let Some(p) = provider else {
137        findings.push(DriftFinding::ProviderError(
138            "HETZNER_API_TOKEN not set — run `yah cloud secrets` for the contract".into(),
139        ));
140        return MachineReport {
141            name: machine.name.clone(),
142            server,
143            bucket_present,
144            findings,
145        };
146    };
147
148    match p.find_server_by_name(&machine.name).await {
149        Ok(Some(s)) => {
150            if s.server_type != machine.server_type() {
151                findings.push(DriftFinding::WrongMachineType {
152                    declared: machine.server_type().to_string(),
153                    actual: s.server_type.clone(),
154                });
155            }
156            if !s.status.is_running() {
157                findings.push(DriftFinding::NotRunning(s.status.clone()));
158            }
159            server = Some(s);
160        }
161        Ok(None) => findings.push(DriftFinding::MissingServer),
162        Err(e) => findings.push(DriftFinding::ProviderError(format!(
163            "find_server_by_name: {e}"
164        ))),
165    }
166
167    if let (Some(spec), Some(loc)) = (&machine.bucket, &location) {
168        match p.bucket_exists(&spec.name, loc.clone()).await {
169            Ok(true) => bucket_present = Some(true),
170            Ok(false) => {
171                bucket_present = Some(false);
172                findings.push(DriftFinding::MissingBucket {
173                    name: spec.name.clone(),
174                });
175            }
176            Err(e) => findings.push(DriftFinding::BucketUnchecked {
177                name: spec.name.clone(),
178                reason: e.to_string(),
179            }),
180        }
181    }
182
183    // Probe the yubaba iff the caller wired an `AgentProbe` AND there's a
184    // server worth talking to. Pre-provision the agent gap is implied by
185    // `MissingServer` and a second line would just be noise.
186    if let (Some(_), Some(probe)) = (server.as_ref(), agent) {
187        if let Err(reason) = probe.probe(machine).await {
188            findings.push(DriftFinding::AgentUnreachable { reason });
189        }
190    }
191
192    MachineReport {
193        name: machine.name.clone(),
194        server,
195        bucket_present,
196        findings,
197    }
198}
199
200// ─── Tests ───────────────────────────────────────────────────────────────────
201
202#[cfg(test)]
203mod tests {
204    use super::*;
205    use crate::config::BucketSpec;
206    use crate::provider::{BucketAcl, BucketRef, ProjectId, ServerId, ServerSpec};
207    use anyhow::{bail, Result};
208    use async_trait::async_trait;
209    use std::sync::Mutex;
210
211    fn sample_machine() -> MachineConfig {
212        MachineConfig {
213            name: "noisetable-pdx-1".into(),
214            provider: "hetzner".into(),
215            location: Some("pdx".into()),
216            server_type: Some("cpx22".into()),
217            hosts_mirrors: vec!["noisetable".into()],
218            mesh_tags: vec!["region:pdx".into()],
219            region: None,
220            zone: None,
221            arch: None,
222            bucket: Some(BucketSpec {
223                name: "noisetable-assets-pdx-1".into(),
224                public_read: false,
225            }),
226            vendor: None,
227            nickname: None,
228            legacy_hostkey_fingerprint: None,
229            registration: Default::default(),
230            ssh_keys: vec![],
231            cloudflared: None,
232            hosts_operator_bridge: false,
233            connect: None,
234            allocatable: None,
235            taints: vec![],
236            sovereign_group: None,
237            sovereign_role: None,
238            ingress_floating_ip: None,
239        }
240    }
241
242    /// In-memory fake provider for drift-logic tests. Configurable
243    /// per-name server presence and per-bucket existence.
244    #[derive(Default)]
245    struct FakeProvider {
246        servers: Mutex<Vec<(String, ServerSummary)>>,
247        buckets: Mutex<Vec<String>>,               // present bucket names
248        bucket_check_fails: Mutex<Option<String>>, // if Some(reason), bucket_exists errors
249    }
250
251    impl FakeProvider {
252        fn with_server(self, name: &str, summary: ServerSummary) -> Self {
253            self.servers.lock().unwrap().push((name.into(), summary));
254            self
255        }
256        fn with_bucket(self, name: &str) -> Self {
257            self.buckets.lock().unwrap().push(name.into());
258            self
259        }
260        fn fail_bucket_check(self, reason: &str) -> Self {
261            *self.bucket_check_fails.lock().unwrap() = Some(reason.into());
262            self
263        }
264    }
265
266    #[async_trait]
267    impl MachineProvider for FakeProvider {
268        async fn ensure_project(&self, name: &str) -> Result<ProjectId> {
269            Ok(ProjectId(name.into()))
270        }
271        async fn create_server(&self, _: &ProjectId, _: &ServerSpec, _: &str) -> Result<ServerId> {
272            bail!("not used in drift tests")
273        }
274        async fn create_bucket(&self, _: &str, _: Location) -> Result<BucketRef> {
275            bail!("not used in drift tests")
276        }
277        async fn server_status(&self, _: &ServerId) -> Result<ServerStatus> {
278            bail!("not used in drift tests")
279        }
280        async fn find_server_by_name(&self, name: &str) -> Result<Option<ServerSummary>> {
281            Ok(self
282                .servers
283                .lock()
284                .unwrap()
285                .iter()
286                .find(|(n, _)| n == name)
287                .map(|(_, s)| s.clone()))
288        }
289        async fn bucket_exists(&self, name: &str, _: Location) -> Result<bool> {
290            if let Some(reason) = self.bucket_check_fails.lock().unwrap().as_deref() {
291                bail!("{reason}");
292            }
293            Ok(self.buckets.lock().unwrap().iter().any(|b| b == name))
294        }
295        async fn destroy_server(&self, _: &ServerId) -> Result<()> {
296            bail!("not used in drift tests")
297        }
298        async fn delete_bucket(&self, _: &str, _: Location) -> Result<()> {
299            bail!("not used in drift tests")
300        }
301        async fn set_bucket_acl(&self, _: &str, _: Location, _: BucketAcl) -> Result<()> {
302            bail!("not used in drift tests")
303        }
304    }
305
306    /// In-memory `AgentProbe` for tests: configurable success/failure.
307    struct FakeProbe {
308        result: std::result::Result<(), String>,
309    }
310
311    #[async_trait]
312    impl AgentProbe for FakeProbe {
313        async fn probe(&self, _: &MachineConfig) -> std::result::Result<(), String> {
314            self.result.clone()
315        }
316    }
317
318    #[tokio::test]
319    async fn no_provider_reports_provider_error() {
320        let report = collect_machine_report(&sample_machine(), None, None).await;
321        assert_eq!(report.headline(), "unknown");
322        assert!(matches!(
323            report.findings.as_slice(),
324            [DriftFinding::ProviderError(_)]
325        ));
326        // Soft finding — doesn't count as drift.
327        assert!(!report.has_drift());
328    }
329
330    #[tokio::test]
331    async fn missing_server_is_drift() {
332        let provider = FakeProvider::default();
333        let report = collect_machine_report(&sample_machine(), Some(&provider), None).await;
334        assert_eq!(report.headline(), "not provisioned");
335        assert!(report.has_drift());
336        assert!(report.findings.contains(&DriftFinding::MissingServer));
337    }
338
339    #[tokio::test]
340    async fn server_present_with_correct_type_and_bucket_is_in_sync() {
341        let provider = FakeProvider::default()
342            .with_server(
343                "noisetable-pdx-1",
344                ServerSummary {
345                    id: ServerId("12345".into()),
346                    server_type: "cpx22".into(),
347                    status: ServerStatus::Running,
348                    public_ipv4: None,
349                    location: "hil".into(),
350                },
351            )
352            .with_bucket("noisetable-assets-pdx-1");
353
354        let report = collect_machine_report(&sample_machine(), Some(&provider), None).await;
355        assert_eq!(report.headline(), "in sync");
356        assert!(!report.has_drift());
357        assert_eq!(report.bucket_present, Some(true));
358        // No probe configured → no AgentUnreachable noise.
359        assert!(!report
360            .findings
361            .iter()
362            .any(|f| matches!(f, DriftFinding::AgentUnreachable { .. })));
363    }
364
365    #[tokio::test]
366    async fn agent_probe_success_emits_no_finding() {
367        let provider = FakeProvider::default()
368            .with_server(
369                "noisetable-pdx-1",
370                ServerSummary {
371                    id: ServerId("1".into()),
372                    server_type: "cpx22".into(),
373                    status: ServerStatus::Running,
374                    public_ipv4: None,
375                    location: "hil".into(),
376                },
377            )
378            .with_bucket("noisetable-assets-pdx-1");
379        let probe = FakeProbe { result: Ok(()) };
380        let report = collect_machine_report(&sample_machine(), Some(&provider), Some(&probe)).await;
381        assert!(!report
382            .findings
383            .iter()
384            .any(|f| matches!(f, DriftFinding::AgentUnreachable { .. })));
385    }
386
387    #[tokio::test]
388    async fn agent_probe_failure_emits_unreachable() {
389        let provider = FakeProvider::default()
390            .with_server(
391                "noisetable-pdx-1",
392                ServerSummary {
393                    id: ServerId("1".into()),
394                    server_type: "cpx22".into(),
395                    status: ServerStatus::Running,
396                    public_ipv4: None,
397                    location: "hil".into(),
398                },
399            )
400            .with_bucket("noisetable-assets-pdx-1");
401        let probe = FakeProbe {
402            result: Err("connection refused".into()),
403        };
404        let report = collect_machine_report(&sample_machine(), Some(&provider), Some(&probe)).await;
405        let unreachable = report
406            .findings
407            .iter()
408            .find(|f| matches!(f, DriftFinding::AgentUnreachable { .. }));
409        assert!(unreachable.is_some());
410        if let Some(DriftFinding::AgentUnreachable { reason }) = unreachable {
411            assert!(reason.contains("connection refused"));
412        }
413        // AgentUnreachable is a soft finding — doesn't promote to drift.
414        assert!(!report.has_drift());
415    }
416
417    #[tokio::test]
418    async fn agent_probe_skipped_when_server_missing() {
419        // Pre-provision: no server → don't bother probing the agent. The
420        // missing-server finding already tells the operator the next step.
421        let provider = FakeProvider::default();
422        let probe = FakeProbe {
423            result: Err("would fail if called".into()),
424        };
425        let report = collect_machine_report(&sample_machine(), Some(&provider), Some(&probe)).await;
426        assert!(!report
427            .findings
428            .iter()
429            .any(|f| matches!(f, DriftFinding::AgentUnreachable { .. })));
430    }
431
432    #[tokio::test]
433    async fn wrong_machine_type_is_drift() {
434        let provider = FakeProvider::default()
435            .with_server(
436                "noisetable-pdx-1",
437                ServerSummary {
438                    id: ServerId("12345".into()),
439                    server_type: "cx22".into(), // declared cpx22
440                    status: ServerStatus::Running,
441                    public_ipv4: None,
442                    location: "hil".into(),
443                },
444            )
445            .with_bucket("noisetable-assets-pdx-1");
446
447        let report = collect_machine_report(&sample_machine(), Some(&provider), None).await;
448        assert!(report.has_drift());
449        assert!(report.findings.iter().any(|f| matches!(
450            f,
451            DriftFinding::WrongMachineType { declared, actual }
452                if declared == "cpx22" && actual == "cx22"
453        )));
454    }
455
456    #[tokio::test]
457    async fn missing_bucket_is_drift() {
458        let provider = FakeProvider::default().with_server(
459            "noisetable-pdx-1",
460            ServerSummary {
461                id: ServerId("1".into()),
462                server_type: "cpx22".into(),
463                status: ServerStatus::Running,
464                public_ipv4: None,
465                location: "hil".into(),
466            },
467        );
468
469        let report = collect_machine_report(&sample_machine(), Some(&provider), None).await;
470        assert!(report.has_drift());
471        assert_eq!(report.bucket_present, Some(false));
472        assert!(report.findings.iter().any(|f| matches!(
473            f,
474            DriftFinding::MissingBucket { name } if name == "noisetable-assets-pdx-1"
475        )));
476    }
477
478    #[tokio::test]
479    async fn bucket_check_failure_is_soft_finding() {
480        let provider = FakeProvider::default()
481            .with_server(
482                "noisetable-pdx-1",
483                ServerSummary {
484                    id: ServerId("1".into()),
485                    server_type: "cpx22".into(),
486                    status: ServerStatus::Running,
487                    public_ipv4: None,
488                    location: "hil".into(),
489                },
490            )
491            .fail_bucket_check("S3 credentials not configured");
492
493        let report = collect_machine_report(&sample_machine(), Some(&provider), None).await;
494        // Real drift only — soft finding doesn't trip.
495        assert!(!report.has_drift());
496        assert!(report
497            .findings
498            .iter()
499            .any(|f| matches!(f, DriftFinding::BucketUnchecked { .. })));
500    }
501
502    #[tokio::test]
503    async fn server_off_is_drift() {
504        let provider = FakeProvider::default()
505            .with_server(
506                "noisetable-pdx-1",
507                ServerSummary {
508                    id: ServerId("1".into()),
509                    server_type: "cpx22".into(),
510                    status: ServerStatus::Off,
511                    public_ipv4: None,
512                    location: "hil".into(),
513                },
514            )
515            .with_bucket("noisetable-assets-pdx-1");
516
517        let report = collect_machine_report(&sample_machine(), Some(&provider), None).await;
518        assert!(report.has_drift());
519        assert!(report
520            .findings
521            .iter()
522            .any(|f| matches!(f, DriftFinding::NotRunning(ServerStatus::Off))));
523    }
524}