Skip to main content

mtp_mount/
fs.rs

1use std::collections::HashMap;
2use std::ffi::OsStr;
3use std::io;
4use std::io::{Seek, SeekFrom};
5use std::path::PathBuf;
6use std::pin::Pin;
7use std::sync::atomic::{AtomicU64, Ordering};
8use std::sync::{Arc, Mutex};
9use std::time::{Duration, SystemTime, UNIX_EPOCH};
10
11use bytes::Bytes;
12use fuser::{
13    BsdFileFlags, Errno, FileAttr, FileHandle, FileType, Filesystem, FopenFlags, Generation,
14    INodeNo, KernelConfig, LockOwner, MountOption, OpenFlags, RenameFlags, ReplyAttr, ReplyCreate,
15    ReplyData, ReplyDirectory, ReplyEmpty, ReplyEntry, ReplyOpen, ReplyStatfs, ReplyWrite, Request,
16    TimeOrNow, WriteFlags,
17};
18use log::{debug, error, info, warn};
19use mtp_rs::mtp::{DeviceEvent, MtpDevice};
20use mtp_rs::{NewObjectInfo, ObjectHandle, Storage};
21
22use crate::buffer::WriteBuffer;
23use crate::device::{is_link_lost, DeviceOpener, UnplugSwitch};
24use crate::inode::{ChildInfo, InodeEntry, InodeKind, InodeTable, FUSE_ROOT_INODE};
25use crate::reconnect::ReconnectPolicy;
26use crate::shutdown::Shutdown;
27use crate::sparse_cache::SparseCache;
28
29const TTL: Duration = Duration::from_secs(1);
30
31/// How much of a spool file an upload holds in memory at a time.
32const UPLOAD_CHUNK: usize = 65536;
33
34/// How many times an operation is retried across reconnects before it gives up.
35/// One reconnect is the cable glitch we're here for; a second is a device that
36/// keeps dropping mid-operation, and past that the retry is not the answer.
37const MAX_ATTEMPTS: u32 = 3;
38
39/// How many times an operation re-resolves its handles and tries again after
40/// the device says a handle is stale.
41///
42/// One, which is what `mtp-rs` prescribes: the first `StaleHandle` means the
43/// device re-keyed the object and a fresh listing has the new token, so the
44/// retry works. A second one for the same operation means the re-resolved token
45/// died too, which isn't a re-key any more; looping on it would hammer the
46/// device instead of telling the caller.
47///
48/// This budget is deliberately separate from [`MAX_ATTEMPTS`]: the two failures
49/// have nothing in common. A stale handle costs one listing against a healthy
50/// session, a dead link costs a reopen and a backoff wait, so spending one
51/// shouldn't shorten the other, and a stale handle must never fall through into
52/// the reconnect path.
53const MAX_STALE_RETRIES: u32 = 1;
54
55type MtpResult<T> = Result<T, mtp_rs::Error>;
56
57/// Everything the mount needs beyond the device itself.
58pub struct MtpFsConfig {
59    pub read_only: bool,
60    /// Disk-backed directory for write buffers and read caches (see [`crate::spool`]).
61    /// Must exist and be writable; resolve it with [`crate::spool::spool_dir_from_env`]
62    /// and [`crate::spool::prepare_spool_dir`].
63    pub spool_dir: PathBuf,
64    /// How long to wait for a device that went away.
65    pub reconnect: ReconnectPolicy,
66    /// The pretend cable, shared with the [`DeviceOpener`] (tests only).
67    pub unplug: UnplugSwitch,
68}
69
70fn mtp_datetime_to_system_time(dt: &mtp_rs::DateTime) -> SystemTime {
71    fn days_from_civil(y: i64, m: i64, d: i64) -> i64 {
72        let y = if m <= 2 { y - 1 } else { y };
73        let era = if y >= 0 { y } else { y - 399 } / 400;
74        let yoe = (y - era * 400) as u64;
75        let m_adj = if m > 2 { m - 3 } else { m + 9 } as u64;
76        let doy = (153 * m_adj + 2) / 5 + d as u64 - 1;
77        let doe = yoe * 365 + yoe / 4 - yoe / 100 + doy;
78        era * 146097 + doe as i64 - 719468
79    }
80
81    let days = days_from_civil(dt.year as i64, dt.month as i64, dt.day as i64);
82    let secs = days * 86400 + dt.hour as i64 * 3600 + dt.minute as i64 * 60 + dt.second as i64;
83    if secs >= 0 {
84        UNIX_EPOCH + Duration::from_secs(secs as u64)
85    } else {
86        UNIX_EPOCH
87    }
88}
89
90fn inode_to_file_attr(entry: &InodeEntry) -> FileAttr {
91    let uid = unsafe { libc::getuid() };
92    let gid = unsafe { libc::getgid() };
93    FileAttr {
94        ino: INodeNo(entry.inode),
95        size: entry.size,
96        blocks: entry.size.div_ceil(512),
97        atime: entry.atime,
98        mtime: entry.mtime,
99        ctime: entry.mtime,
100        crtime: entry.mtime,
101        kind: if entry.is_dir() {
102            FileType::Directory
103        } else {
104            FileType::RegularFile
105        },
106        perm: if entry.is_dir() { 0o755 } else { 0o644 },
107        nlink: if entry.is_dir() { 2 } else { 1 },
108        uid,
109        gid,
110        rdev: 0,
111        blksize: 4096,
112        flags: 0,
113    }
114}
115
116/// The name a storage shows up under in the mount.
117fn storage_name(storage: &Storage) -> String {
118    if storage.info().description.is_empty() {
119        format!("Storage_{}", storage.id().0)
120    } else {
121        storage.info().description.clone()
122    }
123}
124
125/// The temporary name a safe flush uploads under before renaming into place.
126fn temp_upload_name(name: &str) -> String {
127    format!(".~tmp~{name}")
128}
129
130/// Wraps a local I/O failure as an MTP error so spool problems flow through the
131/// same result type as device problems.
132fn io_error(e: io::Error) -> mtp_rs::Error {
133    mtp_rs::Error::Io {
134        message: e.to_string(),
135    }
136}
137
138/// Helper to create an `Unpin` stream from a `Vec<u8>`.
139fn bytes_stream(
140    data: Vec<u8>,
141) -> futures::stream::Iter<std::vec::IntoIter<Result<Bytes, io::Error>>> {
142    let chunks = if data.is_empty() {
143        vec![Ok(Bytes::new())]
144    } else {
145        vec![Ok(Bytes::from(data))]
146    };
147    futures::stream::iter(chunks)
148}
149
150/// Streams a spool file in [`UPLOAD_CHUNK`]-sized pieces, reading each one only
151/// when the consumer asks for it. That's what keeps an upload's memory flat: a
152/// 4 GB file costs one chunk of RAM, not 4 GB.
153///
154/// The read is a blocking `std::fs` read, which is fine because the only caller
155/// polls this from `Handle::block_on` on a FUSE callback thread that has nothing
156/// else to do. Don't spawn this stream onto the runtime: it would park a worker.
157fn file_stream(
158    file: std::fs::File,
159) -> Pin<Box<dyn futures::Stream<Item = Result<Bytes, io::Error>> + Send>> {
160    use std::io::Read as _;
161    // The `Option` is the terminator. Handing the file back only after a
162    // successful read means EOF and errors both end the stream; a version that
163    // kept the file after an error would re-emit that same error forever.
164    Box::pin(futures::stream::unfold(Some(file), |state| async move {
165        let mut file = state?;
166        let mut buf = vec![0u8; UPLOAD_CHUNK];
167        match file.read(&mut buf) {
168            Ok(0) => None,
169            Ok(n) => {
170                buf.truncate(n);
171                Some((Ok(Bytes::from(buf)), Some(file)))
172            }
173            Err(e) => Some((Err(e), None)),
174        }
175    }))
176}
177
178/// Mutable state protected by `RefCell` so fuser's `&self` callbacks can mutate it.
179struct Inner {
180    storages: Vec<Storage>,
181    /// Disk-backed directory for write buffers and read caches (see [`crate::spool`]).
182    spool_dir: PathBuf,
183    inodes: InodeTable,
184    write_buf: WriteBuffer,
185    read_cache: HashMap<u64, SparseCache>,
186    dirs_loaded: HashMap<u64, bool>,
187    fh_to_inode: HashMap<u64, u64>,
188}
189
190/// FUSE filesystem backed by an MTP device.
191pub struct MtpFs {
192    rt: tokio::runtime::Handle,
193    device: Mutex<MtpDevice>,
194    /// How to reopen the same device after it goes away.
195    opener: Arc<dyn DeviceOpener>,
196    policy: ReconnectPolicy,
197    unplug: UnplugSwitch,
198    /// Raised when the device is gone for good, so whoever owns the mount takes
199    /// it down instead of leaving something that answers every call with EIO.
200    shutdown: Arc<Shutdown>,
201    /// Bumped on every reconnect so the event loop from the previous session
202    /// notices it's been superseded and exits.
203    event_epoch: Arc<AtomicU64>,
204    inner: Arc<Mutex<Inner>>,
205    next_fh: AtomicU64,
206    read_only: bool,
207    /// Counter incremented on every MTP partial-read fetch. Used by integration
208    /// tests to verify that the sparse cache prevents redundant fetches.
209    fetch_counter: Arc<AtomicU64>,
210}
211
212impl MtpFs {
213    /// Builds a filesystem over an already-open device.
214    ///
215    /// `opener` must resolve to that same device; it's what the mount uses to
216    /// come back after a disconnect.
217    pub fn new(
218        device: MtpDevice,
219        opener: Arc<dyn DeviceOpener>,
220        rt: tokio::runtime::Handle,
221        config: MtpFsConfig,
222    ) -> Self {
223        let MtpFsConfig {
224            read_only,
225            spool_dir,
226            reconnect,
227            unplug,
228        } = config;
229        Self {
230            rt,
231            device: Mutex::new(device),
232            opener,
233            policy: reconnect,
234            unplug,
235            shutdown: Arc::new(Shutdown::default()),
236            event_epoch: Arc::new(AtomicU64::new(0)),
237            inner: Arc::new(Mutex::new(Inner {
238                storages: Vec::new(),
239                spool_dir: spool_dir.clone(),
240                inodes: InodeTable::new(),
241                write_buf: WriteBuffer::new(spool_dir),
242                read_cache: HashMap::new(),
243                dirs_loaded: HashMap::new(),
244                fh_to_inode: HashMap::new(),
245            })),
246            next_fh: AtomicU64::new(1),
247            read_only,
248            fetch_counter: Arc::new(AtomicU64::new(0)),
249        }
250    }
251
252    /// The signal that asks for this mount to be taken down.
253    ///
254    /// Whoever mounted the filesystem owns the unmount handle, so it has to
255    /// watch this and unmount when a reason shows up. See [`crate::shutdown`].
256    pub fn shutdown(&self) -> Arc<Shutdown> {
257        Arc::clone(&self.shutdown)
258    }
259
260    /// Returns a shared handle to the MTP fetch counter.
261    ///
262    /// The counter increments each time a partial-read operation is issued to
263    /// the device. Primarily used by integration tests to verify cache behavior.
264    #[allow(dead_code)] // used by integration tests via lib.rs, not by the bin
265    pub fn fetch_counter(&self) -> Arc<AtomicU64> {
266        Arc::clone(&self.fetch_counter)
267    }
268
269    fn alloc_fh(&self) -> u64 {
270        self.next_fh.fetch_add(1, Ordering::Relaxed)
271    }
272
273    /// Find the storage index that owns a given inode by walking up the tree.
274    fn find_storage_index(inner: &Inner, inode: u64) -> Option<usize> {
275        let mut current = inode;
276        loop {
277            let entry = inner.inodes.get(current)?;
278            if let InodeKind::Storage { storage_id } = &entry.kind {
279                return inner
280                    .storages
281                    .iter()
282                    .position(|s: &Storage| s.id() == *storage_id);
283            }
284            if current == entry.parent {
285                return None;
286            }
287            current = entry.parent;
288        }
289    }
290
291    /// Runs an MTP operation, riding out the two ways its handles can die under
292    /// it: the device went away, or the device re-keyed the object.
293    ///
294    /// `attempt` is re-run from scratch after either recovery, so it must
295    /// resolve its own handles (through [`Self::file_handle`] and friends) rather
296    /// than close over handles from the previous try.
297    ///
298    /// The two paths are siblings, not variations. A dead link needs a reopen
299    /// and a wait; a stale handle needs neither, because the session is fine and
300    /// only the token is dead, so it re-resolves by path and retries straight
301    /// away. Reopening a healthy device would be a real regression: on Android
302    /// a reopen is expensive and can wedge the device.
303    fn with_recovery<T>(
304        &self,
305        inner: &mut Inner,
306        mut attempt: impl FnMut(&Self, &mut Inner) -> MtpResult<T>,
307    ) -> MtpResult<T> {
308        for _ in 0..MAX_ATTEMPTS {
309            // Stale-handle retries happen against the current session, on their
310            // own budget, and never fall through to the reconnect below.
311            let mut stale_retries = MAX_STALE_RETRIES;
312            loop {
313                if self.unplug.is_unplugged() {
314                    break;
315                }
316                match attempt(self, inner) {
317                    Ok(value) => return Ok(value),
318                    Err(e) if e.is_stale_handle() => {
319                        if stale_retries == 0 {
320                            error!("Operation still hit a stale object handle after re-resolving");
321                            return Err(e);
322                        }
323                        stale_retries -= 1;
324                        debug!("Operation hit a stale object handle, re-resolving by path");
325                        Self::invalidate_handles(inner);
326                    }
327                    Err(e) if !is_link_lost(&e) => return Err(e),
328                    Err(e) => {
329                        debug!("Operation hit a dead session: {e}");
330                        break;
331                    }
332                }
333            }
334            self.reconnect(inner)?;
335        }
336        Err(mtp_rs::Error::Disconnected)
337    }
338
339    /// Marks every cached object handle stale so the next use re-resolves it by
340    /// path, and drops the cached listings that produced them.
341    ///
342    /// Whole-table rather than just the inode that failed, for two reasons. A
343    /// device that re-keys re-keys in batches (Android's MediaProvider does it
344    /// across a whole media rescan), so the neighbours are suspect too. And the
345    /// generation counter makes marking free: re-resolution is lazy, so the only
346    /// inodes that pay for a listing are the ones something actually touches
347    /// again. This is the same mechanism a reconnect uses, minus the reopen.
348    fn invalidate_handles(inner: &mut Inner) {
349        inner.inodes.bump_generation();
350        inner.dirs_loaded.clear();
351        inner.dirs_loaded.insert(FUSE_ROOT_INODE, true);
352    }
353
354    /// Waits for the device to come back and rebuilds the session on top of the
355    /// existing inode tree. Gives up (and unmounts) when the window runs out.
356    fn reconnect(&self, inner: &mut Inner) -> MtpResult<()> {
357        let who = self.opener.describe();
358
359        if self.policy.is_disabled() {
360            self.give_up(&format!(
361                "{who} disconnected and reconnect is off (--reconnect-timeout 0)"
362            ));
363            return Err(mtp_rs::Error::Disconnected);
364        }
365
366        let secs = self.policy.timeout().as_secs();
367        info!("{who} disconnected, waiting up to {secs}s for it to come back...");
368        eprintln!("mtp-mount: {who} disconnected, waiting up to {secs}s for it to come back...");
369
370        for delay in self.policy.schedule() {
371            std::thread::sleep(delay);
372            if self.unplug.is_unplugged() {
373                continue;
374            }
375            let device = match self.opener.open(&self.rt) {
376                Ok(device) => device,
377                Err(e) => {
378                    debug!("Reconnect attempt failed: {e}");
379                    continue;
380                }
381            };
382            match self.adopt(inner, device) {
383                Ok(()) => {
384                    eprintln!("mtp-mount: {who} is back, carrying on.");
385                    return Ok(());
386                }
387                Err(e) => {
388                    warn!("Reopened {who} but couldn't resume the mount: {e}");
389                    continue;
390                }
391            }
392        }
393
394        self.give_up(&format!("{who} didn't come back within {secs}s"));
395        Err(mtp_rs::Error::Disconnected)
396    }
397
398    /// Takes over a freshly opened device: re-maps storage IDs, marks every
399    /// cached object handle stale, and starts a new event loop.
400    ///
401    /// Inode numbers, names, the tree shape, open file handles, read caches, and
402    /// write spools all survive untouched. Only the session-scoped tokens change.
403    fn adopt(&self, inner: &mut Inner, device: MtpDevice) -> MtpResult<()> {
404        let storages = self.rt.block_on(device.storages())?;
405
406        // Storage IDs are session-scoped too. Match the new storages to the
407        // storage inodes by name, falling back to position when a device
408        // reports no description (the old name embeds the old ID).
409        let storage_inodes = inner.inodes.children(FUSE_ROOT_INODE);
410        for (position, storage_ino) in storage_inodes.iter().enumerate() {
411            let name = match inner.inodes.get(*storage_ino) {
412                Some(entry) => entry.name.clone(),
413                None => continue,
414            };
415            let matched = storages
416                .iter()
417                .find(|s| storage_name(s) == name)
418                .or_else(|| storages.get(position));
419            match matched {
420                Some(storage) => inner.inodes.set_storage_id(*storage_ino, storage.id()),
421                None => warn!("Storage '{name}' is missing after the reconnect"),
422            }
423        }
424
425        inner.storages = storages;
426        *self.device.lock().unwrap() = device.clone();
427        inner.inodes.bump_generation();
428
429        // Names and sizes are re-read on the next access; the handles behind
430        // them are re-resolved lazily by path.
431        inner.dirs_loaded.clear();
432        inner.dirs_loaded.insert(FUSE_ROOT_INODE, true);
433
434        let epoch = self.event_epoch.fetch_add(1, Ordering::SeqCst) + 1;
435        self.spawn_event_loop(device, epoch);
436        Ok(())
437    }
438
439    /// Says why the mount is going away and asks for it to be taken down. The
440    /// operation that triggered this still returns an error to its caller.
441    fn give_up(&self, reason: &str) {
442        error!("{reason}; unmounting");
443        eprintln!("mtp-mount: {reason}. Unmounting.");
444        self.shutdown.request(reason);
445    }
446
447    /// Re-resolves an inode's MTP handle by path if it came from an older
448    /// session, walking down from the storage root and refreshing each ancestor
449    /// on the way. Inode numbers never change here.
450    fn ensure_fresh(&self, inner: &mut Inner, inode: u64) -> MtpResult<()> {
451        if inner.inodes.is_fresh(inode) {
452            return Ok(());
453        }
454
455        // Collect the chain from the storage root down to this inode.
456        let mut chain = Vec::new();
457        let mut current = inode;
458        loop {
459            let entry = inner.inodes.get(current).ok_or(mtp_rs::Error::NotFound)?;
460            match entry.kind {
461                InodeKind::Root | InodeKind::Storage { .. } => break,
462                _ => {
463                    chain.push(current);
464                    if current == entry.parent {
465                        return Err(mtp_rs::Error::NotFound);
466                    }
467                    current = entry.parent;
468                }
469            }
470        }
471        chain.reverse();
472
473        let storage_idx = Self::find_storage_index(inner, inode).ok_or(mtp_rs::Error::NotFound)?;
474
475        let mut parent_handle: Option<ObjectHandle> = None;
476        for node in chain {
477            let entry = inner.inodes.get(node).ok_or(mtp_rs::Error::NotFound)?;
478            let name = entry.name.clone();
479            if inner.inodes.is_fresh(node) {
480                parent_handle = match inner.inodes.get(node).map(|e| &e.kind) {
481                    Some(InodeKind::Directory { handle } | InodeKind::File { handle }) => {
482                        Some(*handle)
483                    }
484                    _ => return Err(mtp_rs::Error::NotFound),
485                };
486                continue;
487            }
488
489            let objects = self
490                .rt
491                .block_on(inner.storages[storage_idx].list_objects(parent_handle))?;
492            let found = objects
493                .into_iter()
494                .find(|obj| obj.filename == name)
495                .ok_or(mtp_rs::Error::NotFound)?;
496            inner.inodes.set_handle(node, found.handle);
497            parent_handle = Some(found.handle);
498        }
499
500        Ok(())
501    }
502
503    /// The current handle of a file inode.
504    fn file_handle(&self, inner: &mut Inner, inode: u64) -> MtpResult<ObjectHandle> {
505        self.ensure_fresh(inner, inode)?;
506        match inner.inodes.get(inode).map(|e| &e.kind) {
507            Some(InodeKind::File { handle }) => Ok(*handle),
508            _ => Err(mtp_rs::Error::NotFound),
509        }
510    }
511
512    /// The current handle of a file or directory inode.
513    fn object_handle(&self, inner: &mut Inner, inode: u64) -> MtpResult<ObjectHandle> {
514        self.ensure_fresh(inner, inode)?;
515        match inner.inodes.get(inode).map(|e| &e.kind) {
516            Some(InodeKind::File { handle } | InodeKind::Directory { handle }) => Ok(*handle),
517            _ => Err(mtp_rs::Error::NotFound),
518        }
519    }
520
521    /// The MTP parent handle to use for operations inside a directory inode
522    /// (`None` means the storage root).
523    fn parent_handle(&self, inner: &mut Inner, inode: u64) -> MtpResult<Option<ObjectHandle>> {
524        self.ensure_fresh(inner, inode)?;
525        match inner.inodes.get(inode).map(|e| &e.kind) {
526            Some(InodeKind::Storage { .. }) => Ok(None),
527            Some(InodeKind::Directory { handle }) => Ok(Some(*handle)),
528            _ => Err(mtp_rs::Error::NotFound),
529        }
530    }
531
532    /// Load children of a directory from MTP into the inode table.
533    fn load_dir(&self, inner: &mut Inner, parent_inode: u64) {
534        if inner.dirs_loaded.get(&parent_inode) == Some(&true) {
535            return;
536        }
537
538        if parent_inode == FUSE_ROOT_INODE {
539            inner.dirs_loaded.insert(parent_inode, true);
540            return;
541        }
542
543        match self.with_recovery(inner, |fs, inner| fs.list_into_table(inner, parent_inode)) {
544            Ok(()) => {
545                inner.dirs_loaded.insert(parent_inode, true);
546            }
547            Err(e) => error!("Failed to list MTP objects: {e}"),
548        }
549    }
550
551    /// One listing pass: ask the device, then reconcile the inode table.
552    fn list_into_table(&self, inner: &mut Inner, parent_inode: u64) -> MtpResult<()> {
553        let mtp_parent = self.parent_handle(inner, parent_inode)?;
554        let storage_idx =
555            Self::find_storage_index(inner, parent_inode).ok_or(mtp_rs::Error::NotFound)?;
556
557        let objects = self
558            .rt
559            .block_on(inner.storages[storage_idx].list_objects(mtp_parent))?;
560
561        let children: Vec<ChildInfo> = objects
562            .into_iter()
563            .map(|obj| ChildInfo {
564                handle: obj.handle,
565                is_dir: obj.is_folder(),
566                size: obj.size,
567                mtime: obj
568                    .modified
569                    .as_ref()
570                    .map(mtp_datetime_to_system_time)
571                    .unwrap_or(UNIX_EPOCH),
572                name: obj.filename,
573            })
574            .collect();
575
576        inner.inodes.sync_children(parent_inode, &children);
577        Ok(())
578    }
579
580    /// Flush a dirty write buffer to MTP.
581    ///
582    /// When the device supports rename, uses a safe upload-then-delete-then-rename
583    /// sequence to avoid data loss if the upload fails. Falls back to
584    /// delete-then-upload on devices without rename support.
585    ///
586    /// The spooled bytes live in an unlinked temp file that a disconnect can't
587    /// touch, so an upload interrupted by a cable glitch is retried from the
588    /// start once the device is back. Only the upload is retried: once it lands,
589    /// the data is on the device and a second attempt would just duplicate it.
590    fn flush_to_mtp(&self, inner: &mut Inner, fh: u64) -> MtpResult<()> {
591        let buf = match inner.write_buf.close(fh) {
592            Some(b) => b,
593            None => return Ok(()),
594        };
595
596        if !buf.is_dirty() {
597            return Ok(());
598        }
599
600        let inode = buf.inode;
601        let mut file = buf.into_file();
602        let file_len = file.seek(SeekFrom::End(0)).map_err(io_error)?;
603
604        let entry = match inner.inodes.get(inode) {
605            Some(e) => e.clone(),
606            None => {
607                error!("Flush: inode {inode} not found");
608                return Err(mtp_rs::Error::NotFound);
609            }
610        };
611
612        let mut attempts = 0u32;
613        self.with_recovery(inner, |fs, inner| {
614            let mut attempt = file.try_clone().map_err(io_error)?;
615            attempt.seek(SeekFrom::Start(0)).map_err(io_error)?;
616            attempts += 1;
617            fs.flush_once(inner, inode, &entry, file_len, attempt, attempts > 1)
618        })
619    }
620
621    /// One flush attempt against the current session.
622    ///
623    /// `is_retry` means an earlier attempt died mid-upload, so a half-written
624    /// temp object may be sitting in the target directory; it's cleared out
625    /// before uploading again.
626    #[allow(clippy::too_many_arguments)]
627    fn flush_once(
628        &self,
629        inner: &mut Inner,
630        inode: u64,
631        entry: &InodeEntry,
632        size: u64,
633        file: std::fs::File,
634        is_retry: bool,
635    ) -> MtpResult<()> {
636        let handle = self.file_handle(inner, inode)?;
637        let storage_idx = Self::find_storage_index(inner, inode).ok_or(mtp_rs::Error::NotFound)?;
638        let parent_handle = match self.parent_handle(inner, entry.parent) {
639            Ok(handle) => handle,
640            Err(e) => {
641                error!("Flush: no parent directory for inode {inode}: {e}");
642                return Err(e);
643            }
644        };
645
646        let supports_rename = self.device.lock().unwrap().supports_rename();
647
648        if is_retry && supports_rename {
649            self.purge_leftover(
650                inner,
651                storage_idx,
652                parent_handle,
653                &temp_upload_name(&entry.name),
654            );
655        }
656
657        if supports_rename {
658            self.flush_safe(
659                inner,
660                inode,
661                handle,
662                storage_idx,
663                parent_handle,
664                entry,
665                size,
666                file,
667            )
668        } else {
669            warn!(
670                "Flush: device does not support rename, using delete-then-upload \
671                 (data loss possible if upload fails)"
672            );
673            self.flush_unsafe(
674                inner,
675                inode,
676                handle,
677                storage_idx,
678                parent_handle,
679                entry,
680                size,
681                file,
682            )
683        }
684    }
685
686    /// Safe flush: upload with temp name, delete old, rename new.
687    ///
688    /// Only the upload returns an error to the caller: after it lands, the bytes
689    /// are safe on the device and a retry would upload them twice, so the later
690    /// steps report what happened and return `Ok`.
691    #[allow(clippy::too_many_arguments)]
692    fn flush_safe(
693        &self,
694        inner: &mut Inner,
695        inode: u64,
696        old_handle: ObjectHandle,
697        storage_idx: usize,
698        parent_handle: Option<ObjectHandle>,
699        entry: &InodeEntry,
700        size: u64,
701        file: std::fs::File,
702    ) -> MtpResult<()> {
703        let storage = &inner.storages[storage_idx];
704        let temp_name = temp_upload_name(&entry.name);
705
706        // Step 1: Upload new data with a temp name.
707        let info = NewObjectInfo::file(&temp_name, size);
708        let stream = file_stream(file);
709        let new_handle = match self
710            .rt
711            .block_on(storage.upload(parent_handle, info, stream))
712        {
713            Ok(h) => h,
714            Err(e) => {
715                error!("Flush: upload failed (original file untouched): {e}");
716                return Err(e.into());
717            }
718        };
719
720        // Step 2: Delete old object.
721        if let Err(e) = self.rt.block_on(storage.delete(old_handle)) {
722            error!("Flush: failed to delete old object (new data saved as '{temp_name}'): {e}");
723            if let Some(e) = inner.inodes.get_mut(inode) {
724                e.kind = InodeKind::File { handle: new_handle };
725                e.name = temp_name;
726                e.size = size;
727                e.mtime = SystemTime::now();
728            }
729            return Ok(());
730        }
731
732        // Step 3: Rename temp to original name.
733        if let Err(e) = self.rt.block_on(storage.rename(new_handle, &entry.name)) {
734            warn!(
735                "Flush: rename from '{temp_name}' to '{}' failed: {e}",
736                entry.name
737            );
738            if let Some(e) = inner.inodes.get_mut(inode) {
739                e.kind = InodeKind::File { handle: new_handle };
740                e.name = temp_name;
741                e.size = size;
742                e.mtime = SystemTime::now();
743            }
744            return Ok(());
745        }
746
747        if let Some(e) = inner.inodes.get_mut(inode) {
748            e.kind = InodeKind::File { handle: new_handle };
749            e.size = size;
750            e.mtime = SystemTime::now();
751        }
752        Ok(())
753    }
754
755    /// Unsafe flush: delete old object, then upload. Data is lost if upload fails.
756    #[allow(clippy::too_many_arguments)]
757    fn flush_unsafe(
758        &self,
759        inner: &mut Inner,
760        inode: u64,
761        old_handle: ObjectHandle,
762        storage_idx: usize,
763        parent_handle: Option<ObjectHandle>,
764        entry: &InodeEntry,
765        size: u64,
766        file: std::fs::File,
767    ) -> MtpResult<()> {
768        let storage = &inner.storages[storage_idx];
769
770        if let Err(e) = self.rt.block_on(storage.delete(old_handle)) {
771            error!("Flush: failed to delete old object: {e}");
772            return Err(e);
773        }
774
775        let info = NewObjectInfo::file(&entry.name, size);
776        let stream = file_stream(file);
777
778        match self
779            .rt
780            .block_on(storage.upload(parent_handle, info, stream))
781        {
782            Ok(new_handle) => {
783                if let Some(e) = inner.inodes.get_mut(inode) {
784                    e.kind = InodeKind::File { handle: new_handle };
785                    e.size = size;
786                    e.mtime = SystemTime::now();
787                }
788                Ok(())
789            }
790            Err(e) => {
791                error!("Flush: upload failed after delete (data lost): {e}");
792                Err(e.into())
793            }
794        }
795    }
796
797    /// Best-effort removal of a leftover object by name, used to clear a
798    /// half-uploaded temp file before retrying a flush.
799    fn purge_leftover(
800        &self,
801        inner: &mut Inner,
802        storage_idx: usize,
803        parent_handle: Option<ObjectHandle>,
804        name: &str,
805    ) {
806        let storage = &inner.storages[storage_idx];
807        let objects = match self.rt.block_on(storage.list_objects(parent_handle)) {
808            Ok(objects) => objects,
809            Err(e) => {
810                debug!("Flush retry: couldn't list the target directory: {e}");
811                return;
812            }
813        };
814        for obj in objects.into_iter().filter(|o| o.filename == name) {
815            if let Err(e) = self.rt.block_on(storage.delete(obj.handle)) {
816                warn!("Flush retry: couldn't remove leftover '{name}': {e}");
817            }
818        }
819    }
820
821    pub fn mount_options(&self) -> Vec<MountOption> {
822        let mut opts = vec![
823            MountOption::FSName("mtp-mount".to_string()),
824            MountOption::Subtype("mtp".to_string()),
825            MountOption::DefaultPermissions,
826            MountOption::NoDev,
827            MountOption::NoSuid,
828        ];
829        if self.read_only {
830            opts.push(MountOption::RO);
831        } else {
832            opts.push(MountOption::RW);
833        }
834        opts
835    }
836
837    /// Starts the event monitor for the current session.
838    fn spawn_event_loop(&self, device: MtpDevice, epoch: u64) {
839        let inner = Arc::clone(&self.inner);
840        let current_epoch = Arc::clone(&self.event_epoch);
841        self.rt.spawn(async move {
842            Self::event_loop(device, inner, current_epoch, epoch).await;
843        });
844    }
845
846    /// Background event loop that polls the device for MTP events and invalidates
847    /// cached directory listings when objects change on the device side.
848    ///
849    /// Exits when its session is gone: either the device stopped answering, or a
850    /// reconnect moved the mount to a newer session (`epoch`), which starts its
851    /// own loop. A disconnect here doesn't tear the mount down; the next FUSE
852    /// operation is what drives the reconnect.
853    async fn event_loop(
854        device: MtpDevice,
855        inner: Arc<Mutex<Inner>>,
856        current_epoch: Arc<AtomicU64>,
857        epoch: u64,
858    ) {
859        loop {
860            if current_epoch.load(Ordering::SeqCst) != epoch {
861                debug!("Event loop: superseded by a reconnect");
862                return;
863            }
864            match tokio::time::timeout(Duration::from_millis(200), device.next_event()).await {
865                Ok(Ok(event)) => {
866                    Self::handle_event(&inner, &event);
867                }
868                Ok(Err(mtp_rs::Error::Timeout)) => continue,
869                Ok(Err(e)) if is_link_lost(&e) => {
870                    debug!("Event loop: device disconnected");
871                    break;
872                }
873                Ok(Err(e)) => {
874                    warn!("Event loop error: {e}");
875                    break;
876                }
877                Err(_) => continue, // tokio timeout elapsed, loop again
878            }
879        }
880    }
881
882    /// Process a single device event by invalidating the relevant cache entries.
883    fn handle_event(inner: &Mutex<Inner>, event: &DeviceEvent) {
884        match event {
885            DeviceEvent::ObjectAdded { handle } => {
886                debug!("Event: object added {:?}", handle);
887                let mut inner = inner.lock().unwrap();
888                // The new object might be in any directory. If we can find its parent
889                // in the inode table (the parent dir was already cached), invalidate
890                // just that directory. Otherwise, invalidate all directories.
891                if let Some(parent_ino) = inner.inodes.find_parent_by_handle(*handle) {
892                    inner.dirs_loaded.remove(&parent_ino);
893                } else {
894                    Self::invalidate_all_dirs(&mut inner);
895                }
896            }
897            DeviceEvent::ObjectRemoved { handle } => {
898                debug!("Event: object removed {:?}", handle);
899                let mut inner = inner.lock().unwrap();
900                if let Some(parent_ino) = inner.inodes.find_parent_by_handle(*handle) {
901                    inner.dirs_loaded.remove(&parent_ino);
902                } else {
903                    Self::invalidate_all_dirs(&mut inner);
904                }
905            }
906            DeviceEvent::ObjectInfoChanged { handle } => {
907                debug!("Event: object info changed {:?}", handle);
908                let mut inner = inner.lock().unwrap();
909                // Invalidate the parent directory and clear any read cache for this file.
910                if let Some(parent_ino) = inner.inodes.find_parent_by_handle(*handle) {
911                    inner.dirs_loaded.remove(&parent_ino);
912                }
913                // Clear read cache entries for file handles pointing to this object.
914                let fhs_to_clear: Vec<u64> = inner
915                    .fh_to_inode
916                    .iter()
917                    .filter_map(|(&fh, &ino)| {
918                        inner.inodes.get(ino).and_then(|e| match &e.kind {
919                            InodeKind::File { handle: h } if *h == *handle => Some(fh),
920                            _ => None,
921                        })
922                    })
923                    .collect();
924                for fh in fhs_to_clear {
925                    inner.read_cache.remove(&fh);
926                }
927            }
928            DeviceEvent::StoreAdded { .. }
929            | DeviceEvent::StoreRemoved { .. }
930            | DeviceEvent::StorageInfoChanged { .. } => {
931                debug!("Event: storage change {:?}", event);
932                // Storage-level changes: invalidate everything.
933                let mut inner = inner.lock().unwrap();
934                Self::invalidate_all_dirs(&mut inner);
935            }
936            _ => {
937                debug!("Event: unhandled {:?}", event);
938            }
939        }
940    }
941
942    /// Mark all cached directories as stale so they're re-fetched on next access.
943    fn invalidate_all_dirs(inner: &mut Inner) {
944        inner.dirs_loaded.retain(|&k, _| k == FUSE_ROOT_INODE);
945    }
946}
947
948impl Filesystem for MtpFs {
949    fn init(&mut self, _req: &Request, _config: &mut KernelConfig) -> io::Result<()> {
950        let storages = self
951            .rt
952            .block_on(self.device.lock().unwrap().storages())
953            .map_err(|e: mtp_rs::Error| io::Error::other(e.to_string()))?;
954
955        let mut inner = self.inner.lock().unwrap();
956        for storage in &storages {
957            inner
958                .inodes
959                .add_storage(storage.id(), storage_name(storage));
960        }
961        inner.dirs_loaded.insert(FUSE_ROOT_INODE, true);
962        inner.storages = storages;
963        drop(inner);
964
965        // Spawn a background task that monitors device events and invalidates
966        // cached directory listings when objects are added, removed, or changed.
967        let event_device = self.device.lock().unwrap().clone();
968        self.spawn_event_loop(event_device, self.event_epoch.load(Ordering::SeqCst));
969
970        debug!(
971            "MtpFs initialized with {} storages + event monitor",
972            self.inner.lock().unwrap().storages.len()
973        );
974        Ok(())
975    }
976
977    fn lookup(&self, _req: &Request, parent: INodeNo, name: &OsStr, reply: ReplyEntry) {
978        let parent_ino = parent.0;
979        let name_str = match name.to_str() {
980            Some(s) => s,
981            None => {
982                reply.error(Errno::ENOENT);
983                return;
984            }
985        };
986
987        let mut inner = self.inner.lock().unwrap();
988        self.load_dir(&mut inner, parent_ino);
989
990        match inner.inodes.lookup(parent_ino, name_str) {
991            Some(ino) => {
992                let entry = inner.inodes.get(ino).unwrap();
993                let attr = inode_to_file_attr(entry);
994                reply.entry(&TTL, &attr, Generation(0));
995            }
996            None => {
997                reply.error(Errno::ENOENT);
998            }
999        }
1000    }
1001
1002    fn getattr(&self, _req: &Request, ino: INodeNo, _fh: Option<FileHandle>, reply: ReplyAttr) {
1003        let inner = self.inner.lock().unwrap();
1004        match inner.inodes.get(ino.0) {
1005            Some(entry) => {
1006                let mut attr = inode_to_file_attr(entry);
1007                for (&fh, &inode) in &inner.fh_to_inode {
1008                    if inode == ino.0 {
1009                        if let Some(size) = inner.write_buf.size(fh) {
1010                            attr.size = size;
1011                            attr.blocks = size.div_ceil(512);
1012                        }
1013                        break;
1014                    }
1015                }
1016                reply.attr(&TTL, &attr);
1017            }
1018            None => {
1019                reply.error(Errno::ENOENT);
1020            }
1021        }
1022    }
1023
1024    fn readdir(
1025        &self,
1026        _req: &Request,
1027        ino: INodeNo,
1028        _fh: FileHandle,
1029        offset: u64,
1030        mut reply: ReplyDirectory,
1031    ) {
1032        let ino_val = ino.0;
1033
1034        let mut inner = self.inner.lock().unwrap();
1035        self.load_dir(&mut inner, ino_val);
1036
1037        let parent_ino = inner
1038            .inodes
1039            .get(ino_val)
1040            .map(|e| e.parent)
1041            .unwrap_or(FUSE_ROOT_INODE);
1042
1043        let mut entries: Vec<(u64, INodeNo, FileType, String)> = vec![
1044            (1, INodeNo(ino_val), FileType::Directory, ".".to_string()),
1045            (
1046                2,
1047                INodeNo(parent_ino),
1048                FileType::Directory,
1049                "..".to_string(),
1050            ),
1051        ];
1052
1053        let children = inner.inodes.children(ino_val);
1054        for (i, child_ino) in children.iter().enumerate() {
1055            if let Some(child) = inner.inodes.get(*child_ino) {
1056                let kind = if child.is_dir() {
1057                    FileType::Directory
1058                } else {
1059                    FileType::RegularFile
1060                };
1061                entries.push((i as u64 + 3, INodeNo(*child_ino), kind, child.name.clone()));
1062            }
1063        }
1064
1065        for (i, (off, ino, kind, name)) in entries.iter().enumerate() {
1066            if i as u64 >= offset && reply.add(*ino, *off, *kind, name) {
1067                break;
1068            }
1069        }
1070        reply.ok();
1071    }
1072
1073    fn open(&self, _req: &Request, ino: INodeNo, _flags: OpenFlags, reply: ReplyOpen) {
1074        let mut inner = self.inner.lock().unwrap();
1075        match inner.inodes.get(ino.0) {
1076            Some(entry) if !entry.is_dir() => {
1077                let fh = self.alloc_fh();
1078                inner.fh_to_inode.insert(fh, ino.0);
1079                reply.opened(FileHandle(fh), FopenFlags::empty());
1080            }
1081            Some(_) => {
1082                reply.error(Errno::EISDIR);
1083            }
1084            None => {
1085                reply.error(Errno::ENOENT);
1086            }
1087        }
1088    }
1089
1090    fn read(
1091        &self,
1092        _req: &Request,
1093        ino: INodeNo,
1094        fh: FileHandle,
1095        offset: u64,
1096        size: u32,
1097        _flags: OpenFlags,
1098        _lock_owner: Option<LockOwner>,
1099        reply: ReplyData,
1100    ) {
1101        let fh_val = fh.0;
1102        let mut inner = self.inner.lock().unwrap();
1103
1104        // If there's a write buffer open for this fh, read from it.
1105        if inner.write_buf.is_open(fh_val) {
1106            match inner.write_buf.read(fh_val, offset as i64, size) {
1107                Ok(data) => reply.data(&data),
1108                Err(e) => {
1109                    error!("Read from write buffer failed: {e}");
1110                    reply.error(Errno::EIO);
1111                }
1112            }
1113            return;
1114        }
1115
1116        // Resolve the MTP object and its storage.
1117        let entry = match inner.inodes.get(ino.0) {
1118            Some(e) => e.clone(),
1119            None => {
1120                reply.error(Errno::ENOENT);
1121                return;
1122            }
1123        };
1124
1125        if entry.is_dir() {
1126            reply.error(Errno::EISDIR);
1127            return;
1128        }
1129
1130        // Lazily create a sparse cache for this file handle.
1131        use std::collections::hash_map::Entry;
1132        let spool_dir = inner.spool_dir.clone();
1133        if let Entry::Vacant(slot) = inner.read_cache.entry(fh_val) {
1134            let cache = match SparseCache::new(entry.size, &spool_dir) {
1135                Ok(c) => c,
1136                Err(e) => {
1137                    error!("Failed to create sparse cache: {e}");
1138                    reply.error(Errno::EIO);
1139                    return;
1140                }
1141            };
1142            slot.insert(cache);
1143        }
1144
1145        // Figure out which byte ranges still need to be fetched from MTP.
1146        let missing = {
1147            let cache = inner.read_cache.get(&fh_val).unwrap();
1148            cache.missing_ranges(offset, size as u64)
1149        };
1150
1151        // Fetch missing ranges. `read_range` uses the 64-bit partial-read op to
1152        // support offsets beyond 4 GB. Each USB transfer is capped at 1 MB to keep
1153        // latency reasonable.
1154        // Each chunk resolves the object handle again inside `with_recovery`,
1155        // so a read that spans a cable glitch picks up the new session's handle
1156        // and carries on from the byte it stopped at.
1157        const CHUNK: u64 = 1024 * 1024;
1158        for range in missing {
1159            let mut cursor = range.start;
1160            while cursor < range.end {
1161                let chunk_size = (range.end - cursor).min(CHUNK) as u32;
1162                self.fetch_counter.fetch_add(1, Ordering::Relaxed);
1163                let bytes = match self.with_recovery(&mut inner, |fs, inner| {
1164                    let handle = fs.file_handle(inner, ino.0)?;
1165                    let storage_idx =
1166                        Self::find_storage_index(inner, ino.0).ok_or(mtp_rs::Error::NotFound)?;
1167                    fs.rt.block_on(
1168                        inner.storages[storage_idx].read_range(handle, cursor, chunk_size),
1169                    )
1170                }) {
1171                    Ok(b) => b,
1172                    Err(e) => {
1173                        error!("MTP read_range failed at offset {cursor}: {e}");
1174                        reply.error(Errno::EIO);
1175                        return;
1176                    }
1177                };
1178                let bytes_len = bytes.len() as u64;
1179                let cache = inner.read_cache.get_mut(&fh_val).unwrap();
1180                if let Err(e) = cache.write_at(cursor, &bytes) {
1181                    error!("Sparse cache write failed: {e}");
1182                    reply.error(Errno::EIO);
1183                    return;
1184                }
1185                // Short read from device — the object is smaller than reported;
1186                // stop fetching to avoid an infinite loop.
1187                if bytes_len == 0 {
1188                    break;
1189                }
1190                cursor += bytes_len;
1191            }
1192        }
1193
1194        // Serve the requested slice from the cache.
1195        let cache = inner.read_cache.get_mut(&fh_val).unwrap();
1196        match cache.read_at(offset, size as u64) {
1197            Ok(buf) => reply.data(&buf),
1198            Err(e) => {
1199                error!("Sparse cache read failed: {e}");
1200                reply.error(Errno::EIO);
1201            }
1202        }
1203    }
1204
1205    fn release(
1206        &self,
1207        _req: &Request,
1208        _ino: INodeNo,
1209        fh: FileHandle,
1210        _flags: OpenFlags,
1211        _lock_owner: Option<LockOwner>,
1212        _flush: bool,
1213        reply: ReplyEmpty,
1214    ) {
1215        let fh_val = fh.0;
1216        let mut inner = self.inner.lock().unwrap();
1217
1218        let flushed = if inner.write_buf.is_open(fh_val) {
1219            self.flush_to_mtp(&mut inner, fh_val)
1220        } else {
1221            Ok(())
1222        };
1223
1224        inner.read_cache.remove(&fh_val);
1225        inner.fh_to_inode.remove(&fh_val);
1226
1227        // A failed flush means the bytes never reached the device, so `close()`
1228        // has to say so instead of pretending the write worked.
1229        match flushed {
1230            Ok(()) => reply.ok(),
1231            Err(e) => {
1232                error!("Flush on close failed: {e}");
1233                reply.error(Errno::EIO);
1234            }
1235        }
1236    }
1237
1238    fn write(
1239        &self,
1240        _req: &Request,
1241        ino: INodeNo,
1242        fh: FileHandle,
1243        offset: u64,
1244        data: &[u8],
1245        _write_flags: WriteFlags,
1246        _flags: OpenFlags,
1247        _lock_owner: Option<LockOwner>,
1248        reply: ReplyWrite,
1249    ) {
1250        if self.read_only {
1251            reply.error(Errno::EROFS);
1252            return;
1253        }
1254
1255        let fh_val = fh.0;
1256        let mut inner = self.inner.lock().unwrap();
1257
1258        if !inner.write_buf.is_open(fh_val) {
1259            let original_size = inner.inodes.get(ino.0).map(|e| e.size).unwrap_or(0);
1260            if let Err(e) = inner.write_buf.open(fh_val, ino.0, original_size) {
1261                error!("Failed to open write buffer: {e}");
1262                reply.error(Errno::EIO);
1263                return;
1264            }
1265        }
1266
1267        match inner.write_buf.write(fh_val, offset as i64, data) {
1268            Ok(written) => reply.written(written),
1269            Err(e) => {
1270                error!("Write failed: {e}");
1271                reply.error(Errno::EIO);
1272            }
1273        }
1274    }
1275
1276    fn create(
1277        &self,
1278        _req: &Request,
1279        parent: INodeNo,
1280        name: &OsStr,
1281        _mode: u32,
1282        _umask: u32,
1283        _flags: i32,
1284        reply: ReplyCreate,
1285    ) {
1286        if self.read_only {
1287            reply.error(Errno::EROFS);
1288            return;
1289        }
1290
1291        let name_str = match name.to_str() {
1292            Some(s) => s,
1293            None => {
1294                reply.error(Errno::EINVAL);
1295                return;
1296            }
1297        };
1298
1299        let parent_ino = parent.0;
1300        let mut inner = self.inner.lock().unwrap();
1301
1302        if !inner.inodes.get(parent_ino).is_some_and(|e| e.is_dir()) {
1303            reply.error(Errno::ENOTDIR);
1304            return;
1305        }
1306
1307        let handle = match self.with_recovery(&mut inner, |fs, inner| {
1308            let mtp_parent = fs.parent_handle(inner, parent_ino)?;
1309            let storage_idx =
1310                Self::find_storage_index(inner, parent_ino).ok_or(mtp_rs::Error::NotFound)?;
1311            let info = NewObjectInfo::file(name_str, 0);
1312            let stream = bytes_stream(Vec::new());
1313            fs.rt
1314                .block_on(inner.storages[storage_idx].upload(mtp_parent, info, stream))
1315                .map_err(mtp_rs::Error::from)
1316        }) {
1317            Ok(h) => h,
1318            Err(e) => {
1319                error!("MTP create failed: {e}");
1320                reply.error(Errno::EIO);
1321                return;
1322            }
1323        };
1324
1325        let now = SystemTime::now();
1326        let ino = inner
1327            .inodes
1328            .add_object(parent_ino, handle, name_str.to_string(), false, 0, now);
1329
1330        let fh = self.alloc_fh();
1331        inner.fh_to_inode.insert(fh, ino);
1332        if let Err(e) = inner.write_buf.open(fh, ino, 0) {
1333            error!("Failed to open write buffer: {e}");
1334            reply.error(Errno::EIO);
1335            return;
1336        }
1337
1338        let entry = inner.inodes.get(ino).unwrap();
1339        let attr = inode_to_file_attr(entry);
1340        reply.created(
1341            &TTL,
1342            &attr,
1343            Generation(0),
1344            FileHandle(fh),
1345            FopenFlags::empty(),
1346        );
1347    }
1348
1349    fn mkdir(
1350        &self,
1351        _req: &Request,
1352        parent: INodeNo,
1353        name: &OsStr,
1354        _mode: u32,
1355        _umask: u32,
1356        reply: ReplyEntry,
1357    ) {
1358        if self.read_only {
1359            reply.error(Errno::EROFS);
1360            return;
1361        }
1362
1363        let name_str = match name.to_str() {
1364            Some(s) => s,
1365            None => {
1366                reply.error(Errno::EINVAL);
1367                return;
1368            }
1369        };
1370
1371        let parent_ino = parent.0;
1372        let mut inner = self.inner.lock().unwrap();
1373
1374        if !inner.inodes.get(parent_ino).is_some_and(|e| e.is_dir()) {
1375            reply.error(Errno::ENOTDIR);
1376            return;
1377        }
1378
1379        let handle = match self.with_recovery(&mut inner, |fs, inner| {
1380            let mtp_parent = fs.parent_handle(inner, parent_ino)?;
1381            let storage_idx =
1382                Self::find_storage_index(inner, parent_ino).ok_or(mtp_rs::Error::NotFound)?;
1383            fs.rt
1384                .block_on(inner.storages[storage_idx].create_folder(mtp_parent, name_str))
1385        }) {
1386            Ok(h) => h,
1387            Err(e) => {
1388                error!("MTP mkdir failed: {e}");
1389                reply.error(Errno::EIO);
1390                return;
1391            }
1392        };
1393
1394        let now = SystemTime::now();
1395        let ino = inner
1396            .inodes
1397            .add_object(parent_ino, handle, name_str.to_string(), true, 0, now);
1398
1399        let entry = inner.inodes.get(ino).unwrap();
1400        let attr = inode_to_file_attr(entry);
1401        reply.entry(&TTL, &attr, Generation(0));
1402    }
1403
1404    fn unlink(&self, _req: &Request, parent: INodeNo, name: &OsStr, reply: ReplyEmpty) {
1405        if self.read_only {
1406            reply.error(Errno::EROFS);
1407            return;
1408        }
1409
1410        let name_str = match name.to_str() {
1411            Some(s) => s,
1412            None => {
1413                reply.error(Errno::ENOENT);
1414                return;
1415            }
1416        };
1417
1418        let parent_ino = parent.0;
1419        let mut inner = self.inner.lock().unwrap();
1420
1421        let child_ino = match inner.inodes.lookup(parent_ino, name_str) {
1422            Some(i) => i,
1423            None => {
1424                reply.error(Errno::ENOENT);
1425                return;
1426            }
1427        };
1428
1429        if inner.inodes.get(child_ino).is_some_and(|e| e.is_dir()) {
1430            reply.error(Errno::EISDIR);
1431            return;
1432        }
1433
1434        if let Err(e) = self.with_recovery(&mut inner, |fs, inner| {
1435            let handle = fs.file_handle(inner, child_ino)?;
1436            let storage_idx =
1437                Self::find_storage_index(inner, child_ino).ok_or(mtp_rs::Error::NotFound)?;
1438            fs.rt.block_on(inner.storages[storage_idx].delete(handle))
1439        }) {
1440            error!("MTP delete failed: {e}");
1441            reply.error(Errno::EIO);
1442            return;
1443        }
1444
1445        inner.inodes.remove(child_ino);
1446        reply.ok();
1447    }
1448
1449    fn rmdir(&self, _req: &Request, parent: INodeNo, name: &OsStr, reply: ReplyEmpty) {
1450        if self.read_only {
1451            reply.error(Errno::EROFS);
1452            return;
1453        }
1454
1455        let name_str = match name.to_str() {
1456            Some(s) => s,
1457            None => {
1458                reply.error(Errno::ENOENT);
1459                return;
1460            }
1461        };
1462
1463        let parent_ino = parent.0;
1464        let mut inner = self.inner.lock().unwrap();
1465
1466        let child_ino = match inner.inodes.lookup(parent_ino, name_str) {
1467            Some(i) => i,
1468            None => {
1469                reply.error(Errno::ENOENT);
1470                return;
1471            }
1472        };
1473
1474        if !matches!(
1475            inner.inodes.get(child_ino).map(|e| &e.kind),
1476            Some(InodeKind::Directory { .. })
1477        ) {
1478            reply.error(Errno::ENOTDIR);
1479            return;
1480        }
1481
1482        if let Err(e) = self.with_recovery(&mut inner, |fs, inner| {
1483            let handle = fs.object_handle(inner, child_ino)?;
1484            let storage_idx =
1485                Self::find_storage_index(inner, child_ino).ok_or(mtp_rs::Error::NotFound)?;
1486            fs.rt.block_on(inner.storages[storage_idx].delete(handle))
1487        }) {
1488            error!("MTP rmdir failed: {e}");
1489            reply.error(Errno::EIO);
1490            return;
1491        }
1492
1493        inner.inodes.remove(child_ino);
1494        reply.ok();
1495    }
1496
1497    fn rename(
1498        &self,
1499        _req: &Request,
1500        parent: INodeNo,
1501        name: &OsStr,
1502        newparent: INodeNo,
1503        newname: &OsStr,
1504        _flags: RenameFlags,
1505        reply: ReplyEmpty,
1506    ) {
1507        if self.read_only {
1508            reply.error(Errno::EROFS);
1509            return;
1510        }
1511
1512        let name_str = match name.to_str() {
1513            Some(s) => s,
1514            None => {
1515                reply.error(Errno::ENOENT);
1516                return;
1517            }
1518        };
1519        let newname_str = match newname.to_str() {
1520            Some(s) => s,
1521            None => {
1522                reply.error(Errno::EINVAL);
1523                return;
1524            }
1525        };
1526
1527        let parent_ino = parent.0;
1528        let newparent_ino = newparent.0;
1529        let mut inner = self.inner.lock().unwrap();
1530
1531        let child_ino = match inner.inodes.lookup(parent_ino, name_str) {
1532            Some(i) => i,
1533            None => {
1534                reply.error(Errno::ENOENT);
1535                return;
1536            }
1537        };
1538
1539        if !matches!(
1540            inner.inodes.get(child_ino).map(|e| &e.kind),
1541            Some(InodeKind::File { .. } | InodeKind::Directory { .. })
1542        ) {
1543            reply.error(Errno::EINVAL);
1544            return;
1545        }
1546        if parent_ino != newparent_ino
1547            && !inner.inodes.get(newparent_ino).is_some_and(|e| e.is_dir())
1548        {
1549            reply.error(Errno::ENOTDIR);
1550            return;
1551        }
1552
1553        if name_str != newname_str {
1554            if let Err(e) = self.with_recovery(&mut inner, |fs, inner| {
1555                let handle = fs.object_handle(inner, child_ino)?;
1556                let storage_idx =
1557                    Self::find_storage_index(inner, child_ino).ok_or(mtp_rs::Error::NotFound)?;
1558                fs.rt
1559                    .block_on(inner.storages[storage_idx].rename(handle, newname_str))
1560            }) {
1561                error!("MTP rename failed: {e}");
1562                reply.error(Errno::EIO);
1563                return;
1564            }
1565        }
1566
1567        if parent_ino != newparent_ino {
1568            if let Err(e) = self.with_recovery(&mut inner, |fs, inner| {
1569                let handle = fs.object_handle(inner, child_ino)?;
1570                let storage_idx =
1571                    Self::find_storage_index(inner, child_ino).ok_or(mtp_rs::Error::NotFound)?;
1572                let new_mtp_parent = fs
1573                    .parent_handle(inner, newparent_ino)?
1574                    .unwrap_or(ObjectHandle::ROOT);
1575                fs.rt.block_on(inner.storages[storage_idx].move_object(
1576                    handle,
1577                    new_mtp_parent,
1578                    None,
1579                ))
1580            }) {
1581                error!("MTP move failed: {e}");
1582                reply.error(Errno::EIO);
1583                return;
1584            }
1585        }
1586
1587        inner
1588            .inodes
1589            .rename(child_ino, newparent_ino, newname_str.to_string());
1590        reply.ok();
1591    }
1592
1593    fn setattr(
1594        &self,
1595        _req: &Request,
1596        ino: INodeNo,
1597        _mode: Option<u32>,
1598        _uid: Option<u32>,
1599        _gid: Option<u32>,
1600        size: Option<u64>,
1601        _atime: Option<TimeOrNow>,
1602        _mtime: Option<TimeOrNow>,
1603        _ctime: Option<SystemTime>,
1604        fh: Option<FileHandle>,
1605        _crtime: Option<SystemTime>,
1606        _chgtime: Option<SystemTime>,
1607        _bkuptime: Option<SystemTime>,
1608        _flags: Option<BsdFileFlags>,
1609        reply: ReplyAttr,
1610    ) {
1611        if let Some(new_size) = size {
1612            if self.read_only {
1613                reply.error(Errno::EROFS);
1614                return;
1615            }
1616
1617            if let Some(fh) = fh {
1618                let fh_val = fh.0;
1619                let mut inner = self.inner.lock().unwrap();
1620
1621                if !inner.write_buf.is_open(fh_val) {
1622                    let original_size = inner.inodes.get(ino.0).map(|e| e.size).unwrap_or(0);
1623                    if let Err(e) = inner.write_buf.open(fh_val, ino.0, original_size) {
1624                        error!("Failed to open write buffer: {e}");
1625                        reply.error(Errno::EIO);
1626                        return;
1627                    }
1628                }
1629
1630                if new_size == 0 {
1631                    inner.write_buf.close(fh_val);
1632                    if let Err(e) = inner.write_buf.open(fh_val, ino.0, 0) {
1633                        error!("Failed to open write buffer: {e}");
1634                        reply.error(Errno::EIO);
1635                        return;
1636                    }
1637                }
1638            }
1639        }
1640
1641        let inner = self.inner.lock().unwrap();
1642        match inner.inodes.get(ino.0) {
1643            Some(entry) => {
1644                let mut attr = inode_to_file_attr(entry);
1645                if let Some(new_size) = size {
1646                    attr.size = new_size;
1647                    attr.blocks = new_size.div_ceil(512);
1648                }
1649                reply.attr(&TTL, &attr);
1650            }
1651            None => {
1652                reply.error(Errno::ENOENT);
1653            }
1654        }
1655    }
1656
1657    fn statfs(&self, _req: &Request, _ino: INodeNo, reply: ReplyStatfs) {
1658        let inner = self.inner.lock().unwrap();
1659        let block_size: u64 = 4096;
1660
1661        let mut total_bytes: u64 = 0;
1662        let mut free_bytes: u64 = 0;
1663        for storage in &inner.storages {
1664            total_bytes = total_bytes.saturating_add(storage.info().total_capacity);
1665            free_bytes = free_bytes.saturating_add(storage.info().free_space);
1666        }
1667
1668        let blocks = total_bytes / block_size;
1669        let bfree = free_bytes / block_size;
1670
1671        reply.statfs(blocks, bfree, bfree, 0, 0, block_size as u32, 255, 0);
1672    }
1673
1674    fn opendir(&self, _req: &Request, ino: INodeNo, _flags: OpenFlags, reply: ReplyOpen) {
1675        let mut inner = self.inner.lock().unwrap();
1676        match inner.inodes.get(ino.0) {
1677            Some(entry) if entry.is_dir() => {
1678                let fh = self.alloc_fh();
1679                inner.dirs_loaded.remove(&ino.0);
1680                reply.opened(FileHandle(fh), FopenFlags::empty());
1681            }
1682            Some(_) => {
1683                reply.error(Errno::ENOTDIR);
1684            }
1685            None => {
1686                reply.error(Errno::ENOENT);
1687            }
1688        }
1689    }
1690
1691    fn releasedir(
1692        &self,
1693        _req: &Request,
1694        _ino: INodeNo,
1695        _fh: FileHandle,
1696        _flags: OpenFlags,
1697        reply: ReplyEmpty,
1698    ) {
1699        reply.ok();
1700    }
1701}
1702
1703#[cfg(test)]
1704mod tests {
1705    use super::*;
1706    use futures::StreamExt as _;
1707    use std::io::Write as _;
1708
1709    const CHUNK: usize = UPLOAD_CHUNK;
1710
1711    /// A temp file holding `content`, rewound to the start.
1712    fn spool_file(content: &[u8]) -> std::fs::File {
1713        let mut file = tempfile::tempfile().unwrap();
1714        file.write_all(content).unwrap();
1715        file.seek(SeekFrom::Start(0)).unwrap();
1716        file
1717    }
1718
1719    #[test]
1720    fn file_stream_reads_lazily() {
1721        let file = spool_file(&vec![0xABu8; CHUNK * 4]);
1722        // `try_clone` dups the fd, so the clone shares the read cursor and shows
1723        // how far the stream has actually read.
1724        let mut cursor = file.try_clone().unwrap();
1725
1726        let mut stream = file_stream(file);
1727        let first = futures::executor::block_on(stream.next()).unwrap().unwrap();
1728
1729        assert_eq!(first.len(), CHUNK);
1730        assert_eq!(
1731            cursor.stream_position().unwrap(),
1732            CHUNK as u64,
1733            "one poll must read one chunk, not the whole file"
1734        );
1735    }
1736
1737    #[test]
1738    fn file_stream_yields_the_whole_file_then_ends() {
1739        // Two full chunks plus a short tail, so the partial-read path is covered.
1740        let content = vec![0xCDu8; CHUNK * 2 + 17];
1741
1742        let chunks: Vec<_> =
1743            futures::executor::block_on(file_stream(spool_file(&content)).collect());
1744
1745        let sizes: Vec<_> = chunks.iter().map(|c| c.as_ref().unwrap().len()).collect();
1746        assert_eq!(sizes, vec![CHUNK, CHUNK, 17]);
1747        let joined: Vec<u8> = chunks
1748            .into_iter()
1749            .flat_map(|c| c.unwrap().to_vec())
1750            .collect();
1751        assert_eq!(joined, content);
1752    }
1753
1754    #[test]
1755    fn file_stream_ends_after_a_read_error() {
1756        // A write-only fd fails every read with EBADF, and the stream has to stop
1757        // there rather than re-emitting the error forever.
1758        let path = tempfile::NamedTempFile::new().unwrap().into_temp_path();
1759        let file = std::fs::OpenOptions::new().write(true).open(&path).unwrap();
1760
1761        let mut stream = file_stream(file);
1762        assert!(futures::executor::block_on(stream.next()).unwrap().is_err());
1763        assert!(futures::executor::block_on(stream.next()).is_none());
1764    }
1765}