Skip to main content

microsandbox_image/cache/
lease.rs

1//! Operation-scoped protection spanning file access and catalog publication.
2
3use std::collections::BTreeMap;
4use std::path::PathBuf;
5use std::sync::{Arc, Mutex};
6
7use crate::storage_lease::StorageLease;
8use crate::{CachedImageMetadata, Digest, GlobalCache, ImageResult};
9
10//--------------------------------------------------------------------------------------------------
11// Methods
12//--------------------------------------------------------------------------------------------------
13
14impl GlobalCache {
15    /// Start an operation scope. Its clones retain admitted entries until all clones are dropped.
16    /// Keep this scope through publication of durable catalog ownership, including sandbox rows.
17    pub fn operation(&self) -> Self {
18        let mut cache = self.clone();
19        cache.operation = Some(Arc::new(Mutex::new(BTreeMap::new())));
20        cache
21    }
22
23    /// Retain an existing scope, or start one for a standalone image operation.
24    pub(crate) fn operation_or_new(&self) -> Self {
25        if self.operation.is_some() {
26            self.clone()
27        } else {
28            self.operation()
29        }
30    }
31
32    /// Protect entries in stable path order. The returned leases cover this call's reader;
33    /// an operation scope additionally retains them through its caller's publication.
34    pub fn lease_paths(&self, mut paths: Vec<PathBuf>) -> ImageResult<Vec<StorageLease>> {
35        paths = paths
36            .iter()
37            .map(|path| StorageLease::key(path))
38            .collect::<std::io::Result<_>>()?;
39        paths.sort();
40        paths.dedup();
41        if let Some(operation) = &self.operation {
42            let mut retained = operation
43                .lock()
44                .map_err(|_| std::io::Error::other("cache lease scope poisoned"))?;
45            // Reopening a flock creates another open file description. Share the existing
46            // guard instead, including across calls and aliases of the cache directory.
47            return paths
48                .into_iter()
49                .map(|path| {
50                    if let Some(lease) = retained.get(&path) {
51                        return Ok(lease.clone());
52                    }
53                    let lease = StorageLease::shared(&path)?;
54                    retained.insert(path, lease.clone());
55                    Ok(lease)
56                })
57                .collect();
58        }
59        paths
60            .iter()
61            .map(|path| StorageLease::shared(path).map_err(Into::into))
62            .collect()
63    }
64
65    /// Async counterpart of [`Self::lease_paths`].
66    pub async fn lease_paths_async(&self, paths: Vec<PathBuf>) -> ImageResult<Vec<StorageLease>> {
67        let cache = self.clone();
68        tokio::task::spawn_blocking(move || cache.lease_paths(paths))
69            .await
70            .map_err(|error| std::io::Error::other(error.to_string()))?
71    }
72
73    /// All independently reclaimable entries belonging to this immutable manifest.
74    pub fn metadata_paths(&self, metadata: &CachedImageMetadata) -> ImageResult<Vec<PathBuf>> {
75        let digest: Digest = metadata.manifest_digest.parse()?;
76        let mut paths = vec![self.fsmeta_erofs_path(&digest), self.vmdk_path(&digest)];
77        for layer in &metadata.layers {
78            paths.push(self.layer_erofs_path(&layer.diff_id.parse()?));
79        }
80        Ok(paths)
81    }
82}
83
84//--------------------------------------------------------------------------------------------------
85// Tests
86//--------------------------------------------------------------------------------------------------
87
88#[cfg(test)]
89mod tests {
90    use super::*;
91    use crate::{CachedLayerMetadata, ImageConfig, Reference};
92
93    fn metadata(id: char) -> CachedImageMetadata {
94        CachedImageMetadata {
95            manifest_digest: format!("sha256:{}", id.to_string().repeat(64)),
96            config_digest: format!("sha256:{}", id.to_string().repeat(64)),
97            raw_manifest_json: "{}".into(),
98            raw_config_json: "{}".into(),
99            config: ImageConfig::default(),
100            layers: vec![CachedLayerMetadata {
101                digest: format!("sha256:{}", id.to_string().repeat(64)),
102                diff_id: format!("sha256:{}", id.to_string().repeat(64)),
103                media_type: None,
104                size_bytes: Some(1),
105            }],
106        }
107    }
108
109    #[test]
110    fn scope_reuses_pins_across_calls() {
111        let home = tempfile::tempdir().unwrap();
112        let cache = GlobalCache::new(home.path()).unwrap().operation();
113        let path = cache.layer_erofs_path(&metadata('a').manifest_digest.parse().unwrap());
114        for _ in 0..500 {
115            cache.lease_paths(vec![path.clone()]).unwrap();
116        }
117        assert_eq!(cache.operation.as_ref().unwrap().lock().unwrap().len(), 1);
118        assert!(StorageLease::try_exclusive(&path).unwrap().is_none());
119        drop(cache);
120        assert!(StorageLease::try_exclusive(&path).unwrap().is_some());
121    }
122
123    #[test]
124    fn retag_keeps_the_readers_original_dependencies_admitted() {
125        let home = tempfile::tempdir().unwrap();
126        let cache = GlobalCache::new(home.path()).unwrap();
127        let reference: Reference = "example.com/test:latest".parse().unwrap();
128        let original = metadata('a');
129        cache.write_image_metadata(&reference, &original).unwrap();
130        let reader = cache.operation();
131        assert_eq!(
132            reader
133                .read_image_metadata(&reference)
134                .unwrap()
135                .unwrap()
136                .manifest_digest,
137            original.manifest_digest
138        );
139        // Readers no longer need the mutable tag once its exact generation is admitted.
140        cache
141            .write_image_metadata(&reference, &metadata('b'))
142            .unwrap();
143        for path in cache.metadata_paths(&original).unwrap() {
144            assert!(StorageLease::try_exclusive(&path).unwrap().is_none());
145        }
146        drop(reader);
147        for path in cache.metadata_paths(&original).unwrap() {
148            assert!(StorageLease::try_exclusive(&path).unwrap().is_some());
149        }
150    }
151
152    #[tokio::test]
153    async fn digest_scan_retains_only_matching_images() {
154        let home = tempfile::tempdir().unwrap();
155        let cache = GlobalCache::new(home.path()).unwrap();
156        let wanted = metadata('b');
157        let reader = cache.operation();
158        for index in 0..80 {
159            let reference: Reference = format!("example.com/image-{index}:latest").parse().unwrap();
160            cache
161                .write_image_metadata(&reference, &metadata('a'))
162                .unwrap();
163            assert!(
164                reader
165                    .read_image_metadata_matching_async(
166                        cache.image_metadata_path(&reference),
167                        wanted.manifest_digest.clone(),
168                    )
169                    .await
170                    .unwrap()
171                    .is_none()
172            );
173        }
174        assert!(
175            reader
176                .operation
177                .as_ref()
178                .unwrap()
179                .lock()
180                .unwrap()
181                .is_empty()
182        );
183        let reference: Reference = "example.com/wanted:latest".parse().unwrap();
184        cache.write_image_metadata(&reference, &wanted).unwrap();
185        assert!(
186            reader
187                .read_image_metadata_matching_async(
188                    cache.image_metadata_path(&reference),
189                    wanted.manifest_digest.clone(),
190                )
191                .await
192                .unwrap()
193                .is_some()
194        );
195        for path in cache.metadata_paths(&wanted).unwrap() {
196            assert!(StorageLease::try_exclusive(&path).unwrap().is_none());
197        }
198    }
199
200    #[tokio::test]
201    async fn cancelled_waiter_does_not_release_worker_pins() {
202        let home = tempfile::tempdir().unwrap();
203        let cache = GlobalCache::new(home.path()).unwrap().operation();
204        let path = cache.layer_erofs_path(&metadata('a').manifest_digest.parse().unwrap());
205        cache.lease_paths(vec![path.clone()]).unwrap();
206        let (started, ready) = tokio::sync::oneshot::channel();
207        let (release, wait) = std::sync::mpsc::channel();
208        let task = tokio::spawn(async move {
209            tokio::task::spawn_blocking(move || {
210                let _cache = cache;
211                started.send(()).unwrap();
212                wait.recv().unwrap();
213            })
214            .await
215            .unwrap();
216        });
217        ready.await.unwrap();
218        task.abort();
219        assert!(task.await.unwrap_err().is_cancelled());
220        assert!(StorageLease::try_exclusive(&path).unwrap().is_none());
221        release.send(()).unwrap();
222        tokio::time::timeout(std::time::Duration::from_secs(5), async {
223            loop {
224                if StorageLease::try_exclusive(&path).unwrap().is_some() {
225                    break;
226                }
227                tokio::task::yield_now().await;
228            }
229        })
230        .await
231        .unwrap();
232    }
233}