1use crate::{CacheDigest, RemoteActionResult, canonical_json};
2use eyre::{Result, bail};
3use std::fs;
4use std::io::Write;
5use std::path::{Path, PathBuf};
6
7#[derive(Debug, Clone)]
9pub struct LocalCas {
10 root: PathBuf,
11}
12
13#[derive(Debug, Clone)]
15pub struct LocalActionCache {
16 root: PathBuf,
17 cas: LocalCas,
18}
19
20impl LocalCas {
21 pub fn new(root: impl Into<PathBuf>) -> Self {
23 Self { root: root.into() }
24 }
25
26 pub fn root(&self) -> &Path {
28 &self.root
29 }
30
31 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 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 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 pub fn store_file(&self, digest: &CacheDigest, source: &Path) -> Result<PathBuf> {
70 self.store_file_inner(digest, source, true)
71 }
72
73 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 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 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 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 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 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 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 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 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 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}