Skip to main content

mbx_cache_core/
local.rs

1use crate::{CacheDigest, RemoteActionResult, canonical_json};
2use eyre::{Result, bail};
3use std::fs;
4use std::io::Write;
5use std::path::{Path, PathBuf};
6
7/// A validated, content-addressed store on the local filesystem.
8#[derive(Debug, Clone)]
9pub struct LocalCas {
10    root: PathBuf,
11}
12
13/// A local index from action digests to their referenced cache objects.
14#[derive(Debug, Clone)]
15pub struct LocalActionCache {
16    root: PathBuf,
17    cas: LocalCas,
18}
19
20impl LocalCas {
21    /// Create a local content-addressed store beneath `root`.
22    pub fn new(root: impl Into<PathBuf>) -> Self {
23        Self { root: root.into() }
24    }
25
26    /// Return the root shared by this store and its action-result index.
27    pub fn root(&self) -> &Path {
28        &self.root
29    }
30
31    /// Resolve the storage path for a validated digest.
32    pub fn path_for(&self, digest: &CacheDigest) -> Result<PathBuf> {
33        digest.validate()?;
34        Ok(self
35            .root
36            .join("cas/v1")
37            .join(&digest.algorithm)
38            .join(&digest.hash[..2])
39            .join(format!("{}-{}", digest.hash, digest.size)))
40    }
41
42    /// Find and verify a stored object.
43    pub fn find(&self, digest: &CacheDigest) -> Result<Option<PathBuf>> {
44        let path = self.path_for(digest)?;
45        if !path.exists() {
46            return Ok(None);
47        }
48        if !digest.matches_file(&path)? {
49            bail!(
50                "local CAS blob failed digest verification: {}",
51                path.display()
52            );
53        }
54        Ok(Some(path))
55    }
56
57    /// Atomically store bytes after verifying their declared digest.
58    pub fn store_bytes(&self, digest: &CacheDigest, bytes: &[u8]) -> Result<PathBuf> {
59        if !digest.matches_bytes(bytes)? {
60            bail!("bytes do not match the declared CAS digest");
61        }
62        self.store_with(digest, |temporary| {
63            temporary.write_all(bytes)?;
64            Ok(())
65        })
66    }
67
68    /// Atomically store a file after verifying its declared digest.
69    pub fn store_file(&self, digest: &CacheDigest, source: &Path) -> Result<PathBuf> {
70        self.store_file_inner(digest, source, true)
71    }
72
73    /// Store a file whose digest was already verified by this crate.
74    pub(crate) fn store_verified_file(
75        &self,
76        digest: &CacheDigest,
77        source: &Path,
78    ) -> Result<PathBuf> {
79        self.store_file_inner(digest, source, false)
80    }
81
82    /// Move an already-verified temporary file into the CAS when possible.
83    ///
84    /// Remote downloads are staged beneath the cache root, so the usual path
85    /// is an atomic same-filesystem rename. If a caller supplies a path on a
86    /// different filesystem, fall back to the copy-based store operation.
87    pub fn adopt_verified_file(&self, digest: &CacheDigest, source: &Path) -> Result<PathBuf> {
88        if fs::metadata(source)?.len() != digest.size {
89            bail!("staged blob size does not match the declared CAS digest");
90        }
91        let destination = self.path_for(digest)?;
92        match self.find(digest) {
93            Ok(Some(existing)) => return Ok(existing),
94            // Preserve the established repair path for a poisoned destination;
95            // replacing it portably needs the copy-based temporary file flow.
96            Err(_) => return self.store_file_inner(digest, source, false),
97            Ok(None) => {}
98        }
99        fs::create_dir_all(destination.parent().expect("CAS path has a parent"))?;
100        make_owner_writable(source)?;
101        match fs::rename(source, &destination) {
102            Ok(()) => Ok(destination),
103            Err(error) if error.kind() == std::io::ErrorKind::AlreadyExists => self
104                .find(digest)?
105                .ok_or_else(|| eyre::eyre!("concurrent CAS write did not publish a valid blob")),
106            Err(_) => self.store_file_inner(digest, source, false),
107        }
108    }
109
110    fn store_file_inner(
111        &self,
112        digest: &CacheDigest,
113        source: &Path,
114        verify: bool,
115    ) -> Result<PathBuf> {
116        let destination = self.path_for(digest)?;
117        // A blob that fails verification cannot be restored from, and nothing
118        // else repairs it: the read path reports an error rather than a miss,
119        // so without republishing over it the digest stays poisoned until
120        // eviction happens to reclaim it. `LocalActionCache::store` already
121        // recovers this way one layer up.
122        let replace_invalid = match self.find(digest) {
123            Ok(Some(existing)) => return Ok(existing),
124            Ok(None) => false,
125            Err(_) => true,
126        };
127        let parent = destination.parent().expect("CAS path has a parent");
128        fs::create_dir_all(parent)?;
129        let staging = tempfile::tempdir_in(parent)?;
130        let temporary = staging.path().join("blob");
131        reflink_copy::reflink_or_copy(source, &temporary)?;
132        let temporary = tempfile::TempPath::try_from_path(temporary)?;
133        make_owner_writable(&temporary)?;
134        // Not fsynced: every read verifies the digest, so a blob torn by a
135        // crash is detected and treated as absent rather than trusted.
136        if verify && !digest.matches_file(&temporary)? {
137            bail!("staged blob does not match the declared CAS digest");
138        }
139        if fs::metadata(&temporary)?.len() != digest.size {
140            bail!("staged blob size does not match the declared CAS digest");
141        }
142        if replace_invalid {
143            temporary
144                .persist(&destination)
145                .map_err(|error| error.error)?;
146            return Ok(destination);
147        }
148        match temporary.persist_noclobber(&destination) {
149            Ok(()) => Ok(destination),
150            Err(error) if error.error.kind() == std::io::ErrorKind::AlreadyExists => self
151                .find(digest)?
152                .ok_or_else(|| eyre::eyre!("concurrent CAS write did not publish a valid blob")),
153            Err(error) => Err(error.error.into()),
154        }
155    }
156
157    fn store_with(
158        &self,
159        digest: &CacheDigest,
160        write: impl FnOnce(&mut tempfile::NamedTempFile) -> Result<()>,
161    ) -> Result<PathBuf> {
162        let destination = self.path_for(digest)?;
163        // A blob that fails verification cannot be restored from, and nothing
164        // else repairs it: the read path reports an error rather than a miss,
165        // so without republishing over it the digest stays poisoned until
166        // eviction happens to reclaim it. `LocalActionCache::store` already
167        // recovers this way one layer up.
168        let replace_invalid = match self.find(digest) {
169            Ok(Some(existing)) => return Ok(existing),
170            Ok(None) => false,
171            Err(_) => true,
172        };
173        let parent = destination.parent().expect("CAS path has a parent");
174        fs::create_dir_all(parent)?;
175        let mut temporary = tempfile::NamedTempFile::new_in(parent)?;
176        write(&mut temporary)?;
177        temporary.flush()?;
178        if !digest.matches_file(temporary.path())? {
179            bail!("staged blob does not match the declared CAS digest");
180        }
181        if replace_invalid {
182            temporary
183                .persist(&destination)
184                .map_err(|error| error.error)?;
185            return Ok(destination);
186        }
187        match temporary.persist_noclobber(&destination) {
188            Ok(_) => Ok(destination),
189            Err(error) if error.error.kind() == std::io::ErrorKind::AlreadyExists => self
190                .find(digest)?
191                .ok_or_else(|| eyre::eyre!("concurrent CAS write did not publish a valid blob")),
192            Err(error) => Err(error.error.into()),
193        }
194    }
195}
196
197#[cfg(unix)]
198fn make_owner_writable(path: &Path) -> Result<()> {
199    use std::os::unix::fs::PermissionsExt as _;
200    let mut permissions = fs::metadata(path)?.permissions();
201    permissions.set_mode(permissions.mode() | 0o200);
202    fs::set_permissions(path, permissions)?;
203    Ok(())
204}
205
206#[cfg(windows)]
207fn make_owner_writable(path: &Path) -> Result<()> {
208    let mut permissions = fs::metadata(path)?.permissions();
209    permissions.set_readonly(false);
210    fs::set_permissions(path, permissions)?;
211    Ok(())
212}
213
214impl LocalActionCache {
215    /// Create an action-result index beneath `root`.
216    pub fn new(root: impl Into<PathBuf>) -> Self {
217        let root = root.into();
218        Self {
219            cas: LocalCas::new(root.clone()),
220            root,
221        }
222    }
223
224    /// Resolve the storage path for an action digest.
225    pub fn path_for(&self, action: &CacheDigest) -> Result<PathBuf> {
226        action.validate()?;
227        if action.algorithm != "blake3" {
228            bail!("local action keys must use blake3");
229        }
230        Ok(self
231            .root
232            .join("action-results/v1")
233            .join(&action.algorithm)
234            .join(&action.hash[..2])
235            .join(format!("{}-{}.json", action.hash, action.size)))
236    }
237
238    /// Find and strictly validate a canonical action result.
239    pub fn find(&self, action: &CacheDigest) -> Result<Option<RemoteActionResult>> {
240        let path = self.path_for(action)?;
241        if !path.exists() {
242            return Ok(None);
243        }
244        let bytes = fs::read(&path)?;
245        let result: RemoteActionResult = serde_json::from_slice(&bytes)?;
246        if result.version != 1 || result.action != *action || canonical_json(&result)? != bytes {
247            bail!("local action result is invalid: {}", path.display());
248        }
249        Ok(Some(result))
250    }
251
252    /// Atomically publish an action result after validating all referenced objects.
253    pub fn store(&self, result: &RemoteActionResult) -> Result<PathBuf> {
254        if result.version != 1 {
255            bail!("unsupported local action result version");
256        }
257        for digest in [
258            Some(&result.action),
259            result.metadata.as_ref(),
260            result.output_root.as_ref(),
261        ]
262        .into_iter()
263        .flatten()
264        {
265            if self.cas.find(digest)?.is_none() {
266                bail!("cannot publish an action result with a missing blob");
267            }
268        }
269        let destination = self.path_for(&result.action)?;
270        let replace_invalid = match self.find(&result.action) {
271            Ok(Some(existing)) => {
272                if existing == *result {
273                    return Ok(destination);
274                }
275                bail!("local action key already has a different result");
276            }
277            Ok(None) => false,
278            Err(_) => true,
279        };
280        let parent = destination
281            .parent()
282            .expect("action-result path has a parent");
283        fs::create_dir_all(parent)?;
284        let mut temporary = tempfile::NamedTempFile::new_in(parent)?;
285        temporary.write_all(&canonical_json(result)?)?;
286        temporary.flush()?;
287        if replace_invalid {
288            temporary
289                .persist(&destination)
290                .map_err(|error| error.error)?;
291            return Ok(destination);
292        }
293        match temporary.persist_noclobber(&destination) {
294            Ok(_) => Ok(destination),
295            Err(error) if error.error.kind() == std::io::ErrorKind::AlreadyExists => self
296                .find(&result.action)?
297                .filter(|existing| existing == result)
298                .map(|_| destination)
299                .ok_or_else(|| eyre::eyre!("concurrent action write was invalid or conflicting")),
300            Err(error) => Err(error.error.into()),
301        }
302    }
303}
304
305#[cfg(test)]
306mod tests {
307    use super::*;
308
309    #[test]
310    fn stores_and_validates_blobs_atomically() {
311        let directory = tempfile::tempdir().unwrap();
312        let cas = LocalCas::new(directory.path());
313        let digest = CacheDigest::blake3(b"cached object");
314
315        let path = cas.store_bytes(&digest, b"cached object").unwrap();
316        assert_eq!(cas.find(&digest).unwrap(), Some(path.clone()));
317        assert_eq!(fs::read(&path).unwrap(), b"cached object");
318        assert_eq!(cas.store_bytes(&digest, b"cached object").unwrap(), path);
319        assert!(cas.store_bytes(&digest, b"other object").is_err());
320    }
321
322    #[test]
323    fn stored_files_are_independent_from_the_source() {
324        let directory = tempfile::tempdir().unwrap();
325        let cas = LocalCas::new(directory.path().join("cache"));
326        let source = directory.path().join("source");
327        fs::write(&source, b"cached object").unwrap();
328        let digest = CacheDigest::blake3(b"cached object");
329
330        let stored = cas.store_file(&digest, &source).unwrap();
331        fs::write(source, b"other object!").unwrap();
332
333        assert_eq!(fs::read(stored).unwrap(), b"cached object");
334        assert!(cas.find(&digest).unwrap().is_some());
335    }
336
337    #[test]
338    fn adopts_verified_files_without_leaving_the_staging_copy() {
339        let directory = tempfile::tempdir().unwrap();
340        let cas = LocalCas::new(directory.path().join("cache"));
341        let staging = directory.path().join("remote");
342        fs::create_dir(&staging).unwrap();
343        let source = staging.join("blob");
344        fs::write(&source, b"cached object").unwrap();
345        let digest = CacheDigest::blake3(b"cached object");
346
347        let stored = cas.adopt_verified_file(&digest, &source).unwrap();
348
349        assert!(!source.exists());
350        assert_eq!(fs::read(&stored).unwrap(), b"cached object");
351        assert_eq!(cas.find(&digest).unwrap(), Some(stored));
352    }
353
354    #[test]
355    fn rejects_files_with_the_wrong_digest() {
356        let directory = tempfile::tempdir().unwrap();
357        let cas = LocalCas::new(directory.path().join("cache"));
358        let source = directory.path().join("source");
359        fs::write(&source, b"other object").unwrap();
360        let digest = CacheDigest::blake3(b"cached object");
361
362        assert!(cas.store_file(&digest, &source).is_err());
363        assert!(!cas.path_for(&digest).unwrap().exists());
364    }
365
366    #[test]
367    fn stores_read_only_source_files() {
368        let directory = tempfile::tempdir().unwrap();
369        let cas = LocalCas::new(directory.path().join("cache"));
370        let source = directory.path().join("source");
371        fs::write(&source, b"cached object").unwrap();
372        let mut permissions = fs::metadata(&source).unwrap().permissions();
373        permissions.set_readonly(true);
374        fs::set_permissions(&source, permissions).unwrap();
375        let digest = CacheDigest::blake3(b"cached object");
376
377        let stored = cas.store_file(&digest, &source).unwrap();
378
379        assert_eq!(fs::read(stored).unwrap(), b"cached object");
380        assert!(fs::metadata(&source).unwrap().permissions().readonly());
381        make_owner_writable(&source).unwrap();
382    }
383
384    #[test]
385    fn rejects_corrupt_existing_blobs() {
386        let directory = tempfile::tempdir().unwrap();
387        let cas = LocalCas::new(directory.path());
388        let digest = CacheDigest::blake3(b"cached object");
389        let path = cas.store_bytes(&digest, b"cached object").unwrap();
390        fs::write(path, b"corrupt").unwrap();
391
392        assert!(cas.find(&digest).is_err());
393    }
394
395    #[test]
396    fn republishes_over_a_corrupt_blob() {
397        let directory = tempfile::tempdir().unwrap();
398        let cas = LocalCas::new(directory.path());
399        let digest = CacheDigest::blake3(b"cached object");
400        let path = cas.store_bytes(&digest, b"cached object").unwrap();
401        fs::write(&path, b"corrupt").unwrap();
402
403        assert_eq!(cas.store_bytes(&digest, b"cached object").unwrap(), path);
404        assert_eq!(fs::read(&path).unwrap(), b"cached object");
405        assert_eq!(cas.find(&digest).unwrap(), Some(path));
406    }
407
408    #[test]
409    fn republishes_a_file_over_a_corrupt_blob() {
410        let directory = tempfile::tempdir().unwrap();
411        let cas = LocalCas::new(directory.path().join("cache"));
412        let source = directory.path().join("source");
413        fs::write(&source, b"cached object").unwrap();
414        let digest = CacheDigest::blake3(b"cached object");
415        let path = cas.store_file(&digest, &source).unwrap();
416        fs::write(&path, b"corrupt").unwrap();
417
418        assert_eq!(cas.store_file(&digest, &source).unwrap(), path);
419        assert_eq!(fs::read(&path).unwrap(), b"cached object");
420        assert_eq!(cas.find(&digest).unwrap(), Some(path));
421    }
422
423    #[test]
424    fn publishes_action_results_after_referenced_blobs() {
425        let directory = tempfile::tempdir().unwrap();
426        let cas = LocalCas::new(directory.path());
427        let actions = LocalActionCache::new(directory.path());
428        let action = CacheDigest::blake3(b"action");
429        let metadata = CacheDigest::blake3(b"metadata");
430        let output_root = CacheDigest::blake3(b"directory");
431        let result = RemoteActionResult {
432            action: action.clone(),
433            metadata: Some(metadata.clone()),
434            output_root: Some(output_root.clone()),
435            version: 1,
436        };
437
438        assert!(actions.store(&result).is_err());
439        cas.store_bytes(&action, b"action").unwrap();
440        cas.store_bytes(&metadata, b"metadata").unwrap();
441        cas.store_bytes(&output_root, b"directory").unwrap();
442        actions.store(&result).unwrap();
443        assert_eq!(actions.find(&action).unwrap(), Some(result));
444    }
445
446    #[test]
447    fn atomically_replaces_a_corrupt_action_result() {
448        let directory = tempfile::tempdir().unwrap();
449        let cas = LocalCas::new(directory.path());
450        let actions = LocalActionCache::new(directory.path());
451        let action = CacheDigest::blake3(b"action");
452        let result = RemoteActionResult {
453            action: action.clone(),
454            metadata: None,
455            output_root: None,
456            version: 1,
457        };
458        cas.store_bytes(&action, b"action").unwrap();
459        let path = actions.path_for(&action).unwrap();
460        fs::create_dir_all(path.parent().unwrap()).unwrap();
461        fs::write(&path, b"truncated").unwrap();
462
463        assert!(actions.find(&action).is_err());
464        assert_eq!(actions.store(&result).unwrap(), path);
465        assert_eq!(actions.find(&action).unwrap(), Some(result));
466    }
467}