Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
5 changes: 2 additions & 3 deletions Cargo.toml
Original file line number Diff line number Diff line change
Expand Up @@ -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"] }
Expand All @@ -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"] }
6 changes: 6 additions & 0 deletions README.md
Original file line number Diff line number Diff line change
@@ -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.
51 changes: 50 additions & 1 deletion examples/example.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand Down Expand Up @@ -175,3 +175,52 @@ async fn main() -> Result<(), txfs::Error> {

Ok(())
}

impl de::FromStream for File {
type Context = ();
async fn from_stream<D: de::Decoder>(_: (), decoder: &mut D) -> Result<Self, D::Error> {
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<E: de::Error>(self, value: String) -> Result<File, E> {
Ok(File::Text(value))
}
async fn visit_seq<A: de::SeqAccess>(self, mut seq: A) -> Result<File, A::Error> {
let mut bytes = Vec::new();
while let Some(byte) = seq.next_element::<u8>(()).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<Self> {
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<u64> {
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)
}
}
14 changes: 7 additions & 7 deletions src/dir.rs
Original file line number Diff line number Diff line change
Expand Up @@ -250,8 +250,8 @@ where
txn_id: TxnId,
) -> Result<impl Stream<Item = Result<(Key, FileVersionRead<TxnId, FE, F>)>> + Send + Unpin + '_>
where
FE: AsType<F>,
F: FileLoad,
FE: FileLoad + AsType<F>,
F: Send + Sync + 'static,
{
let entries = self.entries.iter(txn_id).await?;
let files = entries.filter_map(|(name, entry)| match &*entry {
Expand Down Expand Up @@ -330,7 +330,7 @@ where
contents: F,
) -> Result<File<TxnId, FE>>
where
FE: AsType<F>,
FE: FileLoad + AsType<F>,
F: GetSize + Clone,
{
#[cfg(feature = "logging")]
Expand Down Expand Up @@ -394,8 +394,8 @@ where
name: &Id,
) -> Result<FileVersionRead<TxnId, FE, F>>
where
F: FileLoad,
FE: AsType<F>,
F: Send + Sync + 'static,
FE: FileLoad + AsType<F>,
{
if let Some(file) = self.get_file(txn_id, name).await? {
file.read(txn_id).await
Expand All @@ -412,8 +412,8 @@ where
name: &Id,
) -> Result<FileVersionWrite<TxnId, FE, F>>
where
F: FileLoad + GetSize + Clone,
FE: FileSave + AsType<F>,
F: Send + Sync + 'static + GetSize + Clone,
FE: FileLoad + FileSave + AsType<F>,
{
if let Some(file) = self.get_file(txn_id, name).await? {
file.write(txn_id).await
Expand Down
18 changes: 9 additions & 9 deletions src/file.rs
Original file line number Diff line number Diff line change
Expand Up @@ -78,7 +78,7 @@ where
version: F,
) -> Result<Self>
where
FE: AsType<F>,
FE: FileLoad + AsType<F>,
F: GetSize,
{
debug_assert!(versions
Expand Down Expand Up @@ -154,8 +154,8 @@ where
/// Lock this file for reading at the given `txn_id`.
pub async fn read<F>(&self, txn_id: TxnId) -> Result<FileVersionRead<TxnId, FE, F>>
where
F: FileLoad,
FE: AsType<F>,
F: Send + Sync + 'static,
FE: FileLoad + AsType<F>,
{
let last_modified = self.last_modified.read(txn_id).await?;
let staged = {
Expand Down Expand Up @@ -185,17 +185,17 @@ where
/// Lock this file for reading at the given `txn_id` without borrowing.
pub async fn into_read<F>(self, txn_id: TxnId) -> Result<FileVersionRead<TxnId, FE, F>>
where
F: FileLoad,
FE: AsType<F>,
F: Send + Sync + 'static,
FE: FileLoad + AsType<F>,
{
self.read(txn_id).await
}

/// Lock this file for writing at the given `txn_id`.
pub async fn write<F>(&self, txn_id: TxnId) -> Result<FileVersionWrite<TxnId, FE, F>>
where
F: FileLoad + Clone + GetSize,
FE: AsType<F>,
F: Send + Sync + 'static + Clone + GetSize,
FE: FileLoad + AsType<F>,
{
let mut last_modified = self.last_modified.write(txn_id).await?;
let version = if last_modified
Expand Down Expand Up @@ -245,8 +245,8 @@ where
/// Lock this file for writing at the given `txn_id` without borrowing.
pub async fn into_write<F>(self, txn_id: TxnId) -> Result<FileVersionWrite<TxnId, FE, F>>
where
F: FileLoad + Clone + GetSize,
FE: AsType<F>,
F: Send + Sync + 'static + Clone + GetSize,
FE: FileLoad + AsType<F>,
{
self.write(txn_id).await
}
Expand Down
2 changes: 1 addition & 1 deletion src/lib.rs
Original file line number Diff line number Diff line change
Expand Up @@ -235,7 +235,7 @@ mod tests {
.await?
.map(|(name, _)| (*name).clone())
.collect::<Vec<_>>();
assert_eq!(entries, [name.clone()]);
assert_eq!(entries.as_slice(), std::slice::from_ref(&name));
let file = dir
.get_file(Txn(7), &name)
.await?
Expand Down
Loading