From 71bf0826902940e621e1e72b563ecc8ea83d699a Mon Sep 17 00:00:00 2001 From: dawn <90008@gaze.systems> Date: Fri, 12 Dec 2025 21:51:45 +0300 Subject: [PATCH] dont use vfs crate --- Cargo.lock | 33 ------ Cargo.toml | 1 - src/fuse.rs | 43 ++++---- src/lib.rs | 288 ++++++++++++++++++++++------------------------------ src/main.rs | 12 ++- 5 files changed, 153 insertions(+), 224 deletions(-) diff --git a/Cargo.lock b/Cargo.lock index 8951d83..9400177 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -112,7 +112,6 @@ dependencies = [ "serde_json", "tokio", "url", - "vfs", ] [[package]] @@ -746,18 +745,6 @@ dependencies = [ "subtle", ] -[[package]] -name = "filetime" -version = "0.2.26" -source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "bc0505cd1b6fa6580283f6bdf70a73fcf4aba1184038c90902b92b3dd0df63ed" -dependencies = [ - "cfg-if", - "libc", - "libredox", - "windows-sys 0.60.2", -] - [[package]] name = "find-msvc-tools" version = "0.1.5" @@ -1729,17 +1716,6 @@ version = "0.2.15" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "f9fbbcab51052fe104eb5e5d351cf728d30a5be1fe14d9be8a3b097481fb97de" -[[package]] -name = "libredox" -version = "0.1.10" -source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "416f7e718bdb06000964960ffa43b4335ad4012ae8b99060261aa4a8088d5ccb" -dependencies = [ - "bitflags", - "libc", - "redox_syscall", -] - [[package]] name = "litemap" version = "0.8.1" @@ -3557,15 +3533,6 @@ version = "0.9.5" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "0b928f33d975fc6ad9f86c8f283853ad26bdd5b10b7f1542aa2fa15e2289105a" -[[package]] -name = "vfs" -version = "0.12.2" -source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "9e723b9e1c02a3cf9f9d0de6a4ddb8cdc1df859078902fe0ae0589d615711ae6" -dependencies = [ - "filetime", -] - [[package]] name = "want" version = "0.3.1" diff --git a/Cargo.toml b/Cargo.toml index 4287044..105793d 100644 --- a/Cargo.toml +++ b/Cargo.toml @@ -4,7 +4,6 @@ version = "0.1.0" edition = "2021" [dependencies] -vfs = { version = "0.12" } bpaf = { version = "0.9", features = ["derive"] } anyhow = "1.0" scc = "2.1" diff --git a/src/fuse.rs b/src/fuse.rs index c6b2704..0ddda5e 100644 --- a/src/fuse.rs +++ b/src/fuse.rs @@ -1,16 +1,20 @@ -use super::*; use easy_fuser::{prelude::*, templates::DefaultFuseHandler}; +use tokio::runtime::Handle; use std::{ ffi::{OsStr, OsString}, - io::{Read, Seek, SeekFrom}, + io::{Cursor, Read, Seek, SeekFrom}, path::PathBuf, + sync::Arc, time::UNIX_EPOCH, }; +use crate::{AtpFS, FileType}; + pub struct AtpFuse { pub fs: Arc, pub inner: DefaultFuseHandler, + pub runtime: Handle, } impl AtpFuse { @@ -45,13 +49,13 @@ impl AtpFuse { fn vfs_metadata_attr(&self, vfs_path: &str) -> FuseResult { let meta = self - .fs - .metadata(vfs_path) + .runtime + .block_on(self.fs.metadata(vfs_path)) .map_err(|_| ErrorKind::FileNotFound.to_error("Not found"))?; let (kind, perm, nlink) = match meta.file_type { - VfsFileType::Directory => (FileKind::Directory, 0o755, 2), - VfsFileType::File => (FileKind::RegularFile, 0o644, 1), + FileType::Directory => (FileKind::Directory, 0o755, 2), + FileType::File => (FileKind::RegularFile, 0o644, 1), }; Ok(FileAttribute { @@ -115,8 +119,8 @@ impl FuseHandler for AtpFuse { let vfs_path = self.path_to_str(&file_id); let stream = self - .fs - .read_dir(&vfs_path) + .runtime + .block_on(self.fs.read_dir(&vfs_path)) .map_err(|_| ErrorKind::InputOutputError.to_error("Read dir failed"))?; let mut entries = vec![ @@ -125,10 +129,11 @@ impl FuseHandler for AtpFuse { ]; for name in stream { - let kind = name - .ends_with(".json") - .then_some(FileKind::RegularFile) - .unwrap_or(FileKind::Directory); + let kind = if name.ends_with(".json") { + FileKind::RegularFile + } else { + FileKind::Directory + }; entries.push((OsString::from(name), kind)); } @@ -145,12 +150,6 @@ impl FuseHandler for AtpFuse { _flags: FUSEOpenFlags, _lock_owner: Option, ) -> FuseResult> { - let vfs_path = self.path_to_str(&file_id); - let mut reader = self - .fs - .open_file(&vfs_path) - .map_err(|_| ErrorKind::FileNotFound.to_error("File not found"))?; - // Only support absolute start seeks for now. let pos = match seek { SeekFrom::Start(p) => p, @@ -163,13 +162,19 @@ impl FuseHandler for AtpFuse { return Ok(Vec::new()); } + let vfs_path = self.path_to_str(&file_id); + let data = self + .runtime + .block_on(self.fs.open_file(&vfs_path)) + .map_err(|_| ErrorKind::FileNotFound.to_error("File not found"))?; + let mut reader = Cursor::new(data.as_slice()); + // Seek to the requested position. reader .seek(SeekFrom::Start(pos)) .map_err(|_| ErrorKind::InputOutputError.to_error("Seek failed"))?; // Read up to `size` bytes into the buffer. - // We use take to limit the read, then read_to_end or just read into buffer. let mut buf = vec![0u8; size as usize]; let n = reader .read(&mut buf) diff --git a/src/lib.rs b/src/lib.rs index a2f04c4..02cb404 100644 --- a/src/lib.rs +++ b/src/lib.rs @@ -1,18 +1,20 @@ use anyhow::{anyhow, Result}; -#[cfg(target_arch = "wasm32")] -use futures::executor::block_on; use jacquard::{ api::com_atproto::repo::{describe_repo::DescribeRepo, list_records::ListRecords}, client::{credential_session::CredentialSession, Agent, BasicClient, MemorySessionStore}, identity::{resolver::IdentityResolver, slingshot_resolver_default}, + prelude::*, types::{did::Did, nsid::Nsid, string::Handle}, - xrpc::XrpcClient, }; use scc::{HashMap, HashSet}; use url::Url; -use vfs::{error::VfsErrorKind, FileSystem, SeekAndRead, VfsFileType, VfsMetadata, VfsResult}; -use std::{collections::HashMap as StdHashMap, fmt::Debug, sync::Arc}; +use std::{ + collections::HashMap as StdHashMap, + fmt::Debug, + io::{self, ErrorKind}, + sync::Arc, +}; pub mod cli; #[cfg(target_os = "linux")] @@ -41,9 +43,21 @@ pub async fn resolve_pds(did: &Did<'_>) -> Result { .ok_or_else(|| anyhow!("no pds endpoint in did doc")) } +#[derive(Debug, Clone, Copy, PartialEq, Eq)] +pub enum FileType { + File, + Directory, +} + +#[derive(Debug, Clone)] +pub struct Metadata { + pub file_type: FileType, + pub len: u64, +} + #[derive(Debug)] struct CachedPage { - files: StdHashMap>, + files: StdHashMap>>, next_cursor: Option, } @@ -52,51 +66,30 @@ pub struct AtpFS { client: BasicClient, cache: HashMap>, root_cache: HashSet, - #[cfg(not(target_arch = "wasm32"))] - handle: tokio::runtime::Handle, } impl Debug for AtpFS { fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result { - f.debug_struct("AtProtoFS").field("did", &self.did).finish() + f.debug_struct("AtpFS").field("did", &self.did).finish() } } impl AtpFS { - pub fn new(did: Did<'static>, pds: Url) -> Self { - #[cfg(not(target_arch = "wasm32"))] - let handle = tokio::runtime::Handle::current(); - + pub async fn new(did: Did<'static>, pds: Url) -> Self { let store = MemorySessionStore::default(); let session = CredentialSession::new(Arc::new(store), Arc::new(slingshot_resolver_default())); - #[cfg(not(target_arch = "wasm32"))] - tokio::task::block_in_place(|| handle.block_on(session.set_endpoint(pds))); - - #[cfg(target_arch = "wasm32")] - block_on(session.set_endpoint(pds)); + session.set_endpoint(pds).await; Self { did, client: Agent::new(session), cache: HashMap::default(), root_cache: HashSet::default(), - #[cfg(not(target_arch = "wasm32"))] - handle, } } - #[cfg(not(target_arch = "wasm32"))] - fn block_on(&self, future: F) -> F::Output { - tokio::task::block_in_place(move || self.handle.block_on(future)) - } - - #[cfg(target_arch = "wasm32")] - fn block_on(&self, future: F) -> F::Output { - block_on(future) - } - fn segments<'a, 's>(&'s self, path: &'a str) -> Vec<&'a str> { path.trim_matches('/') .split('/') @@ -104,27 +97,21 @@ impl AtpFS { .collect() } - fn vfs_dir_metadata() -> VfsMetadata { - VfsMetadata { - file_type: VfsFileType::Directory, + fn dir_metadata() -> Metadata { + Metadata { + file_type: FileType::Directory, len: 0, - created: None, - modified: None, - accessed: None, } } - fn vfs_file_metadata(len: u64) -> VfsMetadata { - VfsMetadata { - file_type: VfsFileType::File, + fn file_metadata(len: u64) -> Metadata { + Metadata { + file_type: FileType::File, len, - created: None, - modified: None, - accessed: None, } } - async fn ensure_root_loaded(&self) -> VfsResult { + async fn ensure_root_loaded(&self) -> io::Result<()> { if self.root_cache.is_empty() { let request = DescribeRepo::new().repo(self.did.clone()).build(); @@ -132,24 +119,25 @@ impl AtpFS { .client .send(request) .await - .map_err(|e| VfsErrorKind::Other(e.to_string()))?; + .map_err(|e| io::Error::new(ErrorKind::Other, e.to_string()))?; let output = response .into_output() - .map_err(|e| VfsErrorKind::Other(e.to_string()))?; + .map_err(|e| io::Error::new(ErrorKind::Other, e.to_string()))?; for col in output.collections { let _ = self.root_cache.insert_async(col.to_string()).await; } } - return Ok("".to_string()); + return Ok(()); } - async fn ensure_loaded(&self, path: &str) -> VfsResult { + async fn ensure_loaded(&self, path: &str) -> io::Result { let segs = self.segments(path); if segs.is_empty() { - return self.ensure_root_loaded().await; + self.ensure_root_loaded().await?; + return Ok("".to_string()); } let collection = segs[0]; @@ -158,7 +146,7 @@ impl AtpFS { } if !self.root_cache.contains(collection) { - return Err(VfsErrorKind::FileNotFound.into()); + return Err(ErrorKind::NotFound.into()); } let mut current_key = collection.to_string(); @@ -174,12 +162,12 @@ impl AtpFS { parent_cursor = Some(cursor); current_key = format!("{}/next", current_key); } else { - return Err(VfsErrorKind::FileNotFound.into()); + return Err(ErrorKind::NotFound.into()); } } else if segment.ends_with(".json") { break; } else { - return Err(VfsErrorKind::FileNotFound.into()); + return Err(ErrorKind::NotFound.into()); } } @@ -188,7 +176,7 @@ impl AtpFS { Ok(current_key) } - async fn fetch_page_if_missing(&self, key: &str, cursor: Option) -> VfsResult<()> { + async fn fetch_page_if_missing(&self, key: &str, cursor: Option) -> io::Result<()> { if self.cache.contains(key) { return Ok(()); } @@ -206,18 +194,18 @@ impl AtpFS { .client .send(request.build()) .await - .map_err(|e| VfsErrorKind::Other(e.to_string()))?; + .map_err(|e| io::Error::new(ErrorKind::Other, e.to_string()))?; let output = response .into_output() - .map_err(|e| VfsErrorKind::Other(e.to_string()))?; + .map_err(|e| io::Error::new(ErrorKind::Other, e.to_string()))?; let mut files = StdHashMap::new(); for rec in output.records { if let Some(rkey) = rec.uri.rkey() { let filename = format!("{}.json", rkey.0); let content = serde_json::to_vec_pretty(&rec.value).unwrap_or_default(); - files.insert(filename, content); + files.insert(filename, Arc::new(content)); } } @@ -234,141 +222,109 @@ impl AtpFS { Ok(()) } -} -impl FileSystem for AtpFS { - fn read_dir(&self, path: &str) -> VfsResult + Send>> { - self.block_on(async { - let segs = self.segments(path); - - if segs.is_empty() { - self.ensure_root_loaded().await?; - let mut keys = Vec::new(); - self.root_cache.scan(|k| keys.push(k.clone())); - return Ok(Box::new(keys.into_iter()) as Box + Send>); - } + pub async fn read_dir(&self, path: &str) -> io::Result> { + let segs = self.segments(path); - let cache_key = self.ensure_loaded(path).await?; + if segs.is_empty() { + self.ensure_root_loaded().await?; + let mut keys = Vec::new(); + self.root_cache.scan(|k| keys.push(k.clone())); + return Ok(keys); + } - if path.ends_with(".json") { - return Err(VfsErrorKind::Other("not a directory".into()).into()); - } + let cache_key = self.ensure_loaded(path).await?; - let page = self - .cache - .read(&cache_key, |_, v| v.clone()) - .ok_or(VfsErrorKind::FileNotFound)?; + if path.ends_with(".json") { + return Err(io::Error::new(ErrorKind::Other, "not a directory")); + } - let mut entries: Vec = page.files.keys().cloned().collect(); - if page.next_cursor.is_some() { - entries.push("next".to_string()); - } + let page = self + .cache + .read(&cache_key, |_, v| v.clone()) + .ok_or(ErrorKind::NotFound)?; - Ok(Box::new(entries.into_iter()) as Box + Send>) - }) - } + let mut entries: Vec = page.files.keys().cloned().collect(); + if page.next_cursor.is_some() { + entries.push("next".to_string()); + } - fn create_dir(&self, _path: &str) -> VfsResult<()> { - Err(VfsErrorKind::NotSupported.into()) + Ok(entries) } - fn open_file(&self, path: &str) -> VfsResult> { - self.block_on(async { - let parent_path = std::path::Path::new(path) - .parent() - .unwrap_or(std::path::Path::new("")) - .to_str() - .unwrap(); - let cache_key = self.ensure_loaded(parent_path).await?; - let filename = path.split('/').last().ok_or(VfsErrorKind::FileNotFound)?; + pub async fn open_file(&self, path: &str) -> io::Result>> { + let parent_path = std::path::Path::new(path) + .parent() + .unwrap_or(std::path::Path::new("")) + .to_str() + .unwrap(); + let cache_key = self.ensure_loaded(parent_path).await?; + let filename = path.split('/').last().ok_or(ErrorKind::NotFound)?; - let content = self - .cache - .read(&cache_key, |_, page| page.files.get(filename).cloned()) - .flatten(); - - if let Some(data) = content { - return Ok(Box::new(std::io::Cursor::new(data)) as Box); - } + let content = self + .cache + .read(&cache_key, |_, page| page.files.get(filename).cloned()) + .flatten(); - Err(VfsErrorKind::FileNotFound.into()) - }) + content.ok_or(ErrorKind::NotFound.into()) } - fn metadata(&self, path: &str) -> VfsResult { - self.block_on(async { - let segs = self.segments(path); - if segs.is_empty() { - return Ok(AtpFS::vfs_dir_metadata()); - } - - if segs.len() == 1 { - self.ensure_root_loaded().await?; - if self.root_cache.contains(segs[0]) { - return Ok(AtpFS::vfs_dir_metadata()); - } else { - return Err(VfsErrorKind::FileNotFound.into()); - } - } + pub async fn metadata(&self, path: &str) -> io::Result { + let segs = self.segments(path); + if segs.is_empty() { + return Ok(Self::dir_metadata()); + } - if let Some(last) = segs.last() { - if *last == "next" { - let parent = &path[0..path.len() - 5]; - let cache_key = self.ensure_loaded(parent).await?; - let has_next = self - .cache - .read(&cache_key, |_, v| v.next_cursor.is_some()) - .unwrap_or(false); - if has_next { - return Ok(AtpFS::vfs_dir_metadata()); - } - return Err(VfsErrorKind::FileNotFound.into()); - } + if segs.len() == 1 { + self.ensure_root_loaded().await?; + if self.root_cache.contains(segs[0]) { + return Ok(Self::dir_metadata()); + } else { + return Err(ErrorKind::NotFound.into()); } + } - if path.ends_with(".json") { - let parent_path = std::path::Path::new(path) - .parent() - .unwrap() - .to_str() - .unwrap(); - let cache_key = self.ensure_loaded(parent_path).await?; - let filename = segs.last().unwrap(); - - let len = self + if let Some(last) = segs.last() { + if *last == "next" { + let parent = &path[0..path.len() - 5]; + let cache_key = self.ensure_loaded(parent).await?; + let has_next = self .cache - .read(&cache_key, |_, page| { - page.files.get(*filename).map(|f| f.len()) - }) - .flatten(); - - if let Some(l) = len { - return Ok(AtpFS::vfs_file_metadata(l as u64)); + .read(&cache_key, |_, v| v.next_cursor.is_some()) + .unwrap_or(false); + if has_next { + return Ok(Self::dir_metadata()); } - return Err(VfsErrorKind::FileNotFound.into()); + return Err(ErrorKind::NotFound.into()); } + } - Err(VfsErrorKind::FileNotFound.into()) - }) - } - - fn exists(&self, path: &str) -> VfsResult { - Ok(self.metadata(path).is_ok()) - } + if path.ends_with(".json") { + let parent_path = std::path::Path::new(path) + .parent() + .unwrap() + .to_str() + .unwrap(); + let cache_key = self.ensure_loaded(parent_path).await?; + let filename = segs.last().unwrap(); - fn create_file(&self, _: &str) -> VfsResult> { - Err(VfsErrorKind::NotSupported.into()) - } + let len = self + .cache + .read(&cache_key, |_, page| { + page.files.get(*filename).map(|f| f.len()) + }) + .flatten(); - fn append_file(&self, _: &str) -> VfsResult> { - Err(VfsErrorKind::NotSupported.into()) - } + if let Some(l) = len { + return Ok(Self::file_metadata(l as u64)); + } + return Err(ErrorKind::NotFound.into()); + } - fn remove_file(&self, _path: &str) -> VfsResult<()> { - Err(VfsErrorKind::NotSupported.into()) + Err(ErrorKind::NotFound.into()) } - fn remove_dir(&self, _path: &str) -> VfsResult<()> { - Err(VfsErrorKind::NotSupported.into()) + pub async fn exists(&self, path: &str) -> io::Result { + Ok(self.metadata(path).await.is_ok()) } } diff --git a/src/main.rs b/src/main.rs index 557dc3c..008c5f0 100644 --- a/src/main.rs +++ b/src/main.rs @@ -3,9 +3,9 @@ use atpfs::{ cli::{opts, SubCommand}, resolve_did, resolve_pds, AtpFS, }; -use vfs::FileSystem; use std::sync::Arc; +use tokio::runtime::Handle; async fn run_app(args: Vec) -> Result<()> { let opts = opts().run_inner(args.as_slice()); @@ -23,13 +23,12 @@ async fn run_app(args: Vec) -> Result<()> { let pds = resolve_pds(&did).await?; println!("resolved PDS: {}", pds); - let fs = Arc::new(AtpFS::new(did, pds)); + let fs = Arc::new(AtpFS::new(did, pds).await); match opts.cmd { SubCommand::Ls { path } => { - println!("Listing: {}", path); - let iterator = fs.read_dir(&path)?; - for item in iterator { + let files = fs.read_dir(&path).await?; + for item in files { println!("{}", item); } } @@ -40,9 +39,12 @@ async fn run_app(args: Vec) -> Result<()> { let options = vec![MountOption::RO, MountOption::FSName("atproto".to_string())]; + let handle = Handle::current(); + let fuse_handler = AtpFuse { fs, inner: DefaultFuseHandler::new(), + runtime: handle, }; println!("mounting at {:?}...", mount_point); -- 2.51.2