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}