Skip to main content

mtp_mount/daemon/
supervisor.rs

1//! The loop that owns every mount the daemon has.
2//!
3//! # The seam
4//!
5//! The supervisor never touches USB. Its whole input is a channel of
6//! [`Command`]s: a device arrived, a device left, a mount gave up on its
7//! device, stop. [`crate::daemon::usb`] is what turns `mtp-rs`'s hotplug stream
8//! into those commands in production.
9//!
10//! That's deliberate, and it's the only way this code is testable. USB hotplug
11//! can't be simulated: there's no way to make a container believe a phone was
12//! plugged in, so a supervisor that called `watch_devices()` itself could only
13//! ever be tested with a person and a cable. With the channel, a test sends
14//! `Command::Device(DeviceChange::Arrived(..))` and a real FUSE mount over a
15//! real (virtual) MTP device appears at a real path, with everything below the
16//! seam being the production code path. The same trick applies to the device
17//! itself through [`DeviceSource`]: production hands back a [`UsbOpener`], the
18//! tests hand back an opener for an `mtp-rs` virtual device.
19//!
20//! [`UsbOpener`]: crate::device::UsbOpener
21//!
22//! # Threading
23//!
24//! [`Supervisor::run`] blocks the thread it's called on and must NOT be called
25//! from inside a tokio runtime: opening a device and mounting it both go
26//! through `Handle::block_on`, the same sync-over-async bridge the FUSE
27//! callbacks use. The daemon runs it on the main thread and keeps the runtime
28//! for the hotplug watch and the signal handler.
29
30use std::collections::HashMap;
31use std::path::PathBuf;
32use std::sync::atomic::{AtomicBool, Ordering};
33use std::sync::mpsc::{Receiver, Sender};
34use std::sync::Arc;
35use std::time::Duration;
36
37use log::{debug, error, info, warn};
38
39use crate::daemon::unmount::{force_unmount, wait_until_unmounted};
40use crate::device::{DeviceOpener, UnplugSwitch};
41use crate::fs::{MtpFs, MtpFsConfig};
42use crate::hints::open_failure_hint;
43use crate::reconnect::ReconnectPolicy;
44
45/// How long a mount gets to leave the filesystem before the daemon complains.
46pub const DEFAULT_UNMOUNT_TIMEOUT: Duration = Duration::from_secs(10);
47
48/// A device the daemon knows about.
49#[derive(Debug, Clone, PartialEq, Eq)]
50pub struct DeviceIdent {
51    /// Identity and directory name in one: see [`crate::daemon::device_dir_name`].
52    ///
53    /// One string does both jobs so that "is this device already mounted?" and
54    /// "where does it go?" can't disagree. It has to be derivable from what an
55    /// arrival *and* a departure report, since a departure is the only thing
56    /// that says which mount to take down.
57    pub key: String,
58    /// What the device is called in the log.
59    pub label: String,
60    /// Serial number, when the device reports one. What a reopen matches on.
61    pub serial: Option<String>,
62}
63
64/// A device showed up or went away.
65#[derive(Debug, Clone)]
66pub enum DeviceChange {
67    Arrived(DeviceIdent),
68    Left(DeviceIdent),
69}
70
71/// Everything the supervisor reacts to.
72#[derive(Debug, Clone)]
73pub enum Command {
74    /// The device set changed.
75    Device(DeviceChange),
76    /// A mount decided its device is gone for good and asked to be taken down.
77    ///
78    /// The filesystem can't unmount itself (see [`crate::shutdown`]), and this
79    /// can beat the hotplug departure: a mount notices a dead session on its
80    /// next operation, while the USB watch only notices on its next poll.
81    GiveUp { key: String, reason: String },
82    /// Unmount everything and return from [`Supervisor::run`].
83    Stop(String),
84}
85
86/// Where the supervisor gets devices from.
87///
88/// Production returns a [`UsbOpener`](crate::device::UsbOpener) matched to the
89/// device's serial; tests return an opener for a virtual device.
90pub trait DeviceSource: Send + Sync {
91    /// The opener for a device the watch reported. It's used to open the device
92    /// now and, if the mount ever needs it, to reopen the same one later, so it
93    /// must resolve to that device and not to "whatever is plugged in".
94    fn opener(&self, ident: &DeviceIdent) -> Arc<dyn DeviceOpener>;
95}
96
97/// How the daemon mounts things.
98pub struct SupervisorConfig {
99    /// Directory that holds one subdirectory per mounted device.
100    pub mount_root: PathBuf,
101    /// Disk-backed spool for write buffers and read caches (see [`crate::spool`]).
102    pub spool_dir: PathBuf,
103    /// Mount every device read-only.
104    pub read_only: bool,
105    /// How long to wait for a mount to actually leave the filesystem.
106    pub unmount_timeout: Duration,
107}
108
109impl SupervisorConfig {
110    /// Config with the default unmount timeout.
111    pub fn new(mount_root: PathBuf, spool_dir: PathBuf, read_only: bool) -> Self {
112        Self {
113            mount_root,
114            spool_dir,
115            read_only,
116            unmount_timeout: DEFAULT_UNMOUNT_TIMEOUT,
117        }
118    }
119}
120
121/// One device's mount.
122struct ActiveMount {
123    path: PathBuf,
124    label: String,
125    session: fuser::BackgroundSession,
126    /// Tells the give-up watcher to stop, so it doesn't outlive the mount.
127    watcher_stop: Arc<AtomicBool>,
128}
129
130/// Mounts devices as they arrive and unmounts them as they leave.
131pub struct Supervisor {
132    config: SupervisorConfig,
133    source: Arc<dyn DeviceSource>,
134    rt: tokio::runtime::Handle,
135    /// Handed to each mount's give-up watcher so it can report back.
136    commands: Sender<Command>,
137    mounts: HashMap<String, ActiveMount>,
138}
139
140impl Supervisor {
141    /// Build a supervisor. `commands` must be a sender for the same channel
142    /// whose receiver goes to [`Supervisor::run`].
143    pub fn new(
144        config: SupervisorConfig,
145        source: Arc<dyn DeviceSource>,
146        rt: tokio::runtime::Handle,
147        commands: Sender<Command>,
148    ) -> Self {
149        Self {
150            config,
151            source,
152            rt,
153            commands,
154            mounts: HashMap::new(),
155        }
156    }
157
158    /// Handle commands until [`Command::Stop`] arrives or every sender is
159    /// dropped, then unmount everything.
160    ///
161    /// Blocks the calling thread. See the module docs on why it can't run
162    /// inside a tokio runtime.
163    pub fn run(mut self, commands: Receiver<Command>) {
164        // A mount root that can't be created isn't recoverable, but it also
165        // isn't worth stopping the process over before the caller has logged
166        // anything: every mount attempt reports it instead.
167        if let Err(e) = std::fs::create_dir_all(&self.config.mount_root) {
168            error!(
169                "Can't create the mount root {}: {e}",
170                self.config.mount_root.display()
171            );
172        }
173
174        let stop_reason = loop {
175            match commands.recv() {
176                Ok(Command::Device(DeviceChange::Arrived(ident))) => self.mount(ident),
177                Ok(Command::Device(DeviceChange::Left(ident))) => {
178                    self.unmount(&ident.key, "the device was unplugged")
179                }
180                Ok(Command::GiveUp { key, reason }) => self.unmount(&key, &reason),
181                Ok(Command::Stop(reason)) => break reason,
182                // Defensive: the supervisor keeps a sender of its own for the
183                // give-up watchers, so the channel can't actually close while
184                // this loop is running. Breaking beats spinning if that ever
185                // stops being true.
186                Err(_) => break "every event source is gone".to_string(),
187            }
188        };
189
190        info!("Shutting down: {stop_reason}");
191        self.unmount_all();
192    }
193
194    /// Devices currently mounted, by key. Used by tests.
195    pub fn mounted_keys(&self) -> Vec<String> {
196        let mut keys: Vec<String> = self.mounts.keys().cloned().collect();
197        keys.sort();
198        keys
199    }
200
201    fn mount(&mut self, ident: DeviceIdent) {
202        if let Some(existing) = self.mounts.get(&ident.key) {
203            warn!(
204                "{} is already mounted at {}; ignoring this arrival. \
205                 Two devices reporting the same serial number look identical from here.",
206                ident.label,
207                existing.path.display()
208            );
209            return;
210        }
211
212        let path = self.config.mount_root.join(&ident.key);
213        if let Err(e) = std::fs::create_dir_all(&path) {
214            error!("Can't create the mount point {}: {e}", path.display());
215            return;
216        }
217
218        let opener = self.source.opener(&ident);
219        let device = match opener.open(&self.rt) {
220            Ok(device) => device,
221            Err(e) => {
222                error!("Can't open {}: {e}", ident.label);
223                if let Some(hint) = open_failure_hint(&e) {
224                    error!("{hint}");
225                }
226                let _ = std::fs::remove_dir(&path);
227                return;
228            }
229        };
230
231        let mtp_fs = MtpFs::new(
232            device,
233            opener,
234            self.rt.clone(),
235            MtpFsConfig {
236                read_only: self.config.read_only,
237                spool_dir: self.config.spool_dir.clone(),
238                // Reconnect stays off here for the same reason it's off in the
239                // CLI, only more so: waiting blocks every process touching the
240                // mount, and the daemon has a better answer than waiting. A
241                // device that comes back arrives as a fresh hotplug event and
242                // gets mounted again at the same path, with nothing frozen in
243                // between. Turning this on would trade a mount that reappears
244                // for a desktop that hangs.
245                reconnect: ReconnectPolicy::from_secs(0),
246                // The pretend cable is a test seam for the CLI's reconnect path;
247                // nothing here ever flips it.
248                unplug: UnplugSwitch::default(),
249            },
250        );
251
252        let shutdown = mtp_fs.shutdown();
253        let mut fuse_config = fuser::Config::default();
254        fuse_config.mount_options = mtp_fs.mount_options();
255
256        let session = match fuser::spawn_mount2(mtp_fs, &path, &fuse_config) {
257            Ok(session) => session,
258            Err(e) => {
259                error!("Can't mount {} at {}: {e}", ident.label, path.display());
260                let _ = std::fs::remove_dir(&path);
261                return;
262            }
263        };
264
265        // The mount raises its shutdown signal from inside a FUSE callback and
266        // can't act on it (see `crate::shutdown`), so one thread per mount turns
267        // that signal into a command the supervisor can act on.
268        let watcher_stop = Arc::new(AtomicBool::new(false));
269        {
270            let watcher_stop = Arc::clone(&watcher_stop);
271            let commands = self.commands.clone();
272            let key = ident.key.clone();
273            std::thread::spawn(move || loop {
274                if let Some(reason) = shutdown.wait_timeout(Duration::from_millis(200)) {
275                    let _ = commands.send(Command::GiveUp { key, reason });
276                    return;
277                }
278                if watcher_stop.load(Ordering::Relaxed) {
279                    return;
280                }
281            });
282        }
283
284        info!("Mounted {} at {}", ident.label, path.display());
285        self.mounts.insert(
286            ident.key,
287            ActiveMount {
288                path,
289                label: ident.label,
290                session,
291                watcher_stop,
292            },
293        );
294    }
295
296    fn unmount(&mut self, key: &str, reason: &str) {
297        let Some(mount) = self.mounts.remove(key) else {
298            debug!("Nothing mounted for {key}; nothing to unmount ({reason})");
299            return;
300        };
301        let ActiveMount {
302            path,
303            label,
304            session,
305            watcher_stop,
306        } = mount;
307        watcher_stop.store(true, Ordering::Relaxed);
308        info!("Unmounting {label} from {}: {reason}", path.display());
309
310        // The forced unmount comes first, and it's ours, not `fuser`'s.
311        // `fuser`'s own unmount is a plain `umount()`, which fails with EBUSY
312        // whenever anything still holds the mount, and a busy mount is the
313        // normal case here: the reason the device left is usually that someone
314        // yanked the cable mid-copy. [`force_unmount`] is the lazy detach
315        // (`umount2(MNT_DETACH)`, or `fusermount3 -u -z` when the syscall is
316        // refused), which succeeds regardless and leaves the in-flight callers
317        // to get their error.
318        //
319        // `umount_and_join` still runs afterwards, on its own thread, to reap
320        // the session: its unmount is a no-op on an already-detached mount, and
321        // the *join* can take as long as the FUSE callback in progress takes to
322        // notice its device is gone. The daemon must not block on that, because
323        // the next device's mount is queued behind this loop.
324        if let Err(e) = force_unmount(&path) {
325            error!("Can't unmount {}: {e}", path.display());
326        }
327        let joining_path = path.clone();
328        std::thread::spawn(move || match session.umount_and_join() {
329            Ok(()) => debug!("The session for {} ended cleanly", joining_path.display()),
330            Err(e) => debug!("The session for {} ended with {e}", joining_path.display()),
331        });
332
333        if wait_until_unmounted(&path, self.config.unmount_timeout) {
334            // Only now is the directory a plain empty directory again, so this
335            // is also a second check: `remove_dir` on a live mount point fails.
336            if let Err(e) = std::fs::remove_dir(&path) {
337                warn!(
338                    "Unmounted {label}, but {} is still there: {e}",
339                    path.display()
340                );
341            }
342        } else {
343            error!(
344                "{} is STILL mounted {}s after unmounting it. \
345                 Anything touching that path may hang; unmount it by hand with \
346                 `fusermount3 -u -z {}`.",
347                path.display(),
348                self.config.unmount_timeout.as_secs(),
349                path.display()
350            );
351        }
352    }
353
354    fn unmount_all(&mut self) {
355        let keys: Vec<String> = self.mounts.keys().cloned().collect();
356        for key in keys {
357            self.unmount(&key, "the daemon is shutting down");
358        }
359        // Tidy, and a signal to anything watching the root that nothing is
360        // mounted any more. Fails harmlessly if the root isn't empty.
361        let _ = std::fs::remove_dir(&self.config.mount_root);
362    }
363}
364
365#[cfg(test)]
366mod tests {
367    use super::*;
368
369    /// A source that's never asked for anything: the tests here only exercise
370    /// the bookkeeping that happens before a device is opened.
371    struct NoDevices;
372
373    impl DeviceSource for NoDevices {
374        fn opener(&self, _ident: &DeviceIdent) -> Arc<dyn DeviceOpener> {
375            unreachable!("these tests never get as far as opening a device")
376        }
377    }
378
379    fn ident(key: &str) -> DeviceIdent {
380        DeviceIdent {
381            key: key.to_string(),
382            label: format!("device {key}"),
383            serial: Some(key.to_string()),
384        }
385    }
386
387    fn supervisor(root: PathBuf) -> (Supervisor, Sender<Command>, Receiver<Command>) {
388        let (tx, rx) = std::sync::mpsc::channel();
389        let rt = tokio::runtime::Runtime::new().unwrap();
390        let handle = rt.handle().clone();
391        std::mem::forget(rt);
392        let supervisor = Supervisor::new(
393            SupervisorConfig::new(root, std::env::temp_dir(), false),
394            Arc::new(NoDevices),
395            handle,
396            tx.clone(),
397        );
398        (supervisor, tx, rx)
399    }
400
401    #[test]
402    fn a_departure_for_a_device_that_was_never_mounted_is_ignored() {
403        let root = tempfile::tempdir().unwrap();
404        let (mut supervisor, _tx, _rx) = supervisor(root.path().to_path_buf());
405        supervisor.unmount("ABC123", "test");
406        assert!(supervisor.mounted_keys().is_empty());
407    }
408
409    #[test]
410    fn stopping_with_nothing_mounted_returns() {
411        let root = tempfile::tempdir().unwrap();
412        let (supervisor, tx, rx) = supervisor(root.path().to_path_buf());
413        tx.send(Command::Stop("test".into())).unwrap();
414        drop(tx);
415        supervisor.run(rx);
416    }
417
418    #[test]
419    fn the_mount_root_is_created_when_the_loop_starts() {
420        let parent = tempfile::tempdir().unwrap();
421        let root = parent.path().join("runtime").join("mtp");
422        let (supervisor, tx, rx) = supervisor(root.clone());
423        tx.send(Command::Stop("test".into())).unwrap();
424        drop(tx);
425        supervisor.run(rx);
426        // `unmount_all` removes it again on the way out, so check the parent
427        // chain, which is what proves `create_dir_all` ran.
428        assert!(root.parent().unwrap().is_dir());
429    }
430
431    #[test]
432    fn a_device_that_cannot_be_opened_leaves_no_directory_behind() {
433        // The failing-open path is the one that runs when a phone is locked or
434        // gvfs got there first: it must not leave an empty mount point in the
435        // root for a file manager to show as a device that isn't there.
436        struct NeverOpens;
437        impl DeviceSource for NeverOpens {
438            fn opener(&self, _ident: &DeviceIdent) -> Arc<dyn DeviceOpener> {
439                struct Refuses;
440                impl DeviceOpener for Refuses {
441                    fn open(
442                        &self,
443                        _rt: &tokio::runtime::Handle,
444                    ) -> Result<mtp_rs::mtp::MtpDevice, mtp_rs::Error> {
445                        Err(mtp_rs::Error::ExclusiveAccess)
446                    }
447                    fn describe(&self) -> String {
448                        "a device that won't open".into()
449                    }
450                }
451                Arc::new(Refuses)
452            }
453        }
454
455        let root = tempfile::tempdir().unwrap();
456        let (tx, _rx) = std::sync::mpsc::channel();
457        let rt = tokio::runtime::Runtime::new().unwrap();
458        let handle = rt.handle().clone();
459        let mut supervisor = Supervisor::new(
460            SupervisorConfig::new(root.path().to_path_buf(), std::env::temp_dir(), false),
461            Arc::new(NeverOpens),
462            handle,
463            tx,
464        );
465
466        supervisor.mount(ident("ABC123"));
467
468        assert!(supervisor.mounted_keys().is_empty());
469        assert!(!root.path().join("ABC123").exists());
470    }
471}