Skip to main content

cloud/
asset_journal.rs

1//! Append-only JSONL journal for static-asset reconciler decisions (R470-T1).
2//!
3//! Path: `.yah/cloud/status.jsonl`. One record per reconciler decision.
4//! Replay yields `last_state_per_asset` — the source of truth for
5//! `yah cloud status` without `--check`. In-process `subscribe()` lets
6//! the desktop panel react to transitions without polling.
7//!
8//! Wire shape matches `.yah/events.jsonl`: append-only, concurrent writers
9//! are safe (each `append` is one `write` syscall on the line), replayable
10//! to a HashMap keyed by `"<service>:<filename>"`.
11
12use std::collections::HashMap;
13use std::path::{Path, PathBuf};
14use std::sync::Arc;
15
16use anyhow::{Context, Result};
17use chrono::{DateTime, Utc};
18use serde::{Deserialize, Serialize};
19use tokio::sync::broadcast;
20
21/// The 8 canonical asset states. All surfaces (CLI, panel, JSON, RPC)
22/// use these exact kebab-case strings — lock this vocabulary before extending.
23#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
24#[serde(rename_all = "kebab-case")]
25pub enum AssetState {
26    /// Catalog blake3 non-zero; HEAD against bucket matches; fetch.blake3 (if
27    /// present) non-zero. No operator action needed.
28    Published,
29    /// `[[asset]].blake3` is the 64-zero sentinel. Run `yah cloud apply`.
30    PlaceholderOutput,
31    /// `[asset.derive.fetch].blake3` is the 64-zero sentinel. Run `yah cloud apply`.
32    PlaceholderFetch,
33    /// Catalog blake3 non-zero; HEAD misses (bucket key absent). Run `yah cloud apply`.
34    PinnedNotPublished,
35    /// HEAD exists but its recomputed hash ≠ catalog blake3. Investigate first.
36    DriftBucket,
37    /// Fetch URL produced a different blake3 than the pinned fetch.blake3.
38    /// Regenerate transform output and open a new blake3 cycle.
39    DriftUpstream,
40    /// Recipe execution failed during last reconcile. Inspect logs; fix recipe.
41    TransformBroken,
42    /// Declared optional; current host/cfg predicate is false. No action needed.
43    NotRequired,
44}
45
46/// One state transition appended to `.yah/cloud/status.jsonl`.
47///
48/// `asset` is `"<service>:<filename>"`, e.g.
49/// `"yah-desktop:yah-desktop/whisper/distil-large-v3-q5_1.bin"`.
50#[derive(Debug, Clone, Serialize, Deserialize)]
51pub struct AssetStatusEvent {
52    pub at: DateTime<Utc>,
53    pub asset: String,
54    #[serde(skip_serializing_if = "Option::is_none")]
55    pub from: Option<AssetState>,
56    pub to: AssetState,
57    #[serde(skip_serializing_if = "Option::is_none")]
58    pub bytes: Option<u64>,
59    #[serde(skip_serializing_if = "Option::is_none")]
60    pub blake3: Option<String>,
61}
62
63/// Append-only JSONL journal at `.yah/cloud/status.jsonl`.
64///
65/// **Write**: `append` serialises one event as a JSONL line. Non-fatal — a
66/// write failure is logged and the reconciler continues.
67///
68/// **Read**: `last_state_per_asset` replays the journal, returning a
69/// `HashMap<asset_key, AssetState>` keyed by `"<service>:<filename>"`.
70/// Returns an empty map when the journal file doesn't exist yet.
71///
72/// **Watch**: `subscribe` returns a `broadcast::Receiver` for in-process
73/// callers (desktop panel). Cross-process consumers (the `--watch` CLI flag)
74/// tail the file and parse new JSONL lines directly.
75pub struct AssetStatusJournal {
76    path: PathBuf,
77    tx: Arc<broadcast::Sender<AssetStatusEvent>>,
78}
79
80impl AssetStatusJournal {
81    pub fn new(path: PathBuf) -> Self {
82        let (tx, _) = broadcast::channel(256);
83        Self {
84            path,
85            tx: Arc::new(tx),
86        }
87    }
88
89    /// Create a journal rooted at `<workspace_root>/.yah/cloud/status.jsonl`.
90    pub fn at_workspace(workspace_root: &Path) -> Self {
91        Self::new(crate::paths::asset_status_journal(workspace_root))
92    }
93
94    /// Append one event. Best-effort: failures are warned, not propagated.
95    pub async fn append(&self, event: &AssetStatusEvent) {
96        if let Err(e) = self.try_append(event).await {
97            tracing::warn!(
98                asset = %event.asset,
99                error = %e,
100                "asset status journal write failed (non-fatal)"
101            );
102        } else {
103            // Ignore send errors — no active subscribers is fine.
104            let _ = self.tx.send(event.clone());
105        }
106    }
107
108    async fn try_append(&self, event: &AssetStatusEvent) -> Result<()> {
109        use tokio::io::AsyncWriteExt;
110
111        if let Some(parent) = self.path.parent() {
112            tokio::fs::create_dir_all(parent)
113                .await
114                .with_context(|| format!("creating {}", parent.display()))?;
115        }
116        let mut line = serde_json::to_string(event).context("serializing AssetStatusEvent")?;
117        line.push('\n');
118        let mut file = tokio::fs::OpenOptions::new()
119            .create(true)
120            .append(true)
121            .open(&self.path)
122            .await
123            .with_context(|| format!("opening {}", self.path.display()))?;
124        file.write_all(line.as_bytes())
125            .await
126            .with_context(|| format!("writing to {}", self.path.display()))?;
127        Ok(())
128    }
129
130    /// Replay the journal file and return the most recent state per asset key.
131    /// Returns an empty map when the journal doesn't exist (before first apply).
132    pub async fn last_state_per_asset(&self) -> HashMap<String, AssetState> {
133        match self.try_replay().await {
134            Ok(map) => map,
135            Err(e) => {
136                tracing::debug!(
137                    error = %e,
138                    journal = %self.path.display(),
139                    "journal replay failed — returning empty map",
140                );
141                HashMap::new()
142            }
143        }
144    }
145
146    async fn try_replay(&self) -> Result<HashMap<String, AssetState>> {
147        let content = match tokio::fs::read_to_string(&self.path).await {
148            Ok(s) => s,
149            Err(e) if e.kind() == std::io::ErrorKind::NotFound => return Ok(HashMap::new()),
150            Err(e) => return Err(e).with_context(|| format!("reading {}", self.path.display())),
151        };
152
153        let mut map = HashMap::new();
154        for (lineno, line) in content.lines().enumerate() {
155            let line = line.trim();
156            if line.is_empty() {
157                continue;
158            }
159            match serde_json::from_str::<AssetStatusEvent>(line) {
160                Ok(event) => {
161                    map.insert(event.asset.clone(), event.to);
162                }
163                Err(e) => {
164                    tracing::warn!(
165                        line = lineno + 1,
166                        journal = %self.path.display(),
167                        error = %e,
168                        "skipping unparseable journal line",
169                    );
170                }
171            }
172        }
173        Ok(map)
174    }
175
176    /// Subscribe to in-process state-change events. Each call yields an
177    /// independent receiver; up to 256 events are buffered before the
178    /// oldest is dropped for a lagging receiver.
179    ///
180    /// For cross-process `--watch`, tail the journal file and parse JSONL lines.
181    pub fn subscribe(&self) -> broadcast::Receiver<AssetStatusEvent> {
182        self.tx.subscribe()
183    }
184
185    /// Path to the journal file on disk.
186    pub fn path(&self) -> &Path {
187        &self.path
188    }
189}
190
191// ── Tests ──────────────────────────────────────────────────────────────────────
192
193#[cfg(test)]
194mod tests {
195    use super::*;
196    use tempfile::tempdir;
197
198    fn sample_event(asset: &str, to: AssetState) -> AssetStatusEvent {
199        AssetStatusEvent {
200            at: Utc::now(),
201            asset: asset.to_string(),
202            from: None,
203            to,
204            bytes: None,
205            blake3: None,
206        }
207    }
208
209    #[test]
210    fn asset_state_serde_round_trip() {
211        let cases = [
212            (AssetState::Published, "\"published\""),
213            (AssetState::PlaceholderOutput, "\"placeholder-output\""),
214            (AssetState::PlaceholderFetch, "\"placeholder-fetch\""),
215            (AssetState::PinnedNotPublished, "\"pinned-not-published\""),
216            (AssetState::DriftBucket, "\"drift-bucket\""),
217            (AssetState::DriftUpstream, "\"drift-upstream\""),
218            (AssetState::TransformBroken, "\"transform-broken\""),
219            (AssetState::NotRequired, "\"not-required\""),
220        ];
221        for (state, expected_json) in cases {
222            let json = serde_json::to_string(&state).unwrap();
223            assert_eq!(json, expected_json, "serialize {state:?}");
224            let rt: AssetState = serde_json::from_str(&json).unwrap();
225            assert_eq!(rt, state, "round-trip {state:?}");
226        }
227    }
228
229    #[test]
230    fn event_serde_round_trip() {
231        let event = AssetStatusEvent {
232            at: DateTime::parse_from_rfc3339("2026-06-06T21:00:00Z")
233                .unwrap()
234                .with_timezone(&Utc),
235            asset: "yah-desktop:whisper/model.bin".to_string(),
236            from: Some(AssetState::PinnedNotPublished),
237            to: AssetState::Published,
238            bytes: Some(1024),
239            blake3: Some("abc123".to_string()),
240        };
241        let json = serde_json::to_string(&event).unwrap();
242        let rt: AssetStatusEvent = serde_json::from_str(&json).unwrap();
243        assert_eq!(rt.asset, event.asset);
244        assert_eq!(rt.from, event.from);
245        assert_eq!(rt.to, event.to);
246        assert_eq!(rt.bytes, event.bytes);
247        assert_eq!(rt.blake3, event.blake3);
248    }
249
250    #[test]
251    fn event_skips_none_fields_in_json() {
252        let event = sample_event("svc:file.bin", AssetState::PlaceholderOutput);
253        let json = serde_json::to_string(&event).unwrap();
254        assert!(
255            !json.contains("\"from\""),
256            "None from should be omitted: {json}"
257        );
258        assert!(
259            !json.contains("\"bytes\""),
260            "None bytes should be omitted: {json}"
261        );
262        assert!(
263            !json.contains("\"blake3\""),
264            "None blake3 should be omitted: {json}"
265        );
266    }
267
268    #[tokio::test]
269    async fn append_creates_file_and_writes_jsonl() {
270        let dir = tempdir().unwrap();
271        let journal = AssetStatusJournal::new(dir.path().join("cloud/status.jsonl"));
272
273        let event = AssetStatusEvent {
274            at: Utc::now(),
275            asset: "svc:model.bin".to_string(),
276            from: Some(AssetState::PinnedNotPublished),
277            to: AssetState::Published,
278            bytes: Some(512),
279            blake3: Some("deafbeef".to_string()),
280        };
281        journal.append(&event).await;
282
283        let content = tokio::fs::read_to_string(journal.path()).await.unwrap();
284        assert!(
285            !content.is_empty(),
286            "journal file must be non-empty after append"
287        );
288        let parsed: AssetStatusEvent = serde_json::from_str(content.trim()).unwrap();
289        assert_eq!(parsed.asset, event.asset);
290        assert_eq!(parsed.to, AssetState::Published);
291    }
292
293    #[tokio::test]
294    async fn replay_empty_when_no_file() {
295        let dir = tempdir().unwrap();
296        let journal = AssetStatusJournal::new(dir.path().join("cloud/status.jsonl"));
297        let map = journal.last_state_per_asset().await;
298        assert!(map.is_empty(), "missing journal → empty map");
299    }
300
301    #[tokio::test]
302    async fn replay_yields_last_state_per_asset() {
303        let dir = tempdir().unwrap();
304        let journal = AssetStatusJournal::new(dir.path().join("cloud/status.jsonl"));
305        let asset = "svc:model.bin";
306
307        journal
308            .append(&AssetStatusEvent {
309                at: Utc::now(),
310                asset: asset.to_string(),
311                from: Some(AssetState::PinnedNotPublished),
312                to: AssetState::Published,
313                bytes: None,
314                blake3: None,
315            })
316            .await;
317        journal
318            .append(&AssetStatusEvent {
319                at: Utc::now(),
320                asset: asset.to_string(),
321                from: Some(AssetState::Published),
322                to: AssetState::DriftBucket,
323                bytes: None,
324                blake3: None,
325            })
326            .await;
327
328        let map = journal.last_state_per_asset().await;
329        assert_eq!(map.get(asset), Some(&AssetState::DriftBucket));
330    }
331
332    #[tokio::test]
333    async fn replay_tracks_multiple_assets_independently() {
334        let dir = tempdir().unwrap();
335        let journal = AssetStatusJournal::new(dir.path().join("cloud/status.jsonl"));
336
337        journal
338            .append(&sample_event("svc:a.bin", AssetState::Published))
339            .await;
340        journal
341            .append(&sample_event("svc:b.bin", AssetState::PinnedNotPublished))
342            .await;
343        journal
344            .append(&sample_event("svc:a.bin", AssetState::DriftBucket))
345            .await;
346
347        let map = journal.last_state_per_asset().await;
348        assert_eq!(map.get("svc:a.bin"), Some(&AssetState::DriftBucket));
349        assert_eq!(map.get("svc:b.bin"), Some(&AssetState::PinnedNotPublished));
350    }
351
352    #[tokio::test]
353    async fn replay_skips_blank_lines_without_panic() {
354        let dir = tempdir().unwrap();
355        let path = dir.path().join("status.jsonl");
356        // Write a valid line, a blank, and another valid line.
357        let valid = serde_json::to_string(&sample_event("svc:x.bin", AssetState::Published))
358            .unwrap()
359            + "\n";
360        tokio::fs::write(&path, format!("{valid}\n{valid}"))
361            .await
362            .unwrap();
363        let journal = AssetStatusJournal::new(path);
364        let map = journal.last_state_per_asset().await;
365        assert_eq!(map.get("svc:x.bin"), Some(&AssetState::Published));
366    }
367
368    #[tokio::test]
369    async fn subscribe_receives_appended_events() {
370        let dir = tempdir().unwrap();
371        let journal = AssetStatusJournal::new(dir.path().join("cloud/status.jsonl"));
372        let mut rx = journal.subscribe();
373
374        let event = sample_event("svc:w.bin", AssetState::Published);
375        journal.append(&event).await;
376
377        let received = rx.try_recv().expect("event must be received");
378        assert_eq!(received.asset, event.asset);
379        assert_eq!(received.to, AssetState::Published);
380    }
381}