use std::collections::HashMap;
use std::ffi::OsStr;
use std::io;
use std::io::{Seek, SeekFrom};
use std::path::PathBuf;
use std::pin::Pin;
use std::sync::atomic::{AtomicU64, Ordering};
use std::sync::{Arc, Mutex};
use std::time::{Duration, SystemTime, UNIX_EPOCH};
use bytes::Bytes;
use fuser::{
BsdFileFlags, Errno, FileAttr, FileHandle, FileType, Filesystem, FopenFlags, Generation,
INodeNo, KernelConfig, LockOwner, MountOption, OpenFlags, RenameFlags, ReplyAttr, ReplyCreate,
ReplyData, ReplyDirectory, ReplyEmpty, ReplyEntry, ReplyOpen, ReplyStatfs, ReplyWrite, Request,
TimeOrNow, WriteFlags,
};
use log::{debug, error, info, warn};
use mtp_rs::mtp::{DeviceEvent, MtpDevice};
use mtp_rs::{NewObjectInfo, ObjectHandle, Storage};
use crate::buffer::WriteBuffer;
use crate::device::{is_link_lost, DeviceOpener, UnplugSwitch};
use crate::inode::{ChildInfo, InodeEntry, InodeKind, InodeTable, FUSE_ROOT_INODE};
use crate::reconnect::ReconnectPolicy;
use crate::shutdown::Shutdown;
use crate::sparse_cache::SparseCache;
const TTL: Duration = Duration::from_secs(1);
const UPLOAD_CHUNK: usize = 65536;
const MAX_ATTEMPTS: u32 = 3;
const MAX_STALE_RETRIES: u32 = 1;
type MtpResult<T> = Result<T, mtp_rs::Error>;
pub struct MtpFsConfig {
pub read_only: bool,
pub spool_dir: PathBuf,
pub reconnect: ReconnectPolicy,
pub unplug: UnplugSwitch,
}
fn mtp_datetime_to_system_time(dt: &mtp_rs::DateTime) -> SystemTime {
fn days_from_civil(y: i64, m: i64, d: i64) -> i64 {
let y = if m <= 2 { y - 1 } else { y };
let era = if y >= 0 { y } else { y - 399 } / 400;
let yoe = (y - era * 400) as u64;
let m_adj = if m > 2 { m - 3 } else { m + 9 } as u64;
let doy = (153 * m_adj + 2) / 5 + d as u64 - 1;
let doe = yoe * 365 + yoe / 4 - yoe / 100 + doy;
era * 146097 + doe as i64 - 719468
}
let days = days_from_civil(dt.year as i64, dt.month as i64, dt.day as i64);
let secs = days * 86400 + dt.hour as i64 * 3600 + dt.minute as i64 * 60 + dt.second as i64;
if secs >= 0 {
UNIX_EPOCH + Duration::from_secs(secs as u64)
} else {
UNIX_EPOCH
}
}
fn inode_to_file_attr(entry: &InodeEntry) -> FileAttr {
let uid = unsafe { libc::getuid() };
let gid = unsafe { libc::getgid() };
FileAttr {
ino: INodeNo(entry.inode),
size: entry.size,
blocks: entry.size.div_ceil(512),
atime: entry.atime,
mtime: entry.mtime,
ctime: entry.mtime,
crtime: entry.mtime,
kind: if entry.is_dir() {
FileType::Directory
} else {
FileType::RegularFile
},
perm: if entry.is_dir() { 0o755 } else { 0o644 },
nlink: if entry.is_dir() { 2 } else { 1 },
uid,
gid,
rdev: 0,
blksize: 4096,
flags: 0,
}
}
fn storage_name(storage: &Storage) -> String {
if storage.info().description.is_empty() {
format!("Storage_{}", storage.id().0)
} else {
storage.info().description.clone()
}
}
fn temp_upload_name(name: &str) -> String {
format!(".~tmp~{name}")
}
fn io_error(e: io::Error) -> mtp_rs::Error {
mtp_rs::Error::Io {
message: e.to_string(),
}
}
fn bytes_stream(
data: Vec<u8>,
) -> futures::stream::Iter<std::vec::IntoIter<Result<Bytes, io::Error>>> {
let chunks = if data.is_empty() {
vec![Ok(Bytes::new())]
} else {
vec![Ok(Bytes::from(data))]
};
futures::stream::iter(chunks)
}
fn file_stream(
file: std::fs::File,
) -> Pin<Box<dyn futures::Stream<Item = Result<Bytes, io::Error>> + Send>> {
use std::io::Read as _;
Box::pin(futures::stream::unfold(Some(file), |state| async move {
let mut file = state?;
let mut buf = vec![0u8; UPLOAD_CHUNK];
match file.read(&mut buf) {
Ok(0) => None,
Ok(n) => {
buf.truncate(n);
Some((Ok(Bytes::from(buf)), Some(file)))
}
Err(e) => Some((Err(e), None)),
}
}))
}
struct Inner {
storages: Vec<Storage>,
spool_dir: PathBuf,
inodes: InodeTable,
write_buf: WriteBuffer,
read_cache: HashMap<u64, SparseCache>,
dirs_loaded: HashMap<u64, bool>,
fh_to_inode: HashMap<u64, u64>,
}
pub struct MtpFs {
rt: tokio::runtime::Handle,
device: Mutex<MtpDevice>,
opener: Arc<dyn DeviceOpener>,
policy: ReconnectPolicy,
unplug: UnplugSwitch,
shutdown: Arc<Shutdown>,
event_epoch: Arc<AtomicU64>,
inner: Arc<Mutex<Inner>>,
next_fh: AtomicU64,
read_only: bool,
fetch_counter: Arc<AtomicU64>,
}
impl MtpFs {
pub fn new(
device: MtpDevice,
opener: Arc<dyn DeviceOpener>,
rt: tokio::runtime::Handle,
config: MtpFsConfig,
) -> Self {
let MtpFsConfig {
read_only,
spool_dir,
reconnect,
unplug,
} = config;
Self {
rt,
device: Mutex::new(device),
opener,
policy: reconnect,
unplug,
shutdown: Arc::new(Shutdown::default()),
event_epoch: Arc::new(AtomicU64::new(0)),
inner: Arc::new(Mutex::new(Inner {
storages: Vec::new(),
spool_dir: spool_dir.clone(),
inodes: InodeTable::new(),
write_buf: WriteBuffer::new(spool_dir),
read_cache: HashMap::new(),
dirs_loaded: HashMap::new(),
fh_to_inode: HashMap::new(),
})),
next_fh: AtomicU64::new(1),
read_only,
fetch_counter: Arc::new(AtomicU64::new(0)),
}
}
pub fn shutdown(&self) -> Arc<Shutdown> {
Arc::clone(&self.shutdown)
}
#[allow(dead_code)] pub fn fetch_counter(&self) -> Arc<AtomicU64> {
Arc::clone(&self.fetch_counter)
}
fn alloc_fh(&self) -> u64 {
self.next_fh.fetch_add(1, Ordering::Relaxed)
}
fn find_storage_index(inner: &Inner, inode: u64) -> Option<usize> {
let mut current = inode;
loop {
let entry = inner.inodes.get(current)?;
if let InodeKind::Storage { storage_id } = &entry.kind {
return inner
.storages
.iter()
.position(|s: &Storage| s.id() == *storage_id);
}
if current == entry.parent {
return None;
}
current = entry.parent;
}
}
fn with_recovery<T>(
&self,
inner: &mut Inner,
mut attempt: impl FnMut(&Self, &mut Inner) -> MtpResult<T>,
) -> MtpResult<T> {
for _ in 0..MAX_ATTEMPTS {
let mut stale_retries = MAX_STALE_RETRIES;
loop {
if self.unplug.is_unplugged() {
break;
}
match attempt(self, inner) {
Ok(value) => return Ok(value),
Err(e) if e.is_stale_handle() => {
if stale_retries == 0 {
error!("Operation still hit a stale object handle after re-resolving");
return Err(e);
}
stale_retries -= 1;
debug!("Operation hit a stale object handle, re-resolving by path");
Self::invalidate_handles(inner);
}
Err(e) if !is_link_lost(&e) => return Err(e),
Err(e) => {
debug!("Operation hit a dead session: {e}");
break;
}
}
}
self.reconnect(inner)?;
}
Err(mtp_rs::Error::Disconnected)
}
fn invalidate_handles(inner: &mut Inner) {
inner.inodes.bump_generation();
inner.dirs_loaded.clear();
inner.dirs_loaded.insert(FUSE_ROOT_INODE, true);
}
fn reconnect(&self, inner: &mut Inner) -> MtpResult<()> {
let who = self.opener.describe();
if self.policy.is_disabled() {
self.give_up(&format!(
"{who} disconnected and reconnect is off (--reconnect-timeout 0)"
));
return Err(mtp_rs::Error::Disconnected);
}
let secs = self.policy.timeout().as_secs();
info!("{who} disconnected, waiting up to {secs}s for it to come back...");
eprintln!("mtp-mount: {who} disconnected, waiting up to {secs}s for it to come back...");
for delay in self.policy.schedule() {
std::thread::sleep(delay);
if self.unplug.is_unplugged() {
continue;
}
let device = match self.opener.open(&self.rt) {
Ok(device) => device,
Err(e) => {
debug!("Reconnect attempt failed: {e}");
continue;
}
};
match self.adopt(inner, device) {
Ok(()) => {
eprintln!("mtp-mount: {who} is back, carrying on.");
return Ok(());
}
Err(e) => {
warn!("Reopened {who} but couldn't resume the mount: {e}");
continue;
}
}
}
self.give_up(&format!("{who} didn't come back within {secs}s"));
Err(mtp_rs::Error::Disconnected)
}
fn adopt(&self, inner: &mut Inner, device: MtpDevice) -> MtpResult<()> {
let storages = self.rt.block_on(device.storages())?;
let storage_inodes = inner.inodes.children(FUSE_ROOT_INODE);
for (position, storage_ino) in storage_inodes.iter().enumerate() {
let name = match inner.inodes.get(*storage_ino) {
Some(entry) => entry.name.clone(),
None => continue,
};
let matched = storages
.iter()
.find(|s| storage_name(s) == name)
.or_else(|| storages.get(position));
match matched {
Some(storage) => inner.inodes.set_storage_id(*storage_ino, storage.id()),
None => warn!("Storage '{name}' is missing after the reconnect"),
}
}
inner.storages = storages;
*self.device.lock().unwrap() = device.clone();
inner.inodes.bump_generation();
inner.dirs_loaded.clear();
inner.dirs_loaded.insert(FUSE_ROOT_INODE, true);
let epoch = self.event_epoch.fetch_add(1, Ordering::SeqCst) + 1;
self.spawn_event_loop(device, epoch);
Ok(())
}
fn give_up(&self, reason: &str) {
error!("{reason}; unmounting");
eprintln!("mtp-mount: {reason}. Unmounting.");
self.shutdown.request(reason);
}
fn ensure_fresh(&self, inner: &mut Inner, inode: u64) -> MtpResult<()> {
if inner.inodes.is_fresh(inode) {
return Ok(());
}
let mut chain = Vec::new();
let mut current = inode;
loop {
let entry = inner.inodes.get(current).ok_or(mtp_rs::Error::NotFound)?;
match entry.kind {
InodeKind::Root | InodeKind::Storage { .. } => break,
_ => {
chain.push(current);
if current == entry.parent {
return Err(mtp_rs::Error::NotFound);
}
current = entry.parent;
}
}
}
chain.reverse();
let storage_idx = Self::find_storage_index(inner, inode).ok_or(mtp_rs::Error::NotFound)?;
let mut parent_handle: Option<ObjectHandle> = None;
for node in chain {
let entry = inner.inodes.get(node).ok_or(mtp_rs::Error::NotFound)?;
let name = entry.name.clone();
if inner.inodes.is_fresh(node) {
parent_handle = match inner.inodes.get(node).map(|e| &e.kind) {
Some(InodeKind::Directory { handle } | InodeKind::File { handle }) => {
Some(*handle)
}
_ => return Err(mtp_rs::Error::NotFound),
};
continue;
}
let objects = self
.rt
.block_on(inner.storages[storage_idx].list_objects(parent_handle))?;
let found = objects
.into_iter()
.find(|obj| obj.filename == name)
.ok_or(mtp_rs::Error::NotFound)?;
inner.inodes.set_handle(node, found.handle);
parent_handle = Some(found.handle);
}
Ok(())
}
fn file_handle(&self, inner: &mut Inner, inode: u64) -> MtpResult<ObjectHandle> {
self.ensure_fresh(inner, inode)?;
match inner.inodes.get(inode).map(|e| &e.kind) {
Some(InodeKind::File { handle }) => Ok(*handle),
_ => Err(mtp_rs::Error::NotFound),
}
}
fn object_handle(&self, inner: &mut Inner, inode: u64) -> MtpResult<ObjectHandle> {
self.ensure_fresh(inner, inode)?;
match inner.inodes.get(inode).map(|e| &e.kind) {
Some(InodeKind::File { handle } | InodeKind::Directory { handle }) => Ok(*handle),
_ => Err(mtp_rs::Error::NotFound),
}
}
fn parent_handle(&self, inner: &mut Inner, inode: u64) -> MtpResult<Option<ObjectHandle>> {
self.ensure_fresh(inner, inode)?;
match inner.inodes.get(inode).map(|e| &e.kind) {
Some(InodeKind::Storage { .. }) => Ok(None),
Some(InodeKind::Directory { handle }) => Ok(Some(*handle)),
_ => Err(mtp_rs::Error::NotFound),
}
}
fn load_dir(&self, inner: &mut Inner, parent_inode: u64) {
if inner.dirs_loaded.get(&parent_inode) == Some(&true) {
return;
}
if parent_inode == FUSE_ROOT_INODE {
inner.dirs_loaded.insert(parent_inode, true);
return;
}
match self.with_recovery(inner, |fs, inner| fs.list_into_table(inner, parent_inode)) {
Ok(()) => {
inner.dirs_loaded.insert(parent_inode, true);
}
Err(e) => error!("Failed to list MTP objects: {e}"),
}
}
fn list_into_table(&self, inner: &mut Inner, parent_inode: u64) -> MtpResult<()> {
let mtp_parent = self.parent_handle(inner, parent_inode)?;
let storage_idx =
Self::find_storage_index(inner, parent_inode).ok_or(mtp_rs::Error::NotFound)?;
let objects = self
.rt
.block_on(inner.storages[storage_idx].list_objects(mtp_parent))?;
let children: Vec<ChildInfo> = objects
.into_iter()
.map(|obj| ChildInfo {
handle: obj.handle,
is_dir: obj.is_folder(),
size: obj.size,
mtime: obj
.modified
.as_ref()
.map(mtp_datetime_to_system_time)
.unwrap_or(UNIX_EPOCH),
name: obj.filename,
})
.collect();
inner.inodes.sync_children(parent_inode, &children);
Ok(())
}
fn flush_to_mtp(&self, inner: &mut Inner, fh: u64) -> MtpResult<()> {
let buf = match inner.write_buf.close(fh) {
Some(b) => b,
None => return Ok(()),
};
if !buf.is_dirty() {
return Ok(());
}
let inode = buf.inode;
let mut file = buf.into_file();
let file_len = file.seek(SeekFrom::End(0)).map_err(io_error)?;
let entry = match inner.inodes.get(inode) {
Some(e) => e.clone(),
None => {
error!("Flush: inode {inode} not found");
return Err(mtp_rs::Error::NotFound);
}
};
let mut attempts = 0u32;
self.with_recovery(inner, |fs, inner| {
let mut attempt = file.try_clone().map_err(io_error)?;
attempt.seek(SeekFrom::Start(0)).map_err(io_error)?;
attempts += 1;
fs.flush_once(inner, inode, &entry, file_len, attempt, attempts > 1)
})
}
#[allow(clippy::too_many_arguments)]
fn flush_once(
&self,
inner: &mut Inner,
inode: u64,
entry: &InodeEntry,
size: u64,
file: std::fs::File,
is_retry: bool,
) -> MtpResult<()> {
let handle = self.file_handle(inner, inode)?;
let storage_idx = Self::find_storage_index(inner, inode).ok_or(mtp_rs::Error::NotFound)?;
let parent_handle = match self.parent_handle(inner, entry.parent) {
Ok(handle) => handle,
Err(e) => {
error!("Flush: no parent directory for inode {inode}: {e}");
return Err(e);
}
};
let supports_rename = self.device.lock().unwrap().supports_rename();
if is_retry && supports_rename {
self.purge_leftover(
inner,
storage_idx,
parent_handle,
&temp_upload_name(&entry.name),
);
}
if supports_rename {
self.flush_safe(
inner,
inode,
handle,
storage_idx,
parent_handle,
entry,
size,
file,
)
} else {
warn!(
"Flush: device does not support rename, using delete-then-upload \
(data loss possible if upload fails)"
);
self.flush_unsafe(
inner,
inode,
handle,
storage_idx,
parent_handle,
entry,
size,
file,
)
}
}
#[allow(clippy::too_many_arguments)]
fn flush_safe(
&self,
inner: &mut Inner,
inode: u64,
old_handle: ObjectHandle,
storage_idx: usize,
parent_handle: Option<ObjectHandle>,
entry: &InodeEntry,
size: u64,
file: std::fs::File,
) -> MtpResult<()> {
let storage = &inner.storages[storage_idx];
let temp_name = temp_upload_name(&entry.name);
let info = NewObjectInfo::file(&temp_name, size);
let stream = file_stream(file);
let new_handle = match self
.rt
.block_on(storage.upload(parent_handle, info, stream))
{
Ok(h) => h,
Err(e) => {
error!("Flush: upload failed (original file untouched): {e}");
return Err(e.into());
}
};
if let Err(e) = self.rt.block_on(storage.delete(old_handle)) {
error!("Flush: failed to delete old object (new data saved as '{temp_name}'): {e}");
if let Some(e) = inner.inodes.get_mut(inode) {
e.kind = InodeKind::File { handle: new_handle };
e.name = temp_name;
e.size = size;
e.mtime = SystemTime::now();
}
return Ok(());
}
if let Err(e) = self.rt.block_on(storage.rename(new_handle, &entry.name)) {
warn!(
"Flush: rename from '{temp_name}' to '{}' failed: {e}",
entry.name
);
if let Some(e) = inner.inodes.get_mut(inode) {
e.kind = InodeKind::File { handle: new_handle };
e.name = temp_name;
e.size = size;
e.mtime = SystemTime::now();
}
return Ok(());
}
if let Some(e) = inner.inodes.get_mut(inode) {
e.kind = InodeKind::File { handle: new_handle };
e.size = size;
e.mtime = SystemTime::now();
}
Ok(())
}
#[allow(clippy::too_many_arguments)]
fn flush_unsafe(
&self,
inner: &mut Inner,
inode: u64,
old_handle: ObjectHandle,
storage_idx: usize,
parent_handle: Option<ObjectHandle>,
entry: &InodeEntry,
size: u64,
file: std::fs::File,
) -> MtpResult<()> {
let storage = &inner.storages[storage_idx];
if let Err(e) = self.rt.block_on(storage.delete(old_handle)) {
error!("Flush: failed to delete old object: {e}");
return Err(e);
}
let info = NewObjectInfo::file(&entry.name, size);
let stream = file_stream(file);
match self
.rt
.block_on(storage.upload(parent_handle, info, stream))
{
Ok(new_handle) => {
if let Some(e) = inner.inodes.get_mut(inode) {
e.kind = InodeKind::File { handle: new_handle };
e.size = size;
e.mtime = SystemTime::now();
}
Ok(())
}
Err(e) => {
error!("Flush: upload failed after delete (data lost): {e}");
Err(e.into())
}
}
}
fn purge_leftover(
&self,
inner: &mut Inner,
storage_idx: usize,
parent_handle: Option<ObjectHandle>,
name: &str,
) {
let storage = &inner.storages[storage_idx];
let objects = match self.rt.block_on(storage.list_objects(parent_handle)) {
Ok(objects) => objects,
Err(e) => {
debug!("Flush retry: couldn't list the target directory: {e}");
return;
}
};
for obj in objects.into_iter().filter(|o| o.filename == name) {
if let Err(e) = self.rt.block_on(storage.delete(obj.handle)) {
warn!("Flush retry: couldn't remove leftover '{name}': {e}");
}
}
}
pub fn mount_options(&self) -> Vec<MountOption> {
let mut opts = vec![
MountOption::FSName("mtp-mount".to_string()),
MountOption::Subtype("mtp".to_string()),
MountOption::DefaultPermissions,
MountOption::NoDev,
MountOption::NoSuid,
];
if self.read_only {
opts.push(MountOption::RO);
} else {
opts.push(MountOption::RW);
}
opts
}
fn spawn_event_loop(&self, device: MtpDevice, epoch: u64) {
let inner = Arc::clone(&self.inner);
let current_epoch = Arc::clone(&self.event_epoch);
self.rt.spawn(async move {
Self::event_loop(device, inner, current_epoch, epoch).await;
});
}
async fn event_loop(
device: MtpDevice,
inner: Arc<Mutex<Inner>>,
current_epoch: Arc<AtomicU64>,
epoch: u64,
) {
loop {
if current_epoch.load(Ordering::SeqCst) != epoch {
debug!("Event loop: superseded by a reconnect");
return;
}
match tokio::time::timeout(Duration::from_millis(200), device.next_event()).await {
Ok(Ok(event)) => {
Self::handle_event(&inner, &event);
}
Ok(Err(mtp_rs::Error::Timeout)) => continue,
Ok(Err(e)) if is_link_lost(&e) => {
debug!("Event loop: device disconnected");
break;
}
Ok(Err(e)) => {
warn!("Event loop error: {e}");
break;
}
Err(_) => continue, }
}
}
fn handle_event(inner: &Mutex<Inner>, event: &DeviceEvent) {
match event {
DeviceEvent::ObjectAdded { handle } => {
debug!("Event: object added {:?}", handle);
let mut inner = inner.lock().unwrap();
if let Some(parent_ino) = inner.inodes.find_parent_by_handle(*handle) {
inner.dirs_loaded.remove(&parent_ino);
} else {
Self::invalidate_all_dirs(&mut inner);
}
}
DeviceEvent::ObjectRemoved { handle } => {
debug!("Event: object removed {:?}", handle);
let mut inner = inner.lock().unwrap();
if let Some(parent_ino) = inner.inodes.find_parent_by_handle(*handle) {
inner.dirs_loaded.remove(&parent_ino);
} else {
Self::invalidate_all_dirs(&mut inner);
}
}
DeviceEvent::ObjectInfoChanged { handle } => {
debug!("Event: object info changed {:?}", handle);
let mut inner = inner.lock().unwrap();
if let Some(parent_ino) = inner.inodes.find_parent_by_handle(*handle) {
inner.dirs_loaded.remove(&parent_ino);
}
let fhs_to_clear: Vec<u64> = inner
.fh_to_inode
.iter()
.filter_map(|(&fh, &ino)| {
inner.inodes.get(ino).and_then(|e| match &e.kind {
InodeKind::File { handle: h } if *h == *handle => Some(fh),
_ => None,
})
})
.collect();
for fh in fhs_to_clear {
inner.read_cache.remove(&fh);
}
}
DeviceEvent::StoreAdded { .. }
| DeviceEvent::StoreRemoved { .. }
| DeviceEvent::StorageInfoChanged { .. } => {
debug!("Event: storage change {:?}", event);
let mut inner = inner.lock().unwrap();
Self::invalidate_all_dirs(&mut inner);
}
_ => {
debug!("Event: unhandled {:?}", event);
}
}
}
fn invalidate_all_dirs(inner: &mut Inner) {
inner.dirs_loaded.retain(|&k, _| k == FUSE_ROOT_INODE);
}
}
impl Filesystem for MtpFs {
fn init(&mut self, _req: &Request, _config: &mut KernelConfig) -> io::Result<()> {
let storages = self
.rt
.block_on(self.device.lock().unwrap().storages())
.map_err(|e: mtp_rs::Error| io::Error::other(e.to_string()))?;
let mut inner = self.inner.lock().unwrap();
for storage in &storages {
inner
.inodes
.add_storage(storage.id(), storage_name(storage));
}
inner.dirs_loaded.insert(FUSE_ROOT_INODE, true);
inner.storages = storages;
drop(inner);
let event_device = self.device.lock().unwrap().clone();
self.spawn_event_loop(event_device, self.event_epoch.load(Ordering::SeqCst));
debug!(
"MtpFs initialized with {} storages + event monitor",
self.inner.lock().unwrap().storages.len()
);
Ok(())
}
fn lookup(&self, _req: &Request, parent: INodeNo, name: &OsStr, reply: ReplyEntry) {
let parent_ino = parent.0;
let name_str = match name.to_str() {
Some(s) => s,
None => {
reply.error(Errno::ENOENT);
return;
}
};
let mut inner = self.inner.lock().unwrap();
self.load_dir(&mut inner, parent_ino);
match inner.inodes.lookup(parent_ino, name_str) {
Some(ino) => {
let entry = inner.inodes.get(ino).unwrap();
let attr = inode_to_file_attr(entry);
reply.entry(&TTL, &attr, Generation(0));
}
None => {
reply.error(Errno::ENOENT);
}
}
}
fn getattr(&self, _req: &Request, ino: INodeNo, _fh: Option<FileHandle>, reply: ReplyAttr) {
let inner = self.inner.lock().unwrap();
match inner.inodes.get(ino.0) {
Some(entry) => {
let mut attr = inode_to_file_attr(entry);
for (&fh, &inode) in &inner.fh_to_inode {
if inode == ino.0 {
if let Some(size) = inner.write_buf.size(fh) {
attr.size = size;
attr.blocks = size.div_ceil(512);
}
break;
}
}
reply.attr(&TTL, &attr);
}
None => {
reply.error(Errno::ENOENT);
}
}
}
fn readdir(
&self,
_req: &Request,
ino: INodeNo,
_fh: FileHandle,
offset: u64,
mut reply: ReplyDirectory,
) {
let ino_val = ino.0;
let mut inner = self.inner.lock().unwrap();
self.load_dir(&mut inner, ino_val);
let parent_ino = inner
.inodes
.get(ino_val)
.map(|e| e.parent)
.unwrap_or(FUSE_ROOT_INODE);
let mut entries: Vec<(u64, INodeNo, FileType, String)> = vec![
(1, INodeNo(ino_val), FileType::Directory, ".".to_string()),
(
2,
INodeNo(parent_ino),
FileType::Directory,
"..".to_string(),
),
];
let children = inner.inodes.children(ino_val);
for (i, child_ino) in children.iter().enumerate() {
if let Some(child) = inner.inodes.get(*child_ino) {
let kind = if child.is_dir() {
FileType::Directory
} else {
FileType::RegularFile
};
entries.push((i as u64 + 3, INodeNo(*child_ino), kind, child.name.clone()));
}
}
for (i, (off, ino, kind, name)) in entries.iter().enumerate() {
if i as u64 >= offset && reply.add(*ino, *off, *kind, name) {
break;
}
}
reply.ok();
}
fn open(&self, _req: &Request, ino: INodeNo, _flags: OpenFlags, reply: ReplyOpen) {
let mut inner = self.inner.lock().unwrap();
match inner.inodes.get(ino.0) {
Some(entry) if !entry.is_dir() => {
let fh = self.alloc_fh();
inner.fh_to_inode.insert(fh, ino.0);
reply.opened(FileHandle(fh), FopenFlags::empty());
}
Some(_) => {
reply.error(Errno::EISDIR);
}
None => {
reply.error(Errno::ENOENT);
}
}
}
fn read(
&self,
_req: &Request,
ino: INodeNo,
fh: FileHandle,
offset: u64,
size: u32,
_flags: OpenFlags,
_lock_owner: Option<LockOwner>,
reply: ReplyData,
) {
let fh_val = fh.0;
let mut inner = self.inner.lock().unwrap();
if inner.write_buf.is_open(fh_val) {
match inner.write_buf.read(fh_val, offset as i64, size) {
Ok(data) => reply.data(&data),
Err(e) => {
error!("Read from write buffer failed: {e}");
reply.error(Errno::EIO);
}
}
return;
}
let entry = match inner.inodes.get(ino.0) {
Some(e) => e.clone(),
None => {
reply.error(Errno::ENOENT);
return;
}
};
if entry.is_dir() {
reply.error(Errno::EISDIR);
return;
}
use std::collections::hash_map::Entry;
let spool_dir = inner.spool_dir.clone();
if let Entry::Vacant(slot) = inner.read_cache.entry(fh_val) {
let cache = match SparseCache::new(entry.size, &spool_dir) {
Ok(c) => c,
Err(e) => {
error!("Failed to create sparse cache: {e}");
reply.error(Errno::EIO);
return;
}
};
slot.insert(cache);
}
let missing = {
let cache = inner.read_cache.get(&fh_val).unwrap();
cache.missing_ranges(offset, size as u64)
};
const CHUNK: u64 = 1024 * 1024;
for range in missing {
let mut cursor = range.start;
while cursor < range.end {
let chunk_size = (range.end - cursor).min(CHUNK) as u32;
self.fetch_counter.fetch_add(1, Ordering::Relaxed);
let bytes = match self.with_recovery(&mut inner, |fs, inner| {
let handle = fs.file_handle(inner, ino.0)?;
let storage_idx =
Self::find_storage_index(inner, ino.0).ok_or(mtp_rs::Error::NotFound)?;
fs.rt.block_on(
inner.storages[storage_idx].read_range(handle, cursor, chunk_size),
)
}) {
Ok(b) => b,
Err(e) => {
error!("MTP read_range failed at offset {cursor}: {e}");
reply.error(Errno::EIO);
return;
}
};
let bytes_len = bytes.len() as u64;
let cache = inner.read_cache.get_mut(&fh_val).unwrap();
if let Err(e) = cache.write_at(cursor, &bytes) {
error!("Sparse cache write failed: {e}");
reply.error(Errno::EIO);
return;
}
if bytes_len == 0 {
break;
}
cursor += bytes_len;
}
}
let cache = inner.read_cache.get_mut(&fh_val).unwrap();
match cache.read_at(offset, size as u64) {
Ok(buf) => reply.data(&buf),
Err(e) => {
error!("Sparse cache read failed: {e}");
reply.error(Errno::EIO);
}
}
}
fn release(
&self,
_req: &Request,
_ino: INodeNo,
fh: FileHandle,
_flags: OpenFlags,
_lock_owner: Option<LockOwner>,
_flush: bool,
reply: ReplyEmpty,
) {
let fh_val = fh.0;
let mut inner = self.inner.lock().unwrap();
let flushed = if inner.write_buf.is_open(fh_val) {
self.flush_to_mtp(&mut inner, fh_val)
} else {
Ok(())
};
inner.read_cache.remove(&fh_val);
inner.fh_to_inode.remove(&fh_val);
match flushed {
Ok(()) => reply.ok(),
Err(e) => {
error!("Flush on close failed: {e}");
reply.error(Errno::EIO);
}
}
}
fn write(
&self,
_req: &Request,
ino: INodeNo,
fh: FileHandle,
offset: u64,
data: &[u8],
_write_flags: WriteFlags,
_flags: OpenFlags,
_lock_owner: Option<LockOwner>,
reply: ReplyWrite,
) {
if self.read_only {
reply.error(Errno::EROFS);
return;
}
let fh_val = fh.0;
let mut inner = self.inner.lock().unwrap();
if !inner.write_buf.is_open(fh_val) {
let original_size = inner.inodes.get(ino.0).map(|e| e.size).unwrap_or(0);
if let Err(e) = inner.write_buf.open(fh_val, ino.0, original_size) {
error!("Failed to open write buffer: {e}");
reply.error(Errno::EIO);
return;
}
}
match inner.write_buf.write(fh_val, offset as i64, data) {
Ok(written) => reply.written(written),
Err(e) => {
error!("Write failed: {e}");
reply.error(Errno::EIO);
}
}
}
fn create(
&self,
_req: &Request,
parent: INodeNo,
name: &OsStr,
_mode: u32,
_umask: u32,
_flags: i32,
reply: ReplyCreate,
) {
if self.read_only {
reply.error(Errno::EROFS);
return;
}
let name_str = match name.to_str() {
Some(s) => s,
None => {
reply.error(Errno::EINVAL);
return;
}
};
let parent_ino = parent.0;
let mut inner = self.inner.lock().unwrap();
if !inner.inodes.get(parent_ino).is_some_and(|e| e.is_dir()) {
reply.error(Errno::ENOTDIR);
return;
}
let handle = match self.with_recovery(&mut inner, |fs, inner| {
let mtp_parent = fs.parent_handle(inner, parent_ino)?;
let storage_idx =
Self::find_storage_index(inner, parent_ino).ok_or(mtp_rs::Error::NotFound)?;
let info = NewObjectInfo::file(name_str, 0);
let stream = bytes_stream(Vec::new());
fs.rt
.block_on(inner.storages[storage_idx].upload(mtp_parent, info, stream))
.map_err(mtp_rs::Error::from)
}) {
Ok(h) => h,
Err(e) => {
error!("MTP create failed: {e}");
reply.error(Errno::EIO);
return;
}
};
let now = SystemTime::now();
let ino = inner
.inodes
.add_object(parent_ino, handle, name_str.to_string(), false, 0, now);
let fh = self.alloc_fh();
inner.fh_to_inode.insert(fh, ino);
if let Err(e) = inner.write_buf.open(fh, ino, 0) {
error!("Failed to open write buffer: {e}");
reply.error(Errno::EIO);
return;
}
let entry = inner.inodes.get(ino).unwrap();
let attr = inode_to_file_attr(entry);
reply.created(
&TTL,
&attr,
Generation(0),
FileHandle(fh),
FopenFlags::empty(),
);
}
fn mkdir(
&self,
_req: &Request,
parent: INodeNo,
name: &OsStr,
_mode: u32,
_umask: u32,
reply: ReplyEntry,
) {
if self.read_only {
reply.error(Errno::EROFS);
return;
}
let name_str = match name.to_str() {
Some(s) => s,
None => {
reply.error(Errno::EINVAL);
return;
}
};
let parent_ino = parent.0;
let mut inner = self.inner.lock().unwrap();
if !inner.inodes.get(parent_ino).is_some_and(|e| e.is_dir()) {
reply.error(Errno::ENOTDIR);
return;
}
let handle = match self.with_recovery(&mut inner, |fs, inner| {
let mtp_parent = fs.parent_handle(inner, parent_ino)?;
let storage_idx =
Self::find_storage_index(inner, parent_ino).ok_or(mtp_rs::Error::NotFound)?;
fs.rt
.block_on(inner.storages[storage_idx].create_folder(mtp_parent, name_str))
}) {
Ok(h) => h,
Err(e) => {
error!("MTP mkdir failed: {e}");
reply.error(Errno::EIO);
return;
}
};
let now = SystemTime::now();
let ino = inner
.inodes
.add_object(parent_ino, handle, name_str.to_string(), true, 0, now);
let entry = inner.inodes.get(ino).unwrap();
let attr = inode_to_file_attr(entry);
reply.entry(&TTL, &attr, Generation(0));
}
fn unlink(&self, _req: &Request, parent: INodeNo, name: &OsStr, reply: ReplyEmpty) {
if self.read_only {
reply.error(Errno::EROFS);
return;
}
let name_str = match name.to_str() {
Some(s) => s,
None => {
reply.error(Errno::ENOENT);
return;
}
};
let parent_ino = parent.0;
let mut inner = self.inner.lock().unwrap();
let child_ino = match inner.inodes.lookup(parent_ino, name_str) {
Some(i) => i,
None => {
reply.error(Errno::ENOENT);
return;
}
};
if inner.inodes.get(child_ino).is_some_and(|e| e.is_dir()) {
reply.error(Errno::EISDIR);
return;
}
if let Err(e) = self.with_recovery(&mut inner, |fs, inner| {
let handle = fs.file_handle(inner, child_ino)?;
let storage_idx =
Self::find_storage_index(inner, child_ino).ok_or(mtp_rs::Error::NotFound)?;
fs.rt.block_on(inner.storages[storage_idx].delete(handle))
}) {
error!("MTP delete failed: {e}");
reply.error(Errno::EIO);
return;
}
inner.inodes.remove(child_ino);
reply.ok();
}
fn rmdir(&self, _req: &Request, parent: INodeNo, name: &OsStr, reply: ReplyEmpty) {
if self.read_only {
reply.error(Errno::EROFS);
return;
}
let name_str = match name.to_str() {
Some(s) => s,
None => {
reply.error(Errno::ENOENT);
return;
}
};
let parent_ino = parent.0;
let mut inner = self.inner.lock().unwrap();
let child_ino = match inner.inodes.lookup(parent_ino, name_str) {
Some(i) => i,
None => {
reply.error(Errno::ENOENT);
return;
}
};
if !matches!(
inner.inodes.get(child_ino).map(|e| &e.kind),
Some(InodeKind::Directory { .. })
) {
reply.error(Errno::ENOTDIR);
return;
}
if let Err(e) = self.with_recovery(&mut inner, |fs, inner| {
let handle = fs.object_handle(inner, child_ino)?;
let storage_idx =
Self::find_storage_index(inner, child_ino).ok_or(mtp_rs::Error::NotFound)?;
fs.rt.block_on(inner.storages[storage_idx].delete(handle))
}) {
error!("MTP rmdir failed: {e}");
reply.error(Errno::EIO);
return;
}
inner.inodes.remove(child_ino);
reply.ok();
}
fn rename(
&self,
_req: &Request,
parent: INodeNo,
name: &OsStr,
newparent: INodeNo,
newname: &OsStr,
_flags: RenameFlags,
reply: ReplyEmpty,
) {
if self.read_only {
reply.error(Errno::EROFS);
return;
}
let name_str = match name.to_str() {
Some(s) => s,
None => {
reply.error(Errno::ENOENT);
return;
}
};
let newname_str = match newname.to_str() {
Some(s) => s,
None => {
reply.error(Errno::EINVAL);
return;
}
};
let parent_ino = parent.0;
let newparent_ino = newparent.0;
let mut inner = self.inner.lock().unwrap();
let child_ino = match inner.inodes.lookup(parent_ino, name_str) {
Some(i) => i,
None => {
reply.error(Errno::ENOENT);
return;
}
};
if !matches!(
inner.inodes.get(child_ino).map(|e| &e.kind),
Some(InodeKind::File { .. } | InodeKind::Directory { .. })
) {
reply.error(Errno::EINVAL);
return;
}
if parent_ino != newparent_ino
&& !inner.inodes.get(newparent_ino).is_some_and(|e| e.is_dir())
{
reply.error(Errno::ENOTDIR);
return;
}
if name_str != newname_str {
if let Err(e) = self.with_recovery(&mut inner, |fs, inner| {
let handle = fs.object_handle(inner, child_ino)?;
let storage_idx =
Self::find_storage_index(inner, child_ino).ok_or(mtp_rs::Error::NotFound)?;
fs.rt
.block_on(inner.storages[storage_idx].rename(handle, newname_str))
}) {
error!("MTP rename failed: {e}");
reply.error(Errno::EIO);
return;
}
}
if parent_ino != newparent_ino {
if let Err(e) = self.with_recovery(&mut inner, |fs, inner| {
let handle = fs.object_handle(inner, child_ino)?;
let storage_idx =
Self::find_storage_index(inner, child_ino).ok_or(mtp_rs::Error::NotFound)?;
let new_mtp_parent = fs
.parent_handle(inner, newparent_ino)?
.unwrap_or(ObjectHandle::ROOT);
fs.rt.block_on(inner.storages[storage_idx].move_object(
handle,
new_mtp_parent,
None,
))
}) {
error!("MTP move failed: {e}");
reply.error(Errno::EIO);
return;
}
}
inner
.inodes
.rename(child_ino, newparent_ino, newname_str.to_string());
reply.ok();
}
fn setattr(
&self,
_req: &Request,
ino: INodeNo,
_mode: Option<u32>,
_uid: Option<u32>,
_gid: Option<u32>,
size: Option<u64>,
_atime: Option<TimeOrNow>,
_mtime: Option<TimeOrNow>,
_ctime: Option<SystemTime>,
fh: Option<FileHandle>,
_crtime: Option<SystemTime>,
_chgtime: Option<SystemTime>,
_bkuptime: Option<SystemTime>,
_flags: Option<BsdFileFlags>,
reply: ReplyAttr,
) {
if let Some(new_size) = size {
if self.read_only {
reply.error(Errno::EROFS);
return;
}
if let Some(fh) = fh {
let fh_val = fh.0;
let mut inner = self.inner.lock().unwrap();
if !inner.write_buf.is_open(fh_val) {
let original_size = inner.inodes.get(ino.0).map(|e| e.size).unwrap_or(0);
if let Err(e) = inner.write_buf.open(fh_val, ino.0, original_size) {
error!("Failed to open write buffer: {e}");
reply.error(Errno::EIO);
return;
}
}
if new_size == 0 {
inner.write_buf.close(fh_val);
if let Err(e) = inner.write_buf.open(fh_val, ino.0, 0) {
error!("Failed to open write buffer: {e}");
reply.error(Errno::EIO);
return;
}
}
}
}
let inner = self.inner.lock().unwrap();
match inner.inodes.get(ino.0) {
Some(entry) => {
let mut attr = inode_to_file_attr(entry);
if let Some(new_size) = size {
attr.size = new_size;
attr.blocks = new_size.div_ceil(512);
}
reply.attr(&TTL, &attr);
}
None => {
reply.error(Errno::ENOENT);
}
}
}
fn statfs(&self, _req: &Request, _ino: INodeNo, reply: ReplyStatfs) {
let inner = self.inner.lock().unwrap();
let block_size: u64 = 4096;
let mut total_bytes: u64 = 0;
let mut free_bytes: u64 = 0;
for storage in &inner.storages {
total_bytes = total_bytes.saturating_add(storage.info().total_capacity);
free_bytes = free_bytes.saturating_add(storage.info().free_space);
}
let blocks = total_bytes / block_size;
let bfree = free_bytes / block_size;
reply.statfs(blocks, bfree, bfree, 0, 0, block_size as u32, 255, 0);
}
fn opendir(&self, _req: &Request, ino: INodeNo, _flags: OpenFlags, reply: ReplyOpen) {
let mut inner = self.inner.lock().unwrap();
match inner.inodes.get(ino.0) {
Some(entry) if entry.is_dir() => {
let fh = self.alloc_fh();
inner.dirs_loaded.remove(&ino.0);
reply.opened(FileHandle(fh), FopenFlags::empty());
}
Some(_) => {
reply.error(Errno::ENOTDIR);
}
None => {
reply.error(Errno::ENOENT);
}
}
}
fn releasedir(
&self,
_req: &Request,
_ino: INodeNo,
_fh: FileHandle,
_flags: OpenFlags,
reply: ReplyEmpty,
) {
reply.ok();
}
}
#[cfg(test)]
mod tests {
use super::*;
use futures::StreamExt as _;
use std::io::Write as _;
const CHUNK: usize = UPLOAD_CHUNK;
fn spool_file(content: &[u8]) -> std::fs::File {
let mut file = tempfile::tempfile().unwrap();
file.write_all(content).unwrap();
file.seek(SeekFrom::Start(0)).unwrap();
file
}
#[test]
fn file_stream_reads_lazily() {
let file = spool_file(&vec![0xABu8; CHUNK * 4]);
let mut cursor = file.try_clone().unwrap();
let mut stream = file_stream(file);
let first = futures::executor::block_on(stream.next()).unwrap().unwrap();
assert_eq!(first.len(), CHUNK);
assert_eq!(
cursor.stream_position().unwrap(),
CHUNK as u64,
"one poll must read one chunk, not the whole file"
);
}
#[test]
fn file_stream_yields_the_whole_file_then_ends() {
let content = vec![0xCDu8; CHUNK * 2 + 17];
let chunks: Vec<_> =
futures::executor::block_on(file_stream(spool_file(&content)).collect());
let sizes: Vec<_> = chunks.iter().map(|c| c.as_ref().unwrap().len()).collect();
assert_eq!(sizes, vec![CHUNK, CHUNK, 17]);
let joined: Vec<u8> = chunks
.into_iter()
.flat_map(|c| c.unwrap().to_vec())
.collect();
assert_eq!(joined, content);
}
#[test]
fn file_stream_ends_after_a_read_error() {
let path = tempfile::NamedTempFile::new().unwrap().into_temp_path();
let file = std::fs::OpenOptions::new().write(true).open(&path).unwrap();
let mut stream = file_stream(file);
assert!(futures::executor::block_on(stream.next()).unwrap().is_err());
assert!(futures::executor::block_on(stream.next()).is_none());
}
}