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        }
237    }
238
239    /// In-memory fake provider for drift-logic tests. Configurable
240    /// per-name server presence and per-bucket existence.
241    #[derive(Default)]
242    struct FakeProvider {
243        servers: Mutex<Vec<(String, ServerSummary)>>,
244        buckets: Mutex<Vec<String>>,               // present bucket names
245        bucket_check_fails: Mutex<Option<String>>, // if Some(reason), bucket_exists errors
246    }
247
248    impl FakeProvider {
249        fn with_server(self, name: &str, summary: ServerSummary) -> Self {
250            self.servers.lock().unwrap().push((name.into(), summary));
251            self
252        }
253        fn with_bucket(self, name: &str) -> Self {
254            self.buckets.lock().unwrap().push(name.into());
255            self
256        }
257        fn fail_bucket_check(self, reason: &str) -> Self {
258            *self.bucket_check_fails.lock().unwrap() = Some(reason.into());
259            self
260        }
261    }
262
263    #[async_trait]
264    impl MachineProvider for FakeProvider {
265        async fn ensure_project(&self, name: &str) -> Result<ProjectId> {
266            Ok(ProjectId(name.into()))
267        }
268        async fn create_server(&self, _: &ProjectId, _: &ServerSpec, _: &str) -> Result<ServerId> {
269            bail!("not used in drift tests")
270        }
271        async fn create_bucket(&self, _: &str, _: Location) -> Result<BucketRef> {
272            bail!("not used in drift tests")
273        }
274        async fn server_status(&self, _: &ServerId) -> Result<ServerStatus> {
275            bail!("not used in drift tests")
276        }
277        async fn find_server_by_name(&self, name: &str) -> Result<Option<ServerSummary>> {
278            Ok(self
279                .servers
280                .lock()
281                .unwrap()
282                .iter()
283                .find(|(n, _)| n == name)
284                .map(|(_, s)| s.clone()))
285        }
286        async fn bucket_exists(&self, name: &str, _: Location) -> Result<bool> {
287            if let Some(reason) = self.bucket_check_fails.lock().unwrap().as_deref() {
288                bail!("{reason}");
289            }
290            Ok(self.buckets.lock().unwrap().iter().any(|b| b == name))
291        }
292        async fn destroy_server(&self, _: &ServerId) -> Result<()> {
293            bail!("not used in drift tests")
294        }
295        async fn delete_bucket(&self, _: &str, _: Location) -> Result<()> {
296            bail!("not used in drift tests")
297        }
298        async fn set_bucket_acl(&self, _: &str, _: Location, _: BucketAcl) -> Result<()> {
299            bail!("not used in drift tests")
300        }
301    }
302
303    /// In-memory `AgentProbe` for tests: configurable success/failure.
304    struct FakeProbe {
305        result: std::result::Result<(), String>,
306    }
307
308    #[async_trait]
309    impl AgentProbe for FakeProbe {
310        async fn probe(&self, _: &MachineConfig) -> std::result::Result<(), String> {
311            self.result.clone()
312        }
313    }
314
315    #[tokio::test]
316    async fn no_provider_reports_provider_error() {
317        let report = collect_machine_report(&sample_machine(), None, None).await;
318        assert_eq!(report.headline(), "unknown");
319        assert!(matches!(
320            report.findings.as_slice(),
321            [DriftFinding::ProviderError(_)]
322        ));
323        // Soft finding — doesn't count as drift.
324        assert!(!report.has_drift());
325    }
326
327    #[tokio::test]
328    async fn missing_server_is_drift() {
329        let provider = FakeProvider::default();
330        let report = collect_machine_report(&sample_machine(), Some(&provider), None).await;
331        assert_eq!(report.headline(), "not provisioned");
332        assert!(report.has_drift());
333        assert!(report.findings.contains(&DriftFinding::MissingServer));
334    }
335
336    #[tokio::test]
337    async fn server_present_with_correct_type_and_bucket_is_in_sync() {
338        let provider = FakeProvider::default()
339            .with_server(
340                "noisetable-pdx-1",
341                ServerSummary {
342                    id: ServerId("12345".into()),
343                    server_type: "cpx22".into(),
344                    status: ServerStatus::Running,
345                    public_ipv4: None,
346                    location: "hil".into(),
347                },
348            )
349            .with_bucket("noisetable-assets-pdx-1");
350
351        let report = collect_machine_report(&sample_machine(), Some(&provider), None).await;
352        assert_eq!(report.headline(), "in sync");
353        assert!(!report.has_drift());
354        assert_eq!(report.bucket_present, Some(true));
355        // No probe configured → no AgentUnreachable noise.
356        assert!(!report
357            .findings
358            .iter()
359            .any(|f| matches!(f, DriftFinding::AgentUnreachable { .. })));
360    }
361
362    #[tokio::test]
363    async fn agent_probe_success_emits_no_finding() {
364        let provider = FakeProvider::default()
365            .with_server(
366                "noisetable-pdx-1",
367                ServerSummary {
368                    id: ServerId("1".into()),
369                    server_type: "cpx22".into(),
370                    status: ServerStatus::Running,
371                    public_ipv4: None,
372                    location: "hil".into(),
373                },
374            )
375            .with_bucket("noisetable-assets-pdx-1");
376        let probe = FakeProbe { result: Ok(()) };
377        let report = collect_machine_report(&sample_machine(), Some(&provider), Some(&probe)).await;
378        assert!(!report
379            .findings
380            .iter()
381            .any(|f| matches!(f, DriftFinding::AgentUnreachable { .. })));
382    }
383
384    #[tokio::test]
385    async fn agent_probe_failure_emits_unreachable() {
386        let provider = FakeProvider::default()
387            .with_server(
388                "noisetable-pdx-1",
389                ServerSummary {
390                    id: ServerId("1".into()),
391                    server_type: "cpx22".into(),
392                    status: ServerStatus::Running,
393                    public_ipv4: None,
394                    location: "hil".into(),
395                },
396            )
397            .with_bucket("noisetable-assets-pdx-1");
398        let probe = FakeProbe {
399            result: Err("connection refused".into()),
400        };
401        let report = collect_machine_report(&sample_machine(), Some(&provider), Some(&probe)).await;
402        let unreachable = report
403            .findings
404            .iter()
405            .find(|f| matches!(f, DriftFinding::AgentUnreachable { .. }));
406        assert!(unreachable.is_some());
407        if let Some(DriftFinding::AgentUnreachable { reason }) = unreachable {
408            assert!(reason.contains("connection refused"));
409        }
410        // AgentUnreachable is a soft finding — doesn't promote to drift.
411        assert!(!report.has_drift());
412    }
413
414    #[tokio::test]
415    async fn agent_probe_skipped_when_server_missing() {
416        // Pre-provision: no server → don't bother probing the agent. The
417        // missing-server finding already tells the operator the next step.
418        let provider = FakeProvider::default();
419        let probe = FakeProbe {
420            result: Err("would fail if called".into()),
421        };
422        let report = collect_machine_report(&sample_machine(), Some(&provider), Some(&probe)).await;
423        assert!(!report
424            .findings
425            .iter()
426            .any(|f| matches!(f, DriftFinding::AgentUnreachable { .. })));
427    }
428
429    #[tokio::test]
430    async fn wrong_machine_type_is_drift() {
431        let provider = FakeProvider::default()
432            .with_server(
433                "noisetable-pdx-1",
434                ServerSummary {
435                    id: ServerId("12345".into()),
436                    server_type: "cx22".into(), // declared cpx22
437                    status: ServerStatus::Running,
438                    public_ipv4: None,
439                    location: "hil".into(),
440                },
441            )
442            .with_bucket("noisetable-assets-pdx-1");
443
444        let report = collect_machine_report(&sample_machine(), Some(&provider), None).await;
445        assert!(report.has_drift());
446        assert!(report.findings.iter().any(|f| matches!(
447            f,
448            DriftFinding::WrongMachineType { declared, actual }
449                if declared == "cpx22" && actual == "cx22"
450        )));
451    }
452
453    #[tokio::test]
454    async fn missing_bucket_is_drift() {
455        let provider = FakeProvider::default().with_server(
456            "noisetable-pdx-1",
457            ServerSummary {
458                id: ServerId("1".into()),
459                server_type: "cpx22".into(),
460                status: ServerStatus::Running,
461                public_ipv4: None,
462                location: "hil".into(),
463            },
464        );
465
466        let report = collect_machine_report(&sample_machine(), Some(&provider), None).await;
467        assert!(report.has_drift());
468        assert_eq!(report.bucket_present, Some(false));
469        assert!(report.findings.iter().any(|f| matches!(
470            f,
471            DriftFinding::MissingBucket { name } if name == "noisetable-assets-pdx-1"
472        )));
473    }
474
475    #[tokio::test]
476    async fn bucket_check_failure_is_soft_finding() {
477        let provider = FakeProvider::default()
478            .with_server(
479                "noisetable-pdx-1",
480                ServerSummary {
481                    id: ServerId("1".into()),
482                    server_type: "cpx22".into(),
483                    status: ServerStatus::Running,
484                    public_ipv4: None,
485                    location: "hil".into(),
486                },
487            )
488            .fail_bucket_check("S3 credentials not configured");
489
490        let report = collect_machine_report(&sample_machine(), Some(&provider), None).await;
491        // Real drift only — soft finding doesn't trip.
492        assert!(!report.has_drift());
493        assert!(report
494            .findings
495            .iter()
496            .any(|f| matches!(f, DriftFinding::BucketUnchecked { .. })));
497    }
498
499    #[tokio::test]
500    async fn server_off_is_drift() {
501        let provider = FakeProvider::default()
502            .with_server(
503                "noisetable-pdx-1",
504                ServerSummary {
505                    id: ServerId("1".into()),
506                    server_type: "cpx22".into(),
507                    status: ServerStatus::Off,
508                    public_ipv4: None,
509                    location: "hil".into(),
510                },
511            )
512            .with_bucket("noisetable-assets-pdx-1");
513
514        let report = collect_machine_report(&sample_machine(), Some(&provider), None).await;
515        assert!(report.has_drift());
516        assert!(report
517            .findings
518            .iter()
519            .any(|f| matches!(f, DriftFinding::NotRunning(ServerStatus::Off))));
520    }
521}