1use std::collections::BTreeSet;
24use std::path::{Path, PathBuf};
25
26use anyhow::Result;
27use tracing::{info, warn};
28
29use crate::reconciler::static_asset::WORKLOAD_KIND;
30use crate::ServiceConfig;
31
32#[derive(Debug, Default)]
35pub struct DeriveCacheLiveHashes {
36 pub fetch: BTreeSet<String>,
38 pub transform: BTreeSet<String>,
41}
42
43#[derive(Debug, Clone)]
45pub struct DerivePruneCandidate {
46 pub path: PathBuf,
48 pub cache_kind: &'static str,
50 pub hash: String,
52 pub size: u64,
54}
55
56pub fn collect_live_derive_hashes<'a>(
63 workspace_root: &Path,
64 services: impl IntoIterator<Item = &'a ServiceConfig>,
65) -> Result<DeriveCacheLiveHashes> {
66 let mut live = DeriveCacheLiveHashes::default();
67
68 for svc in services {
69 for component in &svc.components {
70 if component.kind != WORKLOAD_KIND {
71 continue;
72 }
73 let workload_path = workspace_root.join(&component.path).join("workload.toml");
74 let src = match std::fs::read_to_string(&workload_path) {
75 Ok(s) => s,
76 Err(e) => {
77 warn!(
78 path = %workload_path.display(),
79 error = %e,
80 "derive-cache prune: skipping unreadable workload",
81 );
82 continue;
83 }
84 };
85 let envelope: workload_spec::Workload = match toml::from_str(&src) {
86 Ok(w) => w,
87 Err(e) => {
88 warn!(
89 path = %workload_path.display(),
90 error = %e,
91 "derive-cache prune: skipping unparseable workload",
92 );
93 continue;
94 }
95 };
96 let workload_spec::Workload::StaticAsset(workload) = envelope else {
97 continue;
98 };
99 for entry in &workload.assets {
100 let Some(derive) = &entry.derive else {
101 continue;
102 };
103 live.fetch.insert(derive.fetch.blake3.0.clone());
104 if derive.transform.is_some() {
105 live.transform.insert(entry.blake3.0.clone());
106 }
107 }
108 }
109 }
110
111 Ok(live)
112}
113
114pub fn compute_derive_cache_candidates(
120 derive_cache_root: &Path,
121 live: &DeriveCacheLiveHashes,
122) -> Result<Vec<DerivePruneCandidate>> {
123 let mut candidates = Vec::new();
124
125 for (subdir, live_set) in [("fetch", &live.fetch), ("transform", &live.transform)] {
126 let dir = derive_cache_root.join(subdir);
127 let entries = match std::fs::read_dir(&dir) {
128 Ok(e) => e,
129 Err(_) => continue, };
131 for entry in entries.flatten() {
132 let path = entry.path();
133 let ext = path.extension().and_then(|e| e.to_str());
134 if ext != Some("bin") {
135 continue; }
137 let Some(hash) = path
138 .file_stem()
139 .and_then(|s| s.to_str())
140 .map(|s| s.to_string())
141 else {
142 continue;
143 };
144 if live_set.contains(hash.as_str()) {
145 continue;
146 }
147 let size = entry.metadata().map(|m| m.len()).unwrap_or(0);
148 candidates.push(DerivePruneCandidate {
149 path,
150 cache_kind: if subdir == "fetch" {
151 "fetch"
152 } else {
153 "transform"
154 },
155 hash,
156 size,
157 });
158 }
159 }
160
161 candidates.sort_by(|a, b| a.path.cmp(&b.path));
163 Ok(candidates)
164}
165
166pub fn execute_derive_cache_prune(candidates: &[DerivePruneCandidate]) -> (usize, u64, usize) {
168 let mut deleted = 0usize;
169 let mut bytes = 0u64;
170 let mut errors = 0usize;
171
172 for c in candidates {
173 match std::fs::remove_file(&c.path) {
174 Ok(()) => {
175 info!(path = %c.path.display(), "pruned derive cache entry");
176 deleted += 1;
177 bytes += c.size;
178 }
179 Err(e) => {
180 warn!(path = %c.path.display(), error = %e, "failed to prune derive cache entry");
181 errors += 1;
182 }
183 }
184 }
185
186 (deleted, bytes, errors)
187}
188
189#[cfg(test)]
190mod tests {
191 use super::*;
192 use crate::config::{ServiceComponent, ServiceConfig};
193 use tempfile::TempDir;
194
195 fn svc(name: &str, components: Vec<ServiceComponent>) -> ServiceConfig {
196 ServiceConfig {
197 schema_version: 1,
198 name: name.into(),
199 domain: String::new(),
200 health_path: None,
201 components,
202 db: crate::DbCatalog::default(),
203 }
204 }
205
206 fn static_asset_component(id: &str, path: &str) -> ServiceComponent {
207 ServiceComponent {
208 mount: None,
209 id: id.into(),
210 kind: "static-asset".into(),
211 path: path.into(),
212 role: String::new(),
213 publishes: None,
214 wave: 0,
215 git: None,
216 deploy: Default::default(),
217 }
218 }
219
220 fn write_workload(dir: &Path, toml: &str) {
221 if let Some(p) = dir.parent() {
222 std::fs::create_dir_all(p).unwrap();
223 }
224 std::fs::create_dir_all(dir).unwrap();
225 std::fs::write(dir.join("workload.toml"), toml).unwrap();
226 }
227
228 const FETCH_HASH: &str = "aaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaa";
229 const OUTPUT_HASH: &str = "bbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbb";
230
231 fn workload_with_fetch_only(fetch_hash: &str, output_hash: &str) -> String {
232 format!(
233 r#"kind = "static-asset"
234schema_version = "V1"
235
236[[asset]]
237filename = "model.bin"
238blake3 = "{output_hash}"
239
240[asset.derive.fetch]
241url = "https://example.com/model.bin"
242blake3 = "{fetch_hash}"
243license = "mit"
244"#
245 )
246 }
247
248 fn workload_with_fetch_and_transform(fetch_hash: &str, output_hash: &str) -> String {
249 format!(
250 r#"kind = "static-asset"
251schema_version = "V1"
252
253[[asset]]
254filename = "model.bin"
255blake3 = "{output_hash}"
256
257[asset.derive.fetch]
258url = "https://example.com/model.bin"
259blake3 = "{fetch_hash}"
260license = "mit"
261
262[asset.derive.transform]
263recipe = "quantize"
264"#
265 )
266 }
267
268 #[test]
269 fn collect_live_hashes_fetch_only() {
270 let tmp = TempDir::new().unwrap();
271 let workload_dir = tmp.path().join("app/web");
272 write_workload(
273 &workload_dir,
274 &workload_with_fetch_only(FETCH_HASH, OUTPUT_HASH),
275 );
276
277 let svc = svc("my-svc", vec![static_asset_component("site", "app/web")]);
278 let live = collect_live_derive_hashes(tmp.path(), [&svc]).unwrap();
279
280 assert!(live.fetch.contains(FETCH_HASH), "fetch hash must be live");
281 assert!(
282 live.transform.is_empty(),
283 "fetch-only asset has no transform cache entry"
284 );
285 }
286
287 #[test]
288 fn collect_live_hashes_fetch_and_transform() {
289 let tmp = TempDir::new().unwrap();
290 let workload_dir = tmp.path().join("app/web");
291 write_workload(
292 &workload_dir,
293 &workload_with_fetch_and_transform(FETCH_HASH, OUTPUT_HASH),
294 );
295
296 let svc = svc("my-svc", vec![static_asset_component("site", "app/web")]);
297 let live = collect_live_derive_hashes(tmp.path(), [&svc]).unwrap();
298
299 assert!(live.fetch.contains(FETCH_HASH));
300 assert!(
301 live.transform.contains(OUTPUT_HASH),
302 "transform cache entry keyed by output blake3"
303 );
304 }
305
306 #[test]
307 fn collect_live_hashes_missing_workload_skips_silently() {
308 let tmp = TempDir::new().unwrap();
309 let svc = svc("my-svc", vec![static_asset_component("site", "app/web")]);
311 let live = collect_live_derive_hashes(tmp.path(), [&svc]).unwrap();
312 assert!(live.fetch.is_empty());
313 assert!(live.transform.is_empty());
314 }
315
316 #[test]
317 fn collect_live_hashes_ignores_non_static_asset_components() {
318 let tmp = TempDir::new().unwrap();
319 let svc = svc(
320 "my-svc",
321 vec![ServiceComponent {
322 mount: None,
323 id: "api".into(),
324 kind: "container".into(),
325 path: "app/api".into(),
326 role: String::new(),
327 publishes: None,
328 wave: 0,
329 git: None,
330 deploy: Default::default(),
331 }],
332 );
333 let live = collect_live_derive_hashes(tmp.path(), [&svc]).unwrap();
334 assert!(live.fetch.is_empty());
335 }
336
337 #[test]
338 fn candidates_exclude_live_hashes() {
339 let tmp = TempDir::new().unwrap();
340 let fetch_dir = tmp.path().join("fetch");
341 std::fs::create_dir_all(&fetch_dir).unwrap();
342
343 std::fs::write(fetch_dir.join(format!("{FETCH_HASH}.bin")), b"live").unwrap();
345 let orphan = "cccccccccccccccccccccccccccccccccccccccccccccccccccccccccccccccc";
347 std::fs::write(fetch_dir.join(format!("{orphan}.bin")), b"orphan").unwrap();
348 std::fs::write(fetch_dir.join(format!("{orphan}.partial")), b"partial").unwrap();
350
351 let mut live = DeriveCacheLiveHashes::default();
352 live.fetch.insert(FETCH_HASH.to_string());
353
354 let candidates = compute_derive_cache_candidates(tmp.path(), &live).unwrap();
355 assert_eq!(candidates.len(), 1, "exactly one orphan");
356 assert_eq!(candidates[0].hash, orphan);
357 assert_eq!(candidates[0].cache_kind, "fetch");
358 }
359
360 #[test]
361 fn execute_prune_deletes_candidates() {
362 let tmp = TempDir::new().unwrap();
363 let fetch_dir = tmp.path().join("fetch");
364 std::fs::create_dir_all(&fetch_dir).unwrap();
365
366 let file = fetch_dir.join("dead.bin");
367 std::fs::write(&file, b"orphan bytes").unwrap();
368 assert!(file.exists());
369
370 let candidate = DerivePruneCandidate {
371 path: file.clone(),
372 cache_kind: "fetch",
373 hash: "dead".into(),
374 size: 12,
375 };
376 let (deleted, bytes, errors) = execute_derive_cache_prune(&[candidate]);
377 assert_eq!(deleted, 1);
378 assert_eq!(bytes, 12);
379 assert_eq!(errors, 0);
380 assert!(!file.exists(), "file must be deleted");
381 }
382}