Something went wrong. Try again.
A user-space CRDT filesystem
Something went wrong. Try again.
123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485486487488489490491492493494495496497498499500501502503504505506507508509510511512513514515516517518519520521522523524525526527528529530531532533534535536537538539540541542543544545546547548549550551552553554555556557558559560561562563564565566567568569570571572573574575576577578579580581582583584585586587588589590591592593594595596597598599600601602603604605606607608609610611612613614615616617618619620621622623624625626627628629630631632633634635636//! [`HostStore`] — a [`Filesystem`] backed by a real host directory.//!//! Every file's bytes live as a real file under `host_dir`. A [`Transducer`]//! is consulted only on the few content-bearing methods//! ([`Filesystem::read`] and [`Filesystem::write`]); every other operation is//! a direct `std::fs` call against the host directory. The host directory is//! the source of truth — there is no in-memory mirror of file contents.//!//! Trade-off: writes are read-modify-write at file granularity (the whole//! file is decoded, spliced, re-encoded, and rewritten). That keeps the//! [`Transducer`] interface trivial (whole-payload `bytes -> bytes`) at the//! cost of O(file size) per write. Fine for a toy.
use std::collections::HashMap;use std::ffi::OsStr;use std::fs::{self, OpenOptions};use std::os::unix::ffi::OsStrExt;use std::os::unix::fs::MetadataExt;use std::path::{Path, PathBuf};use std::sync::{Arc, Mutex};
use async_trait::async_trait;use fskit_rs::directory_entries::Entry;use fskit_rs::{ AccessMask, CaseFormat, DirectoryEntries, Error, Filesystem, Item, ItemAttributes, ItemType, OpenMode, PathConfOperations, PreallocateFlag, ResourceIdentifier, Result, SetXattrPolicy, StatFsResult, SupportedCapabilities, SyncFlags, TaskOptions, VolumeBehavior, VolumeIdentifier, Xattrs,};
use crate::transducer::Transducer;
// FSKit reserves: invalid = 0, parent_of_root = 1, root = 2.const ROOT_ID: u64 = 2;const PARENT_OF_ROOT: u64 = 1;const FIRST_USER_ID: u64 = 3;
#[derive(Clone)]pub struct HostStore<T: Transducer> { host_dir: PathBuf, transducer: T, state: Arc<Mutex<State>>,}
#[derive(Debug, Default)]struct State { paths: HashMap<u64, PathBuf>, ids: HashMap<PathBuf, u64>, next_id: u64,}
impl<T: Transducer> HostStore<T> { pub fn new(host_dir: impl Into<PathBuf>, transducer: T) -> std::io::Result<Self> { let host_dir = host_dir.into(); fs::create_dir_all(&host_dir)?; let mut paths = HashMap::new(); let mut ids = HashMap::new(); paths.insert(ROOT_ID, host_dir.clone()); ids.insert(host_dir.clone(), ROOT_ID); Ok(Self { host_dir, transducer, state: Arc::new(Mutex::new(State { paths, ids, next_id: FIRST_USER_ID, })), }) }
pub fn host_dir(&self) -> &Path { &self.host_dir }
fn path_for(&self, id: u64) -> Result<PathBuf> { self.state .lock() .unwrap() .paths .get(&id) .cloned() .ok_or(Error::Posix(libc::ENOENT)) }
/// Look up an existing id for `path` or assign a new one. fn intern(&self, path: &Path) -> u64 { let mut s = self.state.lock().unwrap(); if let Some(&id) = s.ids.get(path) { return id; } let id = s.next_id; s.next_id += 1; s.paths.insert(id, path.to_path_buf()); s.ids.insert(path.to_path_buf(), id); id }
fn forget(&self, id: u64) { let mut s = self.state.lock().unwrap(); if let Some(path) = s.paths.remove(&id) { s.ids.remove(&path); } }
fn rebind(&self, id: u64, new_path: PathBuf) { let mut s = self.state.lock().unwrap(); if let Some(old) = s.paths.insert(id, new_path.clone()) { s.ids.remove(&old); } s.ids.insert(new_path, id); }
fn parent_id_of(&self, path: &Path) -> u64 { let Some(parent) = path.parent() else { return PARENT_OF_ROOT; }; self.state .lock() .unwrap() .ids .get(parent) .copied() .unwrap_or(ROOT_ID) }}
fn posix<T>(code: i32) -> Result<T> { Err(Error::Posix(code))}
fn io_to_posix(err: &std::io::Error) -> i32 { err.raw_os_error().unwrap_or(libc::EIO)}
fn item_type_of(meta: &fs::Metadata) -> ItemType { let ft = meta.file_type(); if ft.is_dir() { ItemType::Directory } else if ft.is_symlink() { ItemType::Symlink } else { ItemType::File }}
fn attributes_from(id: u64, parent: u64, meta: &fs::Metadata, logical_size: u64) -> ItemAttributes { let kind = item_type_of(meta); ItemAttributes { file_id: Some(id), parent_id: Some(parent), r#type: Some(kind as i32), size: Some(logical_size), alloc_size: Some(meta.blocks() * 512), link_count: Some(meta.nlink() as u32), mode: Some(meta.mode()), uid: Some(meta.uid()), gid: Some(meta.gid()), ..Default::default() }}
#[async_trait]impl<T: Transducer> Filesystem for HostStore<T> { async fn get_resource_identifier(&mut self) -> Result<ResourceIdentifier> { Ok(ResourceIdentifier { name: Some("host-store".into()), container_id: Some("host-store-volume".into()), }) }
async fn get_volume_identifier(&mut self) -> Result<VolumeIdentifier> { Ok(VolumeIdentifier { id: Some("host-store-volume".into()), name: Some("host-store".into()), }) }
async fn get_volume_behavior(&mut self) -> Result<VolumeBehavior> { Ok(VolumeBehavior { xattr_operations_inhibited: Some(true), is_access_check_inhibited: Some(true), is_volume_rename_inhibited: Some(true), is_preallocate_inhibited: Some(true), is_open_close_inhibited: Some(true), ..Default::default() }) }
async fn get_path_conf_operations(&mut self) -> Result<PathConfOperations> { Ok(PathConfOperations { maximum_link_count: 1, maximum_name_length: 255, restricts_ownership_changes: true, truncates_long_names: false, maximum_xattr_size: None, maximum_xattr_size_in_bits: None, maximum_file_size: Some(1 << 30), maximum_file_size_in_bits: None, }) }
async fn get_volume_capabilities(&mut self) -> Result<SupportedCapabilities> { Ok(SupportedCapabilities { supports_persistent_object_ids: Some(true), supports_64bit_object_ids: Some(true), supports_hidden_files: Some(false), case_format: Some(CaseFormat::Sensitive as i32), ..Default::default() }) }
async fn get_volume_statistics(&mut self) -> Result<StatFsResult> { let capacity: u64 = 1024 * 4096; let used = match fs::read_dir(&self.host_dir) { Ok(rd) => rd .filter_map(|e| e.ok().and_then(|e| e.metadata().ok())) .map(|m| m.len()) .sum(), Err(_) => 0, }; Ok(StatFsResult { block_size: 4096, io_size: 4096, total_blocks: 1024, available_blocks: 1024, free_blocks: 1024, used_blocks: 0, total_bytes: capacity, available_bytes: capacity.saturating_sub(used), free_bytes: capacity.saturating_sub(used), used_bytes: used, total_files: u64::MAX, free_files: u64::MAX, }) }
async fn mount(&mut self, _options: TaskOptions) -> Result<()> { Ok(()) }
async fn unmount(&mut self) -> Result<()> { Ok(()) }
async fn synchronize(&mut self, _flags: SyncFlags) -> Result<()> { Ok(()) }
async fn get_attributes(&mut self, item_id: u64) -> Result<ItemAttributes> { let path = self.path_for(item_id)?; let meta = fs::symlink_metadata(&path).map_err(|e| Error::Posix(io_to_posix(&e)))?; let parent = if item_id == ROOT_ID { PARENT_OF_ROOT } else { self.parent_id_of(&path) }; // Whole-file translators are length-preserving for the only impl we // ship today (Reverse / Identity). For a length-changing translator, // size would need a `decode().len()` round-trip — punt until needed. let logical_size = if meta.is_file() { meta.len() } else { 0 }; Ok(attributes_from(item_id, parent, &meta, logical_size)) }
async fn set_attributes( &mut self, item_id: u64, _attributes: ItemAttributes, ) -> Result<ItemAttributes> { // Toy: ignore; echo current. self.get_attributes(item_id).await }
async fn lookup_item(&mut self, name: &OsStr, directory_id: u64) -> Result<Item> { let dir_path = self.path_for(directory_id)?; let path = dir_path.join(name); let meta = match fs::symlink_metadata(&path) { Ok(m) => m, Err(e) if e.kind() == std::io::ErrorKind::NotFound => return posix(libc::ENOENT), Err(e) => return posix(io_to_posix(&e)), }; let id = self.intern(&path); let parent = self.parent_id_of(&path); let logical_size = if meta.is_file() { meta.len() } else { 0 }; Ok(Item { attributes: Some(attributes_from(id, parent, &meta, logical_size)), name: name.as_bytes().to_vec(), }) }
async fn reclaim_item(&mut self, _item_id: u64) -> Result<()> { Ok(()) }
async fn read_symbolic_link(&mut self, _item_id: u64) -> Result<Vec<u8>> { posix(libc::ENOSYS) }
async fn create_item( &mut self, name: &OsStr, kind: ItemType, directory_id: u64, _attributes: ItemAttributes, ) -> Result<Item> { if !matches!(kind, ItemType::File | ItemType::Directory) { return posix(libc::ENOTSUP); } let dir_path = self.path_for(directory_id)?; let path = dir_path.join(name); if path.symlink_metadata().is_ok() { return posix(libc::EEXIST); } match kind { ItemType::Directory => { fs::create_dir(&path).map_err(|e| Error::Posix(io_to_posix(&e)))?; } ItemType::File => { OpenOptions::new() .write(true) .create_new(true) .open(&path) .map_err(|e| Error::Posix(io_to_posix(&e)))?; } _ => unreachable!(), } let id = self.intern(&path); let meta = fs::symlink_metadata(&path).map_err(|e| Error::Posix(io_to_posix(&e)))?; let parent = self.parent_id_of(&path); Ok(Item { attributes: Some(attributes_from(id, parent, &meta, 0)), name: name.as_bytes().to_vec(), }) }
async fn create_symbolic_link( &mut self, _name: &OsStr, _directory_id: u64, _new_attributes: ItemAttributes, _contents: Vec<u8>, ) -> Result<Item> { posix(libc::ENOSYS) }
async fn create_link( &mut self, _item_id: u64, _name: &OsStr, _directory_id: u64, ) -> Result<Vec<u8>> { posix(libc::ENOSYS) }
async fn remove_item(&mut self, item_id: u64, name: &OsStr, directory_id: u64) -> Result<()> { let dir_path = self.path_for(directory_id)?; let path = dir_path.join(name); let meta = match fs::symlink_metadata(&path) { Ok(m) => m, Err(e) if e.kind() == std::io::ErrorKind::NotFound => return posix(libc::ENOENT), Err(e) => return posix(io_to_posix(&e)), }; let result = if meta.is_dir() { fs::remove_dir(&path) } else { fs::remove_file(&path) }; result.map_err(|e| Error::Posix(io_to_posix(&e)))?; self.forget(item_id); Ok(()) }
async fn rename_item( &mut self, item_id: u64, source_directory_id: u64, source_name: &OsStr, destination_name: &OsStr, destination_directory_id: u64, _over_item_id: Option<u64>, ) -> Result<Vec<u8>> { let src_dir = self.path_for(source_directory_id)?; let dst_dir = self.path_for(destination_directory_id)?; let src = src_dir.join(source_name); let dst = dst_dir.join(destination_name); fs::rename(&src, &dst).map_err(|e| Error::Posix(io_to_posix(&e)))?; self.rebind(item_id, dst); Ok(destination_name.as_bytes().to_vec()) }
async fn enumerate_directory( &mut self, directory_id: u64, cookie: u64, _verifier: u64, ) -> Result<DirectoryEntries> { let dir_path = self.path_for(directory_id)?; let meta = fs::symlink_metadata(&dir_path).map_err(|e| Error::Posix(io_to_posix(&e)))?; if !meta.is_dir() { return posix(libc::ENOTDIR); } // Toy pagination: cookie 0 returns everything; any other cookie means EOF. let mut entries = Vec::new(); if cookie == 0 { let rd = fs::read_dir(&dir_path).map_err(|e| Error::Posix(io_to_posix(&e)))?; for (i, dirent) in rd.enumerate() { let dirent = dirent.map_err(|e| Error::Posix(io_to_posix(&e)))?; let path = dirent.path(); let entry_meta = fs::symlink_metadata(&path).map_err(|e| Error::Posix(io_to_posix(&e)))?; let name = dirent.file_name(); let id = self.intern(&path); let parent = directory_id; let logical_size = if entry_meta.is_file() { entry_meta.len() } else { 0 }; entries.push(Entry { item: Some(Item { attributes: Some(attributes_from(id, parent, &entry_meta, logical_size)), name: name.as_bytes().to_vec(), }), next_cookie: (i as u64) + 1, }); } } Ok(DirectoryEntries { entries, verifier: 1, }) }
async fn activate(&mut self, _options: TaskOptions) -> Result<Item> { let path = self.host_dir.clone(); let meta = fs::symlink_metadata(&path).map_err(|e| Error::Posix(io_to_posix(&e)))?; Ok(Item { attributes: Some(attributes_from(ROOT_ID, PARENT_OF_ROOT, &meta, 0)), name: Vec::new(), }) }
async fn deactivate(&mut self) -> Result<()> { Ok(()) }
async fn get_supported_xattr_names(&mut self, _item_id: u64) -> Result<Xattrs> { Ok(Xattrs { names: Vec::new() }) }
async fn get_xattr(&mut self, _name: &OsStr, _item_id: u64) -> Result<Vec<u8>> { posix(libc::ENOATTR) }
async fn set_xattr( &mut self, _name: &OsStr, _value: Option<Vec<u8>>, _item_id: u64, _policy: SetXattrPolicy, ) -> Result<()> { posix(libc::ENOTSUP) }
async fn get_xattrs(&mut self, _item_id: u64) -> Result<Xattrs> { Ok(Xattrs { names: Vec::new() }) }
async fn open_item(&mut self, _item_id: u64, _modes: Vec<OpenMode>) -> Result<()> { Ok(()) }
async fn close_item(&mut self, _item_id: u64, _modes: Vec<OpenMode>) -> Result<()> { Ok(()) }
// ---- the only methods that actually consult the translator ----------
async fn read(&mut self, item_id: u64, offset: i64, length: i64) -> Result<Vec<u8>> { if offset < 0 || length < 0 { return posix(libc::EINVAL); } let path = self.path_for(item_id)?; let on_disk = match fs::read(&path) { Ok(b) => b, Err(e) if e.kind() == std::io::ErrorKind::IsADirectory => return posix(libc::EISDIR), Err(e) => return posix(io_to_posix(&e)), }; let logical = self.transducer.decode(&on_disk); let off = (offset as usize).min(logical.len()); let end = off.saturating_add(length as usize).min(logical.len()); Ok(logical[off..end].to_vec()) }
async fn write(&mut self, contents: Vec<u8>, item_id: u64, offset: i64) -> Result<i64> { if offset < 0 { return posix(libc::EINVAL); } let path = self.path_for(item_id)?; // Refuse to write to a directory before we read it. let meta = fs::symlink_metadata(&path).map_err(|e| Error::Posix(io_to_posix(&e)))?; if meta.is_dir() { return posix(libc::EISDIR); } // Read-modify-write the whole file: decode -> splice -> encode -> overwrite. let on_disk = fs::read(&path).map_err(|e| Error::Posix(io_to_posix(&e)))?; let mut logical = self.transducer.decode(&on_disk); let off = offset as usize; let written = contents.len(); let new_len = (off + written).max(logical.len()); logical.resize(new_len, 0); logical[off..off + written].copy_from_slice(&contents); let new_on_disk = self.transducer.encode(&logical); fs::write(&path, &new_on_disk).map_err(|e| Error::Posix(io_to_posix(&e)))?; Ok(written as i64) }
async fn check_access(&mut self, _item_id: u64, _access: Vec<AccessMask>) -> Result<bool> { Ok(true) }
async fn set_volume_name(&mut self, _name: Vec<u8>) -> Result<Vec<u8>> { posix(libc::ENOSYS) }
async fn preallocate_space( &mut self, _item_id: u64, _offset: i64, length: i64, _flags: Vec<PreallocateFlag>, ) -> Result<i64> { Ok(length) }
async fn deactivate_item(&mut self, _item_id: u64) -> Result<()> { Ok(()) }}
#[cfg(test)]mod tests { use super::*; use crate::transducer::Reverse; use tempfile::TempDir;
fn fs() -> (HostStore<Reverse>, TempDir) { let dir = tempfile::tempdir().unwrap(); let fs = HostStore::new(dir.path(), Reverse).unwrap(); (fs, dir) }
#[tokio::test] async fn write_then_read_round_trip() { let (mut fs, dir) = fs(); let item = fs .create_item( OsStr::new("hello.txt"), ItemType::File, ROOT_ID, ItemAttributes::default(), ) .await .unwrap(); let id = item.attributes.unwrap().file_id.unwrap();
let n = fs.write(b"hello".to_vec(), id, 0).await.unwrap(); assert_eq!(n, 5);
// On the host, bytes are reversed. let on_disk = std::fs::read(dir.path().join("hello.txt")).unwrap(); assert_eq!(on_disk, b"olleh");
// Reading via the FS returns the logical bytes. assert_eq!(fs.read(id, 0, 5).await.unwrap(), b"hello"); assert_eq!(fs.read(id, 1, 3).await.unwrap(), b"ell"); }
#[tokio::test] async fn write_at_offset_zero_fills_hole() { let (mut fs, dir) = fs(); let item = fs .create_item( OsStr::new("a"), ItemType::File, ROOT_ID, ItemAttributes::default(), ) .await .unwrap(); let id = item.attributes.unwrap().file_id.unwrap();
fs.write(b"abc".to_vec(), id, 5).await.unwrap(); let on_disk = std::fs::read(dir.path().join("a")).unwrap(); // Logical [0,0,0,0,0,'a','b','c'] reversed. assert_eq!(on_disk, b"cba\0\0\0\0\0"); assert_eq!(fs.read(id, 0, 8).await.unwrap(), b"\0\0\0\0\0abc"); }
#[tokio::test] async fn enumerate_after_create() { let (mut fs, _dir) = fs(); for name in ["a", "b", "c"] { fs.create_item( OsStr::new(name), ItemType::File, ROOT_ID, ItemAttributes::default(), ) .await .unwrap(); } let entries = fs.enumerate_directory(ROOT_ID, 0, 0).await.unwrap(); let mut names: Vec<_> = entries .entries .iter() .map(|e| String::from_utf8(e.item.as_ref().unwrap().name.clone()).unwrap()) .collect(); names.sort(); assert_eq!(names, vec!["a", "b", "c"]); }
#[tokio::test] async fn remove_then_lookup_misses() { let (mut fs, _dir) = fs(); let item = fs .create_item( OsStr::new("x"), ItemType::File, ROOT_ID, ItemAttributes::default(), ) .await .unwrap(); let id = item.attributes.unwrap().file_id.unwrap(); fs.remove_item(id, OsStr::new("x"), ROOT_ID).await.unwrap(); let err = fs.lookup_item(OsStr::new("x"), ROOT_ID).await.unwrap_err(); assert!(matches!(err, Error::Posix(c) if c == libc::ENOENT)); }}