microsandbox_image/cache/
lease.rs1use 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
10impl GlobalCache {
15 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 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 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 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 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 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#[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 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}