Skip to main content

mj_controller/
image_pull_gate.rs

1//! Coordination between the daemon's background image download and a session
2//! launch that needs the same image.
3//!
4//! The daemon downloads every configured container image shortly after it
5//! starts. A person who creates a session during that download must not start
6//! a second download of the same image: the two would compete for the same
7//! bandwidth and the same layer store. Instead the launch waits for the one
8//! already running, and then finds the image present.
9//!
10//! There is one lock per (host, image) pair. Whoever is allowed to download
11//! holds it; whoever needs the image waits on it.
12//!
13//! The lock also holds a file lock in the instance's data directory, because
14//! `mj doctor --smoke` runs its smoke test in its own
15//! process while the daemon may be downloading the same image. Two engines
16//! pulling one image at once can fail to extract a layer.
17
18use std::collections::BTreeMap;
19use std::fs::File;
20use std::path::PathBuf;
21use std::sync::{Arc, Mutex, MutexGuard};
22use std::time::Duration;
23
24use anyhow::{Context, Result, bail};
25
26use mj_core::targets::{
27    CommandExecutor, ImageHost, ProvisionStage, ProvisionStageGuard, TargetTemplate,
28};
29
30/// How long a waiter sleeps between attempts on the lock.
31///
32/// A download runs for minutes, so polling this slowly costs nothing and keeps
33/// the wait cancellable without a condition variable: both the refresher and a
34/// launch run on blocking threads through the synchronous `CommandExecutor`,
35/// and either can be asked to stop while it waits.
36const POLL_INTERVAL: Duration = Duration::from_millis(250);
37
38/// The lock covering downloads of one image on one host.
39///
40/// Keyed by host label and image reference, which is what identifies a copy of
41/// an image: two targets naming the same image on the same host share one
42/// download, and the same image on two hosts does not.
43pub(crate) fn image_pull_mutex(host: &ImageHost, image: &str) -> Arc<ImagePullLock> {
44    static LOCKS: std::sync::OnceLock<Mutex<BTreeMap<String, std::sync::Weak<ImagePullLock>>>> =
45        std::sync::OnceLock::new();
46    let key = format!("{}|{image}", host.label());
47    let mut locks = LOCKS
48        .get_or_init(Mutex::default)
49        .lock()
50        .unwrap_or_else(std::sync::PoisonError::into_inner);
51    locks.retain(|_, lock| lock.strong_count() > 0);
52    // Unit tests share one process and must not write to the developer's
53    // data directory; the file lock has a test of its own.
54    let file = (!cfg!(test)).then(|| lock_file_path(&key));
55    let slot = locks.entry(key).or_default();
56    if let Some(lock) = slot.upgrade() {
57        return lock;
58    }
59    let lock = Arc::new(ImagePullLock {
60        local: Mutex::new(()),
61        file,
62    });
63    *slot = Arc::downgrade(&lock);
64    lock
65}
66
67/// The file other Mjolnir processes of this instance lock for the same
68/// (host, image) pair. The name is a hash, because an image reference and a
69/// host label can hold characters a file name cannot.
70fn lock_file_path(key: &str) -> PathBuf {
71    use sha2::{Digest, Sha256};
72    let digest = Sha256::digest(key.as_bytes());
73    let name: String = digest
74        .iter()
75        .take(12)
76        .map(|byte| format!("{byte:02x}"))
77        .collect();
78    mj_core::config::data_dir()
79        .join("image-pulls")
80        .join(format!("{name}.lock"))
81}
82
83/// The right to download one image on one host: a lock for the threads of
84/// this process and, when `file` is set, a file lock for other processes.
85pub(crate) struct ImagePullLock {
86    local: Mutex<()>,
87    file: Option<PathBuf>,
88}
89
90impl ImagePullLock {
91    /// A lock that only this process sees.
92    #[cfg(test)]
93    fn in_process() -> Self {
94        Self {
95            local: Mutex::new(()),
96            file: None,
97        }
98    }
99}
100
101/// A held [`ImagePullLock`]. Dropping it releases both locks.
102#[derive(Debug)]
103pub(crate) struct ImagePullGuard<'a> {
104    _local: MutexGuard<'a, ()>,
105    _file: Option<File>,
106}
107
108/// Try the file lock once. `None` means another process holds it.
109fn try_lock_file(path: &std::path::Path) -> Result<Option<File>> {
110    if let Some(parent) = path.parent() {
111        std::fs::create_dir_all(parent).with_context(|| {
112            format!("create image download lock directory {}", parent.display())
113        })?;
114    }
115    let file = std::fs::OpenOptions::new()
116        .create(true)
117        .truncate(false)
118        .write(true)
119        .open(path)
120        .with_context(|| format!("open image download lock {}", path.display()))?;
121    match file.try_lock() {
122        Ok(()) => Ok(Some(file)),
123        Err(std::fs::TryLockError::WouldBlock) => Ok(None),
124        Err(std::fs::TryLockError::Error(error)) => {
125            Err(error).with_context(|| format!("lock image download lock {}", path.display()))
126        }
127    }
128}
129
130/// Take the right to download one image, waiting for whoever holds it.
131///
132/// `on_wait` runs once, the first time the lock is found busy, so a caller can
133/// tell the user it is waiting without saying anything in the common case
134/// where nothing is downloading. `is_cancelled` is polled while waiting so a
135/// quitting daemon or a cancelled Create does not sit here for minutes.
136pub(crate) fn hold_image_pull<'a>(
137    lock: &'a ImagePullLock,
138    is_cancelled: impl Fn() -> bool,
139    on_wait: impl FnOnce(),
140) -> Result<ImagePullGuard<'a>> {
141    let mut on_wait = Some(on_wait);
142    loop {
143        let local = match lock.local.try_lock() {
144            Ok(guard) => Some(guard),
145            // A holder that panicked left no state behind: the lock guards
146            // nothing but the right to run a download.
147            Err(std::sync::TryLockError::Poisoned(poisoned)) => Some(poisoned.into_inner()),
148            Err(std::sync::TryLockError::WouldBlock) => None,
149        };
150        if let Some(local) = local {
151            let file = match &lock.file {
152                Some(path) => try_lock_file(path)?.map(Some),
153                None => Some(None),
154            };
155            if let Some(file) = file {
156                return Ok(ImagePullGuard {
157                    _local: local,
158                    _file: file,
159                });
160            }
161            // Another process is downloading; let this process's other
162            // waiters see the same state while this one sleeps.
163            drop(local);
164        }
165        if let Some(on_wait) = on_wait.take() {
166            on_wait();
167        }
168        if is_cancelled() {
169            bail!("cancelled while waiting for image download");
170        }
171        std::thread::sleep(POLL_INTERVAL);
172    }
173}
174
175/// Run `work` with nothing downloading this target's image underneath it.
176///
177/// For a container target this waits for a background download of the same
178/// image on the same host, and reports the wait as the "Pull image" stage with
179/// a notice saying what it is waiting for. Nothing extra is reported in the
180/// ordinary case where no download is running, and a target that runs no image
181/// just runs `work`.
182pub(crate) fn with_image_ready<T>(
183    target: &TargetTemplate,
184    executor: &impl CommandExecutor,
185    work: impl FnOnce() -> Result<T>,
186) -> Result<T> {
187    let Some((host, container)) = target.image_host() else {
188        return work();
189    };
190    let image = container.image.clone();
191    let lock = image_pull_mutex(&host, &image);
192    // The stage guard is created inside `on_wait`, so it exists only on the
193    // waiting path, and dropped as soon as the wait is over: the stage
194    // describes the wait, not the work that follows it.
195    let mut waiting = None;
196    let guard = hold_image_pull(
197        &lock,
198        || executor.cancellation_requested(),
199        || {
200            waiting = Some(ProvisionStageGuard::new(
201                executor,
202                ProvisionStage::PullingImage,
203            ));
204            executor.notify_notice(&format!("Waiting for image {image} to finish downloading"));
205        },
206    )?;
207    drop(waiting);
208    let result = work();
209    drop(guard);
210    result
211}
212
213#[cfg(test)]
214mod tests {
215    use super::*;
216    use std::sync::atomic::{AtomicBool, AtomicUsize, Ordering};
217
218    /// A Create issued while the daemon is downloading the image waits for
219    /// that download instead of starting a second one, and says so once.
220    #[test]
221    fn a_create_waits_for_the_in_flight_pull_of_its_image() {
222        let lock = Arc::new(ImagePullLock::in_process());
223        let released = Arc::new(AtomicBool::new(false));
224        let waits = Arc::new(AtomicUsize::new(0));
225
226        let downloader = {
227            let lock = lock.clone();
228            let released = released.clone();
229            std::thread::spawn(move || {
230                let guard = lock.local.lock().unwrap();
231                // Long enough that the waiter has to block on it.
232                std::thread::sleep(Duration::from_millis(400));
233                released.store(true, Ordering::Release);
234                drop(guard);
235            })
236        };
237
238        // Make sure the downloader holds the lock before the launch tries.
239        while lock.local.try_lock().is_ok() {
240            std::thread::sleep(Duration::from_millis(10));
241        }
242
243        let guard = hold_image_pull(
244            &lock,
245            || false,
246            || {
247                waits.fetch_add(1, Ordering::Release);
248            },
249        )
250        .expect("the launch takes the lock once the download finishes");
251        assert!(
252            released.load(Ordering::Acquire),
253            "the launch proceeded while the download still held the lock"
254        );
255        assert_eq!(
256            waits.load(Ordering::Acquire),
257            1,
258            "the wait should be announced exactly once"
259        );
260        drop(guard);
261        downloader.join().expect("the download thread finishes");
262    }
263
264    /// A cancelled Create stops waiting instead of sitting behind a
265    /// multi-gigabyte download.
266    #[test]
267    fn a_waiting_create_stops_when_cancelled() {
268        let lock = ImagePullLock::in_process();
269        let held = lock.local.lock().unwrap();
270
271        let error = hold_image_pull(&lock, || true, || {})
272            .expect_err("a cancelled wait must not return a lock it never took");
273        assert!(
274            format!("{error:#}").contains("cancelled while waiting for image download"),
275            "{error:#}"
276        );
277        drop(held);
278    }
279
280    /// `mj doctor --smoke` runs its smoke test in its own process while the daemon may
281    /// be downloading the same image. The file lock makes it wait, as a
282    /// thread of the daemon would: two locks with the same file stand for the
283    /// two processes.
284    #[test]
285    fn a_smoke_test_in_another_process_waits_for_the_daemon_s_download() {
286        let dir = tempfile::tempdir().unwrap();
287        let path = dir.path().join("pulls").join("image.lock");
288        let daemon = ImagePullLock {
289            local: Mutex::new(()),
290            file: Some(path.clone()),
291        };
292        let setup = ImagePullLock {
293            local: Mutex::new(()),
294            file: Some(path),
295        };
296        let holding = AtomicBool::new(false);
297        let released = AtomicBool::new(false);
298        let waits = AtomicUsize::new(0);
299        std::thread::scope(|scope| {
300            scope.spawn(|| {
301                let held = hold_image_pull(&daemon, || false, || {}).unwrap();
302                holding.store(true, Ordering::Release);
303                std::thread::sleep(Duration::from_millis(400));
304                released.store(true, Ordering::Release);
305                drop(held);
306            });
307            while !holding.load(Ordering::Acquire) {
308                std::thread::sleep(Duration::from_millis(10));
309            }
310            let guard = hold_image_pull(
311                &setup,
312                || false,
313                || {
314                    waits.fetch_add(1, Ordering::Release);
315                },
316            )
317            .unwrap();
318            assert!(released.load(Ordering::Acquire));
319            drop(guard);
320        });
321        assert_eq!(waits.load(Ordering::Acquire), 1);
322    }
323
324    /// The same image on the same host is one download; a different host or a
325    /// different image is not.
326    #[test]
327    fn the_pull_lock_is_shared_per_host_and_image() {
328        let first = image_pull_mutex(&ImageHost::LocalPodman, "ghcr.io/example/dev:latest");
329        let again = image_pull_mutex(&ImageHost::LocalPodman, "ghcr.io/example/dev:latest");
330        assert!(Arc::ptr_eq(&first, &again));
331
332        let other_image = image_pull_mutex(&ImageHost::LocalPodman, "ghcr.io/example/other:latest");
333        assert!(!Arc::ptr_eq(&first, &other_image));
334
335        let other_host = image_pull_mutex(&ImageHost::LocalDocker, "ghcr.io/example/dev:latest");
336        assert!(!Arc::ptr_eq(&first, &other_host));
337    }
338}