1use 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#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
24#[serde(rename_all = "kebab-case")]
25pub enum AssetState {
26 Published,
29 PlaceholderOutput,
31 PlaceholderFetch,
33 PinnedNotPublished,
35 DriftBucket,
37 DriftUpstream,
40 TransformBroken,
42 NotRequired,
44}
45
46#[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
63pub 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 pub fn at_workspace(workspace_root: &Path) -> Self {
91 Self::new(crate::paths::asset_status_journal(workspace_root))
92 }
93
94 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 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 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 pub fn subscribe(&self) -> broadcast::Receiver<AssetStatusEvent> {
182 self.tx.subscribe()
183 }
184
185 pub fn path(&self) -> &Path {
187 &self.path
188 }
189}
190
191#[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 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}