1use std::path::{Path, PathBuf};
19use std::time::{SystemTime, UNIX_EPOCH};
20
21use futures_lite::StreamExt as _;
22use serde::{Deserialize, Serialize};
23use stow_types::error::Context;
24use stow_types::index::{ArtifactIndex, STOW_INDEX_MEDIA_TYPE, index_tag};
25use stow_types::registry::GHCR_BASE;
26
27use crate::config::StowConfig;
28use crate::verify;
29
30const POINTER_FILE: &str = "current.json";
32
33#[derive(Debug, Clone, Serialize, Deserialize)]
37struct SlicePointer {
38 manifest_digest: String,
39 fetched_at: u64,
40 row_count: u64,
41}
42
43#[derive(Debug, Clone)]
46pub struct IndexSlice {
47 pub manifest_digest: String,
49 pub index: ArtifactIndex,
52}
53
54#[derive(Debug, Clone)]
56pub struct CachedSliceStatus {
57 pub target: String,
59 pub rustc_version: String,
61 pub manifest_digest: String,
63 pub fetched_at: u64,
65 pub row_count: u64,
67}
68
69pub async fn ensure_slice(
79 config: &StowConfig,
80 target: &str,
81 rustc_version: &str,
82) -> stow_types::error::Result<IndexSlice> {
83 fetch_slice(config, target, rustc_version, false).await
84}
85
86pub async fn cached_slice(
98 config: &StowConfig,
99 target: &str,
100 rustc_version: &str,
101) -> stow_types::error::Result<Option<IndexSlice>> {
102 let dir = slice_dir(config, target, rustc_version);
103 let Some(pointer) = read_pointer(&dir).await else {
104 return Ok(None);
105 };
106 load_cached(&dir, &pointer).await.map(Some)
107}
108
109pub async fn refresh_slice(
117 config: &StowConfig,
118 target: &str,
119 rustc_version: &str,
120) -> stow_types::error::Result<IndexSlice> {
121 fetch_slice(config, target, rustc_version, true).await
122}
123
124pub async fn cached_slices(
130 config: &StowConfig,
131) -> stow_types::error::Result<Vec<CachedSliceStatus>> {
132 let root = config.cache_dir.join("index");
133 let mut slices = Vec::new();
134 let Ok(mut targets) = async_fs::read_dir(&root).await else {
135 return Ok(slices);
136 };
137 while let Some(target_entry) = targets
138 .next()
139 .await
140 .transpose()
141 .wrap_err_with(|| format!("traverse index cache {}", root.display()))?
142 {
143 if !target_entry
144 .file_type()
145 .await
146 .is_ok_and(|kind| kind.is_dir())
147 {
148 continue;
149 }
150 let target = target_entry.file_name().to_string_lossy().into_owned();
151 let mut versions = async_fs::read_dir(target_entry.path())
152 .await
153 .wrap_err_with(|| format!("read index dir {}", target_entry.path().display()))?;
154 while let Some(version_entry) = versions
155 .next()
156 .await
157 .transpose()
158 .wrap_err_with(|| format!("traverse index dir {}", target_entry.path().display()))?
159 {
160 let dir = version_entry.path();
161 let Some(pointer) = read_pointer(&dir).await else {
162 continue;
163 };
164 slices.push(CachedSliceStatus {
165 target: target.clone(),
166 rustc_version: version_entry.file_name().to_string_lossy().into_owned(),
167 manifest_digest: pointer.manifest_digest,
168 fetched_at: pointer.fetched_at,
169 row_count: pointer.row_count,
170 });
171 }
172 }
173 slices.sort_by(|a, b| {
174 a.target
175 .cmp(&b.target)
176 .then_with(|| a.rustc_version.cmp(&b.rustc_version))
177 });
178 Ok(slices)
179}
180
181async fn fetch_slice(
182 config: &StowConfig,
183 target: &str,
184 rustc_version: &str,
185 force: bool,
186) -> stow_types::error::Result<IndexSlice> {
187 let dir = slice_dir(config, target, rustc_version);
188 let pointer = read_pointer(&dir).await;
189 if !force
190 && let Some(pointer) = &pointer
191 && now_secs().saturating_sub(pointer.fetched_at) < config.index_refresh_interval.as_secs()
192 {
193 return load_cached(&dir, pointer).await;
194 }
195
196 let tag = index_tag(target, rustc_version);
197 let base = stow_oci::RegistryBase::parse(&config.registry_base_url)?;
198 let session = base.session();
199 let reference = base.reference(&tag)?;
200
201 match session.fetch_manifest_digest(&reference).await {
202 Ok(remote_digest) => {
203 if let Some(pointer) = &pointer
204 && pointer.manifest_digest == remote_digest
205 {
206 let pointer = SlicePointer {
207 row_count: pointer.row_count,
208 manifest_digest: remote_digest,
209 fetched_at: now_secs(),
210 };
211 write_pointer(&dir, &pointer).await?;
212 return load_cached(&dir, &pointer).await;
213 }
214 let (blob, manifest_digest, index) =
215 download_verified_slice(config, &session, &reference, &tag, target, rustc_version)
216 .await?;
217 store_slice(&dir, &manifest_digest, &blob).await?;
218 let pointer = SlicePointer {
219 row_count: index.rows.len() as u64,
220 manifest_digest: manifest_digest.clone(),
221 fetched_at: now_secs(),
222 };
223 write_pointer(&dir, &pointer).await?;
224 Ok(IndexSlice {
225 manifest_digest,
226 index,
227 })
228 }
229 Err(error) => {
230 if let Some(pointer) = &pointer {
231 tracing::info!(
232 error = %error,
233 tag = %tag,
234 "index refresh failed; serving cached slice"
235 );
236 return load_cached(&dir, pointer).await;
237 }
238 Err(stow_types::stow_error!("fetch index slice {tag}: {error}"))
239 }
240 }
241}
242
243async fn download_verified_slice(
247 config: &StowConfig,
248 session: &stow_oci::RegistrySession,
249 reference: &oci_client::Reference,
250 tag: &str,
251 target: &str,
252 rustc_version: &str,
253) -> stow_types::error::Result<(Vec<u8>, String, ArtifactIndex)> {
254 let (manifest_digest, manifest) = stow_oci::pull_tagged_manifest(session, reference).await?;
255 let [layer] = manifest.layers.as_slice() else {
256 return Err(stow_types::stow_error!(
257 "index manifest {reference} carries {} layers, expected exactly one",
258 manifest.layers.len()
259 ));
260 };
261 if layer.media_type != STOW_INDEX_MEDIA_TYPE {
262 return Err(stow_types::stow_error!(
263 "index manifest {reference} layer is {}, expected {STOW_INDEX_MEDIA_TYPE}",
264 layer.media_type
265 ));
266 }
267 let blob = stow_oci::pull_blob_verified(session, layer).await?;
268 let materials =
269 stow_oci::pull_signature_materials(session, reference, &manifest_digest).await?;
270 let identity_reference = format!("{GHCR_BASE}:{tag}");
273 verify::verify_index_signature(config, &identity_reference, &manifest_digest, &materials)
274 .await?;
275 let index = stow_types::index::decode(&blob).wrap_err("decode index slice")?;
276 if index.header.target.as_str() != target
277 || index.header.rustc_version.as_str() != rustc_version
278 {
279 return Err(stow_types::stow_error!(
280 "index slice {tag} was published for {}@{}",
281 index.header.target.as_str(),
282 index.header.rustc_version.as_str()
283 ));
284 }
285 Ok((blob, manifest_digest, index))
286}
287
288async fn store_slice(
292 dir: &Path,
293 manifest_digest: &str,
294 blob: &[u8],
295) -> stow_types::error::Result<()> {
296 async_fs::create_dir_all(dir)
297 .await
298 .wrap_err_with(|| format!("create index dir {}", dir.display()))?;
299 let keep = blob_name(manifest_digest);
300 let tmp = dir.join(format!(".tmp-{}", std::process::id()));
301 async_fs::write(&tmp, blob)
302 .await
303 .wrap_err_with(|| format!("write index blob {}", tmp.display()))?;
304 async_fs::rename(&tmp, dir.join(&keep))
305 .await
306 .wrap_err_with(|| format!("commit index blob into {}", dir.display()))?;
307 let mut entries = async_fs::read_dir(dir)
308 .await
309 .wrap_err_with(|| format!("read index dir {}", dir.display()))?;
310 while let Some(entry) = entries
311 .next()
312 .await
313 .transpose()
314 .wrap_err_with(|| format!("traverse index dir {}", dir.display()))?
315 {
316 let name = entry.file_name();
317 if name != POINTER_FILE && name != keep.as_str() {
318 async_fs::remove_file(entry.path())
319 .await
320 .wrap_err_with(|| format!("evict stale index blob {}", entry.path().display()))?;
321 }
322 }
323 Ok(())
324}
325
326async fn write_pointer(dir: &Path, pointer: &SlicePointer) -> stow_types::error::Result<()> {
327 async_fs::create_dir_all(dir)
328 .await
329 .wrap_err_with(|| format!("create index dir {}", dir.display()))?;
330 let tmp = dir.join(format!(".{POINTER_FILE}.tmp-{}", std::process::id()));
331 async_fs::write(&tmp, serde_json::to_vec(pointer)?)
332 .await
333 .wrap_err_with(|| format!("write index pointer {}", tmp.display()))?;
334 async_fs::rename(&tmp, dir.join(POINTER_FILE))
335 .await
336 .wrap_err_with(|| format!("commit index pointer in {}", dir.display()))?;
337 Ok(())
338}
339
340async fn read_pointer(dir: &Path) -> Option<SlicePointer> {
341 let bytes = async_fs::read(dir.join(POINTER_FILE)).await.ok()?;
342 serde_json::from_slice(&bytes).ok()
343}
344
345async fn load_cached(dir: &Path, pointer: &SlicePointer) -> stow_types::error::Result<IndexSlice> {
346 let path = dir.join(blob_name(&pointer.manifest_digest));
347 let bytes = async_fs::read(&path)
348 .await
349 .wrap_err_with(|| format!("read cached index {}", path.display()))?;
350 let index = stow_types::index::decode(&bytes).wrap_err("decode cached index slice")?;
351 Ok(IndexSlice {
352 manifest_digest: pointer.manifest_digest.clone(),
353 index,
354 })
355}
356
357fn blob_name(manifest_digest: &str) -> String {
360 manifest_digest.replace(':', "_")
361}
362
363fn slice_dir(config: &StowConfig, target: &str, rustc_version: &str) -> PathBuf {
364 config
365 .cache_dir
366 .join("index")
367 .join(target)
368 .join(rustc_version)
369}
370
371fn now_secs() -> u64 {
372 SystemTime::now()
373 .duration_since(UNIX_EPOCH)
374 .map_or(0, |since| since.as_secs())
375}
376
377#[cfg(test)]
378mod tests {
379 use semver::Version;
380 use stow_types::artifact::{ArtifactKind, RustCrateType};
381 use stow_types::identity::{
382 CMetadata, CrateName, CrateVersion, DependencyCMetadataJson, FeaturesJson, TargetTriple,
383 WireRustcVersion,
384 };
385 use stow_types::index::{ARTIFACT_INDEX_FORMAT_VERSION, ArtifactIndexHeader, ArtifactIndexRow};
386 use stow_types::platform::{PanicStrategy, Profile, StripLevel};
387
388 use super::*;
389 use crate::config::VerifyMode;
390
391 const TARGET: &str = "x86_64-unknown-linux-gnu";
392 const RUSTC: &str = "1.91.1";
393
394 fn test_config(cache_dir: &Path) -> StowConfig {
395 StowConfig {
396 edge_url: "http://127.0.0.1:8787".to_owned(),
397 registry_base_url: "http://127.0.0.1:8787/v2/water-rs/stow-cache".to_owned(),
398 cache_dir: cache_dir.to_path_buf(),
399 request_timeout: std::time::Duration::from_secs(15),
400 negative_cache_ttl: std::time::Duration::from_mins(5),
401 circuit_reset_after: std::time::Duration::from_mins(1),
402 circuit_trip_threshold: 5,
403 artifact_cache_max_bytes: 1024,
404 index_refresh_interval: std::time::Duration::from_mins(10),
405 verify_mode: VerifyMode::GithubCi,
406 state_db_pool: StowConfig::default_state_db_pool(),
407 trust_material: std::sync::Arc::default(),
408 }
409 }
410
411 fn test_row(c_metadata: &str) -> ArtifactIndexRow {
412 ArtifactIndexRow {
413 crate_name: CrateName::parse("serde").expect("crate name"),
414 version: CrateVersion::new(Version::new(1, 0, 219)),
415 features_json: FeaturesJson::canonicalize(vec!["default".to_owned()])
416 .expect("features"),
417 dependency_c_metadata_json: DependencyCMetadataJson::default(),
418 c_metadata: CMetadata::parse(c_metadata).expect("c_metadata"),
419 compile_key: format!("{c_metadata}{c_metadata}"),
420 bundle_digest: format!("sha256:{c_metadata:0>64}"),
421 bundle_size: 1234,
422 artifact_kind: ArtifactKind::Rlib,
423 crate_types: vec![RustCrateType::Rlib],
424 profile: Profile {
425 opt_level: "3".to_owned(),
426 debuginfo: 0,
427 debug_assertions: false,
428 overflow_checks: false,
429 panic: PanicStrategy::Unwind,
430 strip: StripLevel::None,
431 },
432 emit: vec!["link".to_owned(), "metadata".to_owned()],
433 min_glibc: None,
434 }
435 }
436
437 fn test_index(rows: Vec<ArtifactIndexRow>) -> ArtifactIndex {
438 ArtifactIndex {
439 header: ArtifactIndexHeader {
440 format_version: ARTIFACT_INDEX_FORMAT_VERSION,
441 target: TargetTriple::parse(TARGET).expect("target"),
442 rustc_version: WireRustcVersion::parse(RUSTC).expect("rustc"),
443 generated_at: "2026-09-24T12:00:00Z".to_owned(),
444 row_count: rows.len() as u64,
445 },
446 rows,
447 }
448 }
449
450 #[test]
451 fn blob_name_sanitizes_digest_colon() {
452 assert_eq!(blob_name("sha256:ab12"), "sha256_ab12");
453 }
454
455 #[tokio::test]
456 async fn store_pointer_then_load_cached_round_trips() {
457 let tempdir = tempfile::tempdir().expect("tempdir");
458 let config = test_config(tempdir.path());
459 let dir = slice_dir(&config, TARGET, RUSTC);
460 let index = test_index(vec![test_row("aaaa"), test_row("bbbb")]);
461 let blob = stow_types::index::encode(&index).expect("encode");
462 let digest = "sha256:deadbeef".to_owned();
463
464 store_slice(&dir, &digest, &blob).await.expect("store");
465 write_pointer(
466 &dir,
467 &SlicePointer {
468 manifest_digest: digest.clone(),
469 fetched_at: 1234,
470 row_count: 2,
471 },
472 )
473 .await
474 .expect("write pointer");
475
476 let pointer = read_pointer(&dir).await.expect("pointer");
477 assert_eq!(pointer.manifest_digest, digest);
478 let loaded = load_cached(&dir, &pointer).await.expect("load cached");
479 assert_eq!(loaded.manifest_digest, digest);
480 assert_eq!(loaded.index, index);
481 }
482
483 #[tokio::test]
484 async fn store_slice_evicts_stale_blobs_and_pointer_stays() {
485 let tempdir = tempfile::tempdir().expect("tempdir");
486 let config = test_config(tempdir.path());
487 let dir = slice_dir(&config, TARGET, RUSTC);
488 let blob = stow_types::index::encode(&test_index(vec![test_row("aaaa")])).expect("encode");
489
490 store_slice(&dir, "sha256:aaaa", &blob)
491 .await
492 .expect("first store");
493 write_pointer(
494 &dir,
495 &SlicePointer {
496 manifest_digest: "sha256:aaaa".to_owned(),
497 fetched_at: 1,
498 row_count: 1,
499 },
500 )
501 .await
502 .expect("pointer");
503 store_slice(&dir, "sha256:bbbb", &blob)
504 .await
505 .expect("second store");
506
507 let mut names: Vec<String> = std::fs::read_dir(&dir)
508 .expect("read dir")
509 .map(|entry| {
510 entry
511 .expect("entry")
512 .file_name()
513 .to_string_lossy()
514 .into_owned()
515 })
516 .collect();
517 names.sort();
518 assert_eq!(names, vec!["current.json", "sha256_bbbb"]);
519 }
520
521 #[tokio::test]
522 async fn cached_slices_reports_pointer_fields() {
523 let tempdir = tempfile::tempdir().expect("tempdir");
524 let config = test_config(tempdir.path());
525 let dir = slice_dir(&config, TARGET, RUSTC);
526 let index = test_index(vec![test_row("aaaa")]);
527 let blob = stow_types::index::encode(&index).expect("encode");
528 store_slice(&dir, "sha256:aaaa", &blob)
529 .await
530 .expect("store");
531 write_pointer(
532 &dir,
533 &SlicePointer {
534 manifest_digest: "sha256:aaaa".to_owned(),
535 fetched_at: 4242,
536 row_count: 1,
537 },
538 )
539 .await
540 .expect("pointer");
541
542 let slices = cached_slices(&config).await.expect("cached slices");
543 assert_eq!(slices.len(), 1);
544 let slice = &slices[0];
545 assert_eq!(slice.target, TARGET);
546 assert_eq!(slice.rustc_version, RUSTC);
547 assert_eq!(slice.manifest_digest, "sha256:aaaa");
548 assert_eq!(slice.fetched_at, 4242);
549 assert_eq!(slice.row_count, 1);
550 }
551
552 #[tokio::test]
553 async fn cached_slices_empty_without_cache() {
554 let tempdir = tempfile::tempdir().expect("tempdir");
555 let config = test_config(tempdir.path());
556 assert!(
557 cached_slices(&config)
558 .await
559 .expect("cached slices")
560 .is_empty()
561 );
562 }
563}