From 34c96fefedc7ee94c6903b5ee36ba29abbe79213 Mon Sep 17 00:00:00 2001 From: Haydn Vestal Date: Mon, 21 Sep 2026 14:32:51 +0530 Subject: [PATCH] update to the latest freqfs --- Cargo.toml | 5 ++--- README.md | 6 ++++++ examples/example.rs | 51 ++++++++++++++++++++++++++++++++++++++++++++- src/dir.rs | 14 ++++++------- src/file.rs | 18 ++++++++-------- src/lib.rs | 2 +- 6 files changed, 75 insertions(+), 21 deletions(-) diff --git a/Cargo.toml b/Cargo.toml index 960e85a..8b6e7d8 100644 --- a/Cargo.toml +++ b/Cargo.toml @@ -13,12 +13,10 @@ keywords = ["transactional", "versioning", "file", "io", "cache"] [[example]] name = "example" -required-features = ["stream"] [features] -all = ["logging", "stream"] +all = ["logging"] logging = ["log", "freqfs/logging", "txn_lock/logging"] -stream = ["freqfs/stream"] [dependencies] freqfs = { path = "../freqfs", version = "0.13", features = ["id"] } @@ -30,6 +28,7 @@ safecast = "0.2" txn_lock = { path = "../txn_lock", version = "0.11", features = ["all"] } [dev-dependencies] +tbon = { path = "../tbon", version = "0.9", features = ["tokio-io"] } destream = { path = "../destream", version = "0.10" } rand = "0.9" tokio = { version = "1.49", features = ["macros"] } diff --git a/README.md b/README.md index 96d58da..d7c0f00 100644 --- a/README.md +++ b/README.md @@ -1,2 +1,8 @@ # txfs A cached transactional filesystem layer for Rust + +## Filesystem codecs + +Callers implement `freqfs::FileLoad` and `FileSave` for their complete file-entry +type. txfs does not choose a byte codec and no longer has a `stream` feature. +The example explicitly chooses TBON; the library retains codec-independent I/O. diff --git a/examples/example.rs b/examples/example.rs index 94626fe..c2dec40 100644 --- a/examples/example.rs +++ b/examples/example.rs @@ -5,7 +5,7 @@ use std::path::PathBuf; use std::str::FromStr; use std::time::Duration; -use destream::en; +use destream::{de, en}; use freqfs::{Cache, DirLock}; use hr_id::Id; use rand::Rng; @@ -175,3 +175,52 @@ async fn main() -> Result<(), txfs::Error> { Ok(()) } + +impl de::FromStream for File { + type Context = (); + async fn from_stream(_: (), decoder: &mut D) -> Result { + decoder.decode_any(FileVisitor).await + } +} +struct FileVisitor; +impl de::Visitor for FileVisitor { + type Value = File; + fn expecting() -> &'static str { + "a filesystem entry" + } + fn visit_string(self, value: String) -> Result { + Ok(File::Text(value)) + } + async fn visit_seq(self, mut seq: A) -> Result { + let mut bytes = Vec::new(); + while let Some(byte) = seq.next_element::(()).await? { + bytes.push(byte); + } + Ok(File::Bin(bytes)) + } +} + +impl freqfs::FileLoad for File { + async fn load( + _: &std::path::Path, + file: tokio::fs::File, + _: std::fs::Metadata, + ) -> std::io::Result { + tbon::de::read_from((), file) + .await + .map_err(|error| std::io::Error::new(std::io::ErrorKind::InvalidData, error)) + } +} +impl freqfs::FileSave for File { + async fn save(&self, file: &mut tokio::fs::File) -> std::io::Result { + use futures::TryStreamExt; + use tokio::io::AsyncWriteExt; + let mut stream = tbon::en::encode(self).map_err(std::io::Error::other)?; + let mut size = 0; + while let Some(chunk) = stream.try_next().await.map_err(std::io::Error::other)? { + file.write_all(&chunk).await?; + size += chunk.len() as u64; + } + Ok(size) + } +} diff --git a/src/dir.rs b/src/dir.rs index d97bae6..46d7d2a 100644 --- a/src/dir.rs +++ b/src/dir.rs @@ -250,8 +250,8 @@ where txn_id: TxnId, ) -> Result)>> + Send + Unpin + '_> where - FE: AsType, - F: FileLoad, + FE: FileLoad + AsType, + F: Send + Sync + 'static, { let entries = self.entries.iter(txn_id).await?; let files = entries.filter_map(|(name, entry)| match &*entry { @@ -330,7 +330,7 @@ where contents: F, ) -> Result> where - FE: AsType, + FE: FileLoad + AsType, F: GetSize + Clone, { #[cfg(feature = "logging")] @@ -394,8 +394,8 @@ where name: &Id, ) -> Result> where - F: FileLoad, - FE: AsType, + F: Send + Sync + 'static, + FE: FileLoad + AsType, { if let Some(file) = self.get_file(txn_id, name).await? { file.read(txn_id).await @@ -412,8 +412,8 @@ where name: &Id, ) -> Result> where - F: FileLoad + GetSize + Clone, - FE: FileSave + AsType, + F: Send + Sync + 'static + GetSize + Clone, + FE: FileLoad + FileSave + AsType, { if let Some(file) = self.get_file(txn_id, name).await? { file.write(txn_id).await diff --git a/src/file.rs b/src/file.rs index cff4976..2339cf1 100644 --- a/src/file.rs +++ b/src/file.rs @@ -78,7 +78,7 @@ where version: F, ) -> Result where - FE: AsType, + FE: FileLoad + AsType, F: GetSize, { debug_assert!(versions @@ -154,8 +154,8 @@ where /// Lock this file for reading at the given `txn_id`. pub async fn read(&self, txn_id: TxnId) -> Result> where - F: FileLoad, - FE: AsType, + F: Send + Sync + 'static, + FE: FileLoad + AsType, { let last_modified = self.last_modified.read(txn_id).await?; let staged = { @@ -185,8 +185,8 @@ where /// Lock this file for reading at the given `txn_id` without borrowing. pub async fn into_read(self, txn_id: TxnId) -> Result> where - F: FileLoad, - FE: AsType, + F: Send + Sync + 'static, + FE: FileLoad + AsType, { self.read(txn_id).await } @@ -194,8 +194,8 @@ where /// Lock this file for writing at the given `txn_id`. pub async fn write(&self, txn_id: TxnId) -> Result> where - F: FileLoad + Clone + GetSize, - FE: AsType, + F: Send + Sync + 'static + Clone + GetSize, + FE: FileLoad + AsType, { let mut last_modified = self.last_modified.write(txn_id).await?; let version = if last_modified @@ -245,8 +245,8 @@ where /// Lock this file for writing at the given `txn_id` without borrowing. pub async fn into_write(self, txn_id: TxnId) -> Result> where - F: FileLoad + Clone + GetSize, - FE: AsType, + F: Send + Sync + 'static + Clone + GetSize, + FE: FileLoad + AsType, { self.write(txn_id).await } diff --git a/src/lib.rs b/src/lib.rs index 3046029..d2a573d 100644 --- a/src/lib.rs +++ b/src/lib.rs @@ -235,7 +235,7 @@ mod tests { .await? .map(|(name, _)| (*name).clone()) .collect::>(); - assert_eq!(entries, [name.clone()]); + assert_eq!(entries.as_slice(), std::slice::from_ref(&name)); let file = dir .get_file(Txn(7), &name) .await?