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