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 components,
201 db: crate::DbCatalog::default(),
202 }
203 }
204
205 fn static_asset_component(id: &str, path: &str) -> ServiceComponent {
206 ServiceComponent {
207 mount: None,
208 id: id.into(),
209 kind: "static-asset".into(),
210 path: path.into(),
211 role: String::new(),
212 publishes: None,
213 wave: 0,
214 git: None,
215 }
216 }
217
218 fn write_workload(dir: &Path, toml: &str) {
219 if let Some(p) = dir.parent() {
220 std::fs::create_dir_all(p).unwrap();
221 }
222 std::fs::create_dir_all(dir).unwrap();
223 std::fs::write(dir.join("workload.toml"), toml).unwrap();
224 }
225
226 const FETCH_HASH: &str = "aaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaa";
227 const OUTPUT_HASH: &str = "bbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbb";
228
229 fn workload_with_fetch_only(fetch_hash: &str, output_hash: &str) -> String {
230 format!(
231 r#"kind = "static-asset"
232schema_version = "V1"
233
234[[asset]]
235filename = "model.bin"
236blake3 = "{output_hash}"
237
238[asset.derive.fetch]
239url = "https://example.com/model.bin"
240blake3 = "{fetch_hash}"
241license = "mit"
242"#
243 )
244 }
245
246 fn workload_with_fetch_and_transform(fetch_hash: &str, output_hash: &str) -> String {
247 format!(
248 r#"kind = "static-asset"
249schema_version = "V1"
250
251[[asset]]
252filename = "model.bin"
253blake3 = "{output_hash}"
254
255[asset.derive.fetch]
256url = "https://example.com/model.bin"
257blake3 = "{fetch_hash}"
258license = "mit"
259
260[asset.derive.transform]
261recipe = "quantize"
262"#
263 )
264 }
265
266 #[test]
267 fn collect_live_hashes_fetch_only() {
268 let tmp = TempDir::new().unwrap();
269 let workload_dir = tmp.path().join("app/web");
270 write_workload(
271 &workload_dir,
272 &workload_with_fetch_only(FETCH_HASH, OUTPUT_HASH),
273 );
274
275 let svc = svc("my-svc", vec![static_asset_component("site", "app/web")]);
276 let live = collect_live_derive_hashes(tmp.path(), [&svc]).unwrap();
277
278 assert!(live.fetch.contains(FETCH_HASH), "fetch hash must be live");
279 assert!(
280 live.transform.is_empty(),
281 "fetch-only asset has no transform cache entry"
282 );
283 }
284
285 #[test]
286 fn collect_live_hashes_fetch_and_transform() {
287 let tmp = TempDir::new().unwrap();
288 let workload_dir = tmp.path().join("app/web");
289 write_workload(
290 &workload_dir,
291 &workload_with_fetch_and_transform(FETCH_HASH, OUTPUT_HASH),
292 );
293
294 let svc = svc("my-svc", vec![static_asset_component("site", "app/web")]);
295 let live = collect_live_derive_hashes(tmp.path(), [&svc]).unwrap();
296
297 assert!(live.fetch.contains(FETCH_HASH));
298 assert!(
299 live.transform.contains(OUTPUT_HASH),
300 "transform cache entry keyed by output blake3"
301 );
302 }
303
304 #[test]
305 fn collect_live_hashes_missing_workload_skips_silently() {
306 let tmp = TempDir::new().unwrap();
307 let svc = svc("my-svc", vec![static_asset_component("site", "app/web")]);
309 let live = collect_live_derive_hashes(tmp.path(), [&svc]).unwrap();
310 assert!(live.fetch.is_empty());
311 assert!(live.transform.is_empty());
312 }
313
314 #[test]
315 fn collect_live_hashes_ignores_non_static_asset_components() {
316 let tmp = TempDir::new().unwrap();
317 let svc = svc(
318 "my-svc",
319 vec![ServiceComponent {
320 mount: None,
321 id: "api".into(),
322 kind: "container".into(),
323 path: "app/api".into(),
324 role: String::new(),
325 publishes: None,
326 wave: 0,
327 git: None,
328 }],
329 );
330 let live = collect_live_derive_hashes(tmp.path(), [&svc]).unwrap();
331 assert!(live.fetch.is_empty());
332 }
333
334 #[test]
335 fn candidates_exclude_live_hashes() {
336 let tmp = TempDir::new().unwrap();
337 let fetch_dir = tmp.path().join("fetch");
338 std::fs::create_dir_all(&fetch_dir).unwrap();
339
340 std::fs::write(fetch_dir.join(format!("{FETCH_HASH}.bin")), b"live").unwrap();
342 let orphan = "cccccccccccccccccccccccccccccccccccccccccccccccccccccccccccccccc";
344 std::fs::write(fetch_dir.join(format!("{orphan}.bin")), b"orphan").unwrap();
345 std::fs::write(fetch_dir.join(format!("{orphan}.partial")), b"partial").unwrap();
347
348 let mut live = DeriveCacheLiveHashes::default();
349 live.fetch.insert(FETCH_HASH.to_string());
350
351 let candidates = compute_derive_cache_candidates(tmp.path(), &live).unwrap();
352 assert_eq!(candidates.len(), 1, "exactly one orphan");
353 assert_eq!(candidates[0].hash, orphan);
354 assert_eq!(candidates[0].cache_kind, "fetch");
355 }
356
357 #[test]
358 fn execute_prune_deletes_candidates() {
359 let tmp = TempDir::new().unwrap();
360 let fetch_dir = tmp.path().join("fetch");
361 std::fs::create_dir_all(&fetch_dir).unwrap();
362
363 let file = fetch_dir.join("dead.bin");
364 std::fs::write(&file, b"orphan bytes").unwrap();
365 assert!(file.exists());
366
367 let candidate = DerivePruneCandidate {
368 path: file.clone(),
369 cache_kind: "fetch",
370 hash: "dead".into(),
371 size: 12,
372 };
373 let (deleted, bytes, errors) = execute_derive_cache_prune(&[candidate]);
374 assert_eq!(deleted, 1);
375 assert_eq!(bytes, 12);
376 assert_eq!(errors, 0);
377 assert!(!file.exists(), "file must be deleted");
378 }
379}