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
31const UPLOAD_CHUNK: usize = 65536;
33
34const MAX_ATTEMPTS: u32 = 3;
38
39const MAX_STALE_RETRIES: u32 = 1;
54
55type MtpResult<T> = Result<T, mtp_rs::Error>;
56
57pub struct MtpFsConfig {
59 pub read_only: bool,
60 pub spool_dir: PathBuf,
64 pub reconnect: ReconnectPolicy,
66 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
116fn 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
125fn temp_upload_name(name: &str) -> String {
127 format!(".~tmp~{name}")
128}
129
130fn io_error(e: io::Error) -> mtp_rs::Error {
133 mtp_rs::Error::Io {
134 message: e.to_string(),
135 }
136}
137
138fn 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
150fn 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 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
178struct Inner {
180 storages: Vec<Storage>,
181 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
190pub struct MtpFs {
192 rt: tokio::runtime::Handle,
193 device: Mutex<MtpDevice>,
194 opener: Arc<dyn DeviceOpener>,
196 policy: ReconnectPolicy,
197 unplug: UnplugSwitch,
198 shutdown: Arc<Shutdown>,
201 event_epoch: Arc<AtomicU64>,
204 inner: Arc<Mutex<Inner>>,
205 next_fh: AtomicU64,
206 read_only: bool,
207 fetch_counter: Arc<AtomicU64>,
210}
211
212impl MtpFs {
213 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 pub fn shutdown(&self) -> Arc<Shutdown> {
257 Arc::clone(&self.shutdown)
258 }
259
260 #[allow(dead_code)] 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 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 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 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 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 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 fn adopt(&self, inner: &mut Inner, device: MtpDevice) -> MtpResult<()> {
404 let storages = self.rt.block_on(device.storages())?;
405
406 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 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 fn give_up(&self, reason: &str) {
442 error!("{reason}; unmounting");
443 eprintln!("mtp-mount: {reason}. Unmounting.");
444 self.shutdown.request(reason);
445 }
446
447 fn ensure_fresh(&self, inner: &mut Inner, inode: u64) -> MtpResult<()> {
451 if inner.inodes.is_fresh(inode) {
452 return Ok(());
453 }
454
455 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 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 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 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 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 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 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 #[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 #[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 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 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 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 #[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 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 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 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, }
879 }
880 }
881
882 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 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 if let Some(parent_ino) = inner.inodes.find_parent_by_handle(*handle) {
911 inner.dirs_loaded.remove(&parent_ino);
912 }
913 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 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 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 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 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 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 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 let missing = {
1147 let cache = inner.read_cache.get(&fh_val).unwrap();
1148 cache.missing_ranges(offset, size as u64)
1149 };
1150
1151 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 if bytes_len == 0 {
1188 break;
1189 }
1190 cursor += bytes_len;
1191 }
1192 }
1193
1194 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 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 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 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 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 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}