diff --git a/Cargo.toml b/Cargo.toml index efac86b..858f0a5 100644 --- a/Cargo.toml +++ b/Cargo.toml @@ -5,7 +5,7 @@ edition = "2024" authors = ["Boog900"] license = "MIT" repository = "https://github.com/Cuprate/tapes" -description = "A minimal database for storing data in contigious tapes" +description = "A minimal database for storing data in contiguous tapes" [dependencies] slab = { version = "0.4" } diff --git a/README.md b/README.md index dc9e891..ce4cfd0 100644 --- a/README.md +++ b/README.md @@ -2,18 +2,47 @@ A specialised database for storing data in contiguous tapes. -Each environment supports multiple independent tapes, with ACID updates across them. There are 2 types of tapes: fixed-sized and -blob. Fixed-sized tapes store fixed sized values, which allows arbitrary lookup of values by their index, blob tapes however is -just a contiguous slice of bytes so to access data you must keep its index. +A tape is an append only log of data. This database is not a generic key/value store, it is +specialised for data that builds on top of previous data, in a way that old data is only removed +if the data created after it is removed too, and removing data is less common than adding it. +An example of such data would be a typical blockchain. -## Supported Operations +Each [`Tapes`] instance supports multiple independent tapes, with ACID updates across them. -This database is not a generic key/value store. This crate is highly specialised for storing data that builds on top of previous -data, in a way that old data will only be removed if data created after it is also removed and that removing data is less -common than adding data. An example of such data would be an average blockchain. +## Tapes -Instead of there being a reader and writer transactions, write transactions are split into append and pop. This means there -are 3 transaction types: read, append and pop. Splitting the writer transaction like this makes the database more efficient -at the cost of not being able to do a single atomic rewrite of data. It is important to note though that removing data and -then writing more is still ACID, it's just the database could be left with the data being removed without the new data being -written. \ No newline at end of file +There are 2 kinds of tapes: fixed sized tapes store fixed sized values, which allows lookup of +values by their index, blob tapes are a contiguous slice of bytes, so to access data you must +keep its index. + +A [`WholeBlobTape`] stores all its data in a single file. A [`RollingBlobTape`] stores its data +in multiple files of at most a configured size, and deletes the files that hold removed data +once no reader needs them anymore. On a whole tape removing data from the front does not free up +disk space, on a rolling tape it does. Use a rolling tape when the tape grows without bound and +you need to remove old data FIFO, if not use a whole tape. + +A [`FixedSizedTape`] is a handle to a blob tape that reads and writes fixed sized entries instead +of raw bytes, the entry type must be plain data with no padding [`bytemuck::Pod`]. + +A [`CachedBlobTape`] is a wrapper around a blob tape that keeps the top of the tape in memory, +this speeds up access to recent data and reduces disk I/O. These are both wrappers, so they can +be combined, for example a fixed sized tape over a cached rolling tape. + +You probably want to always use a [`CachedBlobTape`] to prevent doing slow direct I/O. + +## Transactions + +Write transactions are split into append and pop, this means there are 3 transaction types: +read, append and pop. Splitting the writer transaction like this makes the database more +efficient at the cost of not being able to do a single atomic rewrite of data. Removing data and +then writing more is still ACID, it's just the database could be left with the data being +removed without the new data being written. + +## Persistence + +Commits are persisted according to the [`Persistence`] mode they are committed with: + +- `Buffer` writes to the OS buffer only, it is not durable, data committed with it can be lost + on a crash until it is flushed to disk by a later commit with `SyncData` or `SyncAll`. +- `SyncData` syncs the file contents to disk. +- `SyncAll` syncs the file contents and file metadata to disk. diff --git a/src/io_helpers.rs b/src/io_helpers.rs new file mode 100644 index 0000000..8853d0e --- /dev/null +++ b/src/io_helpers.rs @@ -0,0 +1,61 @@ +use std::{fs::File, io}; + +pub(crate) fn read_exact_at_file(file: &File, buf: &mut [u8], offset: u64) -> io::Result<()> { + #[cfg(unix)] + { + use std::os::unix::fs::FileExt; + + file.read_exact_at(buf, offset) + } + + #[cfg(windows)] + { + use std::os::windows::fs::FileExt; + + let mut buf = buf; + let mut offset = offset; + while !buf.is_empty() { + match file.seek_read(buf, offset) { + Ok(0) => { + break; + } + Ok(n) => { + buf = &mut buf[n..]; + offset += n as u64; + } + Err(e) => { + return Err(e); + } + } + } + + if !buf.is_empty() { + Err(io::Error::new( + io::ErrorKind::UnexpectedEof, + "failed to fill the whole buffer", + )) + } else { + Ok(()) + } + } +} + +pub(crate) fn write_all_at(file: &File, buf: &[u8], offset: u64) -> io::Result<()> { + #[cfg(unix)] + { + use std::os::unix::fs::FileExt; + + file.write_all_at(buf, offset) + } + #[cfg(windows)] + { + use std::os::windows::fs::FileExt; + + let n = file.seek_write(buf, offset)?; + if n != buf.len() { + return Err(io::Error::other("Failed to write all bytes to tape")); + } + + Ok(()) + } +} diff --git a/src/lib.rs b/src/lib.rs index 7ca7c86..729984c 100644 --- a/src/lib.rs +++ b/src/lib.rs @@ -1,18 +1,28 @@ +#![doc = include_str!("../README.md")] + mod metadata; -mod ring_buffer; +mod io_helpers; mod tapes; mod traits; pub use tapes::{ - BlobTape, FixedSizedTape, TapeOpenOptions, Tapes, TapesAppendTransaction, TapesReadTransaction, - TapesTruncateTransaction, + CachedBlobTape, CachedTapeOpenOptions, FixedSizedTape, RollingBlobTape, RollingTapeOpenOptions, + Tapes, TapesAppendTransaction, TapesReadTransaction, TapesTruncateTransaction, WholeBlobTape, + WholeTapeOpenOptions, }; -pub use traits::{TapesAppend, TapesRead, TapesTruncate}; +pub use traits::{BlobTape, TapesAppend, TapesRead, TapesTruncate}; +/// How a commit is persisted. +/// +/// `Buffer` is not durable, data committed with it can be lost on a crash until it is flushed to disk +/// by a later commit with `SyncData` or `SyncAll`. #[derive(Clone, Copy, Debug, PartialEq, Eq, Hash)] pub enum Persistence { + /// Writes to the OS buffer only, not durable. Buffer, + /// Syncs the file contents to disk. SyncData, + /// Syncs the file contents and file metadata to disk. SyncAll, } diff --git a/src/metadata.rs b/src/metadata.rs index 5fdafc8..9716ed0 100644 --- a/src/metadata.rs +++ b/src/metadata.rs @@ -14,12 +14,19 @@ use slab::Slab; use crate::Persistence; -pub type ActiveMetadata = HashMap, u64>; +pub type ActiveMetadata = HashMap, TapeMetadata>; + +#[derive(Clone, Copy, Default, BorshDeserialize, BorshSerialize)] +pub struct TapeMetadata { + pub len: u64, + pub start: u64, +} pub struct MetadataGuard { active_metadata: Arc, metadata: Arc, reader_slot: usize, + pub epoch: u64, } impl Deref for MetadataGuard { @@ -113,6 +120,7 @@ impl Metadata { active_metadata: inner.active_metadata.clone(), metadata: Arc::clone(self), reader_slot, + epoch, } } @@ -134,6 +142,16 @@ impl Metadata { Ok(()) } + + pub fn oldest_reader_excluding_reader(&self, guard: &MetadataGuard) -> Option { + self.inner + .lock() + .reader_epochs + .iter() + .filter(|(r, _)| *r != guard.reader_slot) + .map(|(_, r)| *r) + .min() + } } struct MetadataBackingFiles { @@ -213,6 +231,7 @@ impl MetadataBackingFiles { #[derive(BorshSerialize, BorshDeserialize)] struct StoredMetadata { + version: u32, hash: [u8; 32], epoch: u64, tapes: Vec, @@ -221,6 +240,7 @@ struct StoredMetadata { impl Default for StoredMetadata { fn default() -> Self { StoredMetadata { + version: 1, hash: [0; 32], epoch: 0, tapes: borsh::to_vec(&HashMap::, u64>::new()).unwrap(), @@ -244,6 +264,7 @@ fn serialise_metadata(epoch: u64, metadata: &ActiveMetadata) -> Vec { let hash = hasher.finalize().into(); borsh::to_vec(&StoredMetadata { + version: 1, hash, epoch, tapes: tapes_bytes, diff --git a/src/tapes.rs b/src/tapes.rs index f1f3ed9..e7096a1 100644 --- a/src/tapes.rs +++ b/src/tapes.rs @@ -1,50 +1,37 @@ use std::{ - cmp::min, + cmp::{max, min}, collections::HashMap, - fs::{File, OpenOptions}, io, marker::PhantomData, ops::Deref, - path::{Path, PathBuf}, + path::Path, sync::Arc, }; -use parking_lot::RwLock; - use crate::{ Persistence, - metadata::{Metadata, MetadataGuard}, - ring_buffer::{RingBuffer, RingBufferFileWriter}, - traits::{TapesAppend, TapesRead, TapesTruncate, read_exact_at}, + metadata::{Metadata, MetadataGuard, TapeMetadata}, + traits::{BlobTape, BlobTapeWriter, OpenConfig, TapesAppend, TapesRead, TapesTruncate}, }; +mod cached_tape; pub(crate) mod fixed_sized_iter; +mod rolling_tape; +mod whole_tape; -/// Configuration options for opening a tape. -pub struct TapeOpenOptions { - /// The size of the top cache in bytes, this amount of data from the top of the tape will be cached in memory. - pub top_cache_size: u64, - /// The directory to store the tapes. - pub dir: PathBuf, -} +pub use cached_tape::{CachedBlobTape, CachedTapeOpenOptions}; + +pub use rolling_tape::{RollingBlobTape, RollingTapeOpenOptions}; +pub use whole_tape::{WholeBlobTape, WholeTapeOpenOptions}; /// A handle to a fixed-sized tape. /// /// Only a single handle to a tape should be opened. -pub struct FixedSizedTape { - pub(crate) inner: BlobTape, +pub struct FixedSizedTape { + pub(crate) inner: B, phantom_data: PhantomData, } -/// A handle to a blob tape. -/// -/// Only a single handle to a tape should be opened. -pub struct BlobTape { - name: &'static str, - pub(crate) file: Arc, - pub top_cache: Arc>, -} - /// A tapes database. pub struct Tapes { metadata: Arc, @@ -53,7 +40,7 @@ pub struct Tapes { impl Tapes { /// Open a tapes database, with metadata stored at `path`. pub fn open(path: &Path) -> io::Result { - let metadata = Metadata::open(path)?; + let metadata = Metadata::open(&path.join("tapes"))?; Ok(Self { metadata: Arc::new(metadata), @@ -66,6 +53,7 @@ impl Tapes { metadata: Arc::clone(&self.metadata), metadata_guard: self.metadata.metadata(true), modified_tapes: HashMap::new(), + deleted_tapes: Vec::new(), committed: false, } } @@ -85,58 +73,36 @@ impl Tapes { modified_tapes: HashMap::new(), } } - - /// Deletes a tape and its backing file. - /// - /// This method returns [`io::ErrorKind::WouldBlock`] if a transaction is still active. - pub fn delete_tape(&mut self, name: &str, options: &TapeOpenOptions) -> io::Result<()> { - if Arc::strong_count(&self.metadata) != 1 { - return Err(io::Error::new( - io::ErrorKind::WouldBlock, - "cannot delete a tape while a transaction is active", - )); - } - - let metadata_guard = self.metadata.metadata(false); - - if metadata_guard.contains_key(name) { - let mut new_metadata = metadata_guard.clone(); - new_metadata.remove(name); - self.metadata - .update_metadata(new_metadata, true, Persistence::SyncAll)?; - } - - drop(metadata_guard); - - match std::fs::remove_file(options.dir.join(name)) { - Ok(()) => Ok(()), - Err(error) if error.kind() == io::ErrorKind::NotFound => Ok(()), - Err(error) => Err(error), - } - } } /// A tapes appender. +#[expect(clippy::type_complexity)] pub struct TapesAppendTransaction { metadata: Arc, metadata_guard: MetadataGuard, - modified_tapes: HashMap<&'static str, RingBufferFileWriter>, + modified_tapes: HashMap<&'static str, (Box, u64)>, + deleted_tapes: Vec<(&'static str, Box io::Result<()>>)>, committed: bool, } impl TapesAppendTransaction { + /// Checks if a tape exists. + pub fn tape_exists(&self, name: &'static str) -> bool { + self.metadata_guard.contains_key(name) || self.modified_tapes.contains_key(name) + } + /// Opens or creates a fixed-sized tape. - pub fn open_fixed_sized_tape( + pub fn open_fixed_sized_tape( &mut self, name: &'static str, - options: &TapeOpenOptions, - ) -> io::Result> { + options: B::OpenConfig, + ) -> io::Result> { let inner = self.open_blob_tape(name, options)?; if self .metadata_guard .get(name) - .is_some_and(|len| !(*len as usize).is_multiple_of(size_of::())) + .is_some_and(|metadata| !(metadata.len as usize).is_multiple_of(size_of::())) { return Err(io::Error::other( "Tape size is not a multiple of entry size", @@ -150,107 +116,86 @@ impl TapesAppendTransaction { } /// Opens or creates a blob tape. - pub fn open_blob_tape( + pub fn open_blob_tape( &mut self, name: &'static str, - options: &TapeOpenOptions, - ) -> io::Result { - match OpenOptions::new() - .write(true) - .read(true) - .open(options.dir.join(name)) - { - Ok(file) => { - let len = *self.metadata_guard.get(name).unwrap_or(&0); - - if file.metadata()?.len() < len { - return Err(io::Error::other("Tape file is too small")); - } - - let mut ring_buffer = RingBuffer::new(options.top_cache_size as usize, 0); - let buf = ring_buffer.reset( - min(len, options.top_cache_size) as usize, - len.saturating_sub(options.top_cache_size) as usize, - ); - - read_exact_at(&file, buf, len - buf.len() as u64)?; - - let top_cache = Arc::new(RwLock::new(ring_buffer)); - let file = Arc::new(file); - - if self.metadata_guard.get(name).is_none() { - self.modified_tapes.insert( - name, - RingBufferFileWriter { - ring_buffer: top_cache.clone(), - bytes_to_flush: 0, - file: file.clone(), - len, - }, - ); - } - - Ok(BlobTape { - name, - file, - top_cache, - }) - } - Err(e) if e.kind() == io::ErrorKind::NotFound => { - if self.metadata_guard.get(name).is_some() { - return Err(io::Error::other( - "tape was in metadata but file was not found.", - )); - } - - let file = Arc::new( - OpenOptions::new() - .write(true) - .read(true) - .create(true) - .truncate(true) - .open(options.dir.join(name))?, - ); - - let top_cache = Arc::new(RwLock::new(RingBuffer::new( - options.top_cache_size as usize, - 0, - ))); - - self.modified_tapes.insert( - name, - RingBufferFileWriter { - ring_buffer: top_cache.clone(), - bytes_to_flush: 0, - file: file.clone(), - len: 0, - }, - ); - - Ok(BlobTape { - name, - file, - top_cache, - }) - } - Err(e) => Err(e), - } + options: B::OpenConfig, + ) -> io::Result { + let metadata = self.metadata_guard.get(name).copied(); + let start_index = metadata.map_or(options.start_index(), |m| m.start); + let len = metadata.map_or(start_index, |m| m.len); + + let tape = B::open(name, metadata, self.metadata_guard.epoch, options)?; + let w = tape.writer(len)?; + + self.modified_tapes.insert(name, (Box::new(w), start_index)); + + Ok(tape) + } + + /// Deletes a tape. + /// + /// The deletion takes effect when the transaction is committed. Dropping the transaction + /// without committing keeps the tape. + /// + /// Deleting a tape that is also modified in this transaction is not supported and may cause an + /// error. + /// + /// No other transactions should be active during a write transaction that deletes a tape. + pub fn delete_tape(&mut self, tape: B) { + self.deleted_tapes + .push((tape.name(), Box::new(move || tape.delete()))); } /// Commit and consume this transaction. + /// + /// # Errors + /// + /// Fails if flushing a tape, updating the metadata, or deleting a tape fails. + /// + /// If several tapes are deleted and one of the deletions fails, the metadata update and the + /// deletions that already succeeded are not rolled back, which can leave orphaned tape files + /// behind. pub fn commit(mut self, persistence: Persistence) -> io::Result<()> { let mut new_metadata = self.metadata_guard.deref().clone(); - for (&name, tape) in &mut self.modified_tapes { + for (&name, (tape, start_index)) in &mut self.modified_tapes { tape.flush(persistence)?; - new_metadata.insert(name.into(), tape.len); + let metadata = TapeMetadata { + len: tape.len(), + start: *start_index, + }; + + new_metadata.insert(name.into(), metadata); + } + + let mut delete_tape_fns = Vec::with_capacity(self.deleted_tapes.len()); + for (name, tape) in self.deleted_tapes.drain(..) { + new_metadata.remove(name); + delete_tape_fns.push(tape); } self.metadata - .update_metadata(new_metadata, false, persistence)?; + .update_metadata(new_metadata.clone(), false, persistence)?; self.committed = true; + for delete in delete_tape_fns.drain(..) { + delete()?; + } + + let oldest_reader = self + .metadata + .oldest_reader_excluding_reader(&self.metadata_guard) + .unwrap_or(self.metadata_guard.epoch + 1); + for (&name, (tape, _)) in &mut self.modified_tapes { + tape.remove_old_files( + *new_metadata.get(name).unwrap(), + self.metadata_guard.epoch, + oldest_reader, + )?; + } + Ok(()) } } @@ -261,103 +206,184 @@ impl Drop for TapesAppendTransaction { return; } - for (&name, tape) in &self.modified_tapes { - let committed_len = self.metadata_guard.get(name).copied().unwrap_or(0); - debug_assert!(tape.len >= committed_len); - - let appended = tape.len.saturating_sub(committed_len); - tape.ring_buffer.write().pop(appended as usize); + for (&name, (tape, _)) in &mut self.modified_tapes { + let committed_len = self + .metadata_guard + .get(name) + .copied() + .unwrap_or_default() + .len; + debug_assert!(tape.len() >= committed_len); + + let appended = tape.len().saturating_sub(committed_len); + tape.revert(appended as usize); } } } impl TapesRead for TapesAppendTransaction { - fn blob_tape_len(&self, tape: &BlobTape) -> Option { + fn blob_tape_len(&self, tape: &B) -> Option { + self.modified_tapes + .get(tape.name()) + .map(|tape| tape.0.len()) + .or_else(|| { + self.metadata_guard + .get(tape.name()) + .map(|metadata| metadata.len) + }) + } + + fn blob_tape_start(&self, tape: &B) -> Option { self.modified_tapes - .get(tape.name) - .map(|tape| tape.len) - .or_else(|| self.metadata_guard.get(tape.name).copied()) + .get(tape.name()) + .map(|tape| tape.1) + .or_else(|| { + self.metadata_guard + .get(tape.name()) + .map(|metadata| metadata.start) + }) } } impl TapesAppend for TapesAppendTransaction { - fn append_bytes(&mut self, blob_tape: &BlobTape, buf: &[u8]) -> io::Result { - let tape = match self.modified_tapes.get_mut(blob_tape.name) { + fn append_bytes(&mut self, blob_tape: &B, buf: &[u8]) -> io::Result { + let tape = match self.modified_tapes.get_mut(blob_tape.name()) { Some(tape) => tape, None => { - let tape_len = *self + let metadata = self .metadata_guard - .get(blob_tape.name) + .get(blob_tape.name()) .ok_or(io::Error::other("Tape does not exist"))?; - self.modified_tapes.insert( - blob_tape.name, - RingBufferFileWriter { - ring_buffer: Arc::clone(&blob_tape.top_cache), - bytes_to_flush: 0, - file: Arc::clone(&blob_tape.file), - len: tape_len, - }, - ); - - self.modified_tapes.get_mut(blob_tape.name).unwrap() + + let w = blob_tape.writer(metadata.len)?; + + self.modified_tapes + .insert(blob_tape.name(), (Box::new(w), metadata.start)); + + self.modified_tapes.get_mut(blob_tape.name()).unwrap() } }; - tape.write(buf) + tape.0.write_bytes(buf) + } + + fn shift_start_idx(&mut self, blob_tape: &B, new_start: u64) -> io::Result<()> { + let tape = match self.modified_tapes.get_mut(blob_tape.name()) { + Some(tape) => tape, + None => { + let metadata = self + .metadata_guard + .get(blob_tape.name()) + .ok_or(io::Error::other("Tape does not exist"))?; + + let w = blob_tape.writer(metadata.len)?; + + self.modified_tapes + .insert(blob_tape.name(), (Box::new(w), metadata.start)); + + self.modified_tapes.get_mut(blob_tape.name()).unwrap() + } + }; + + if tape.0.len() < new_start { + return Err(io::Error::other( + "Start index cannot be set past the length of the tape", + )); + } + + tape.1 = max(new_start, tape.1); + + Ok(()) } } +/// A tapes truncator. pub struct TapesTruncateTransaction { metadata: Arc, metadata_guard: MetadataGuard, - modified_tapes: HashMap<&'static str, TruncatedTape>, -} - -struct TruncatedTape { - new_len: u64, - top_cache: Arc>, + modified_tapes: HashMap<&'static str, Box>, } impl TapesTruncateTransaction { - pub fn commit(self, persistence: Persistence) -> io::Result<()> { + /// Commit and consume this transaction. + pub fn commit(mut self, persistence: Persistence) -> io::Result<()> { let mut new_metadata = self.metadata_guard.deref().clone(); for (&name, tape) in &self.modified_tapes { - new_metadata.insert(name.into(), tape.new_len); + let start = min( + new_metadata.get(name).map(|m| m.start).unwrap_or_default(), + tape.len(), + ); + + new_metadata.insert( + name.into(), + TapeMetadata { + len: tape.len(), + start, + }, + ); + } + + for tape in self.modified_tapes.values_mut() { + tape.flush(persistence)?; } self.metadata .update_metadata(new_metadata, true, persistence)?; - for tape in self.modified_tapes.values() { - tape.top_cache.write().truncate(tape.new_len as usize); - } - Ok(()) } } impl TapesRead for TapesTruncateTransaction { - fn blob_tape_len(&self, tape: &BlobTape) -> Option { + fn blob_tape_len(&self, tape: &B) -> Option { self.modified_tapes - .get(tape.name) - .map(|tape| tape.new_len) - .or_else(|| self.metadata_guard.get(tape.name).copied()) + .get(tape.name()) + .map(|tape| tape.len()) + .or_else(|| { + self.metadata_guard + .get(tape.name()) + .map(|metadata| metadata.len) + }) + } + + fn blob_tape_start(&self, tape: &B) -> Option { + self.metadata_guard + .get(tape.name()) + .map(|metadata| metadata.start) } } impl TapesTruncate for TapesTruncateTransaction { - fn truncate_blob_tape(&mut self, tape: &BlobTape, new_len: u64) { - let old_len = self.blob_tape_len(tape).unwrap(); - assert!(old_len >= new_len); - - self.modified_tapes.insert( - tape.name, - TruncatedTape { - new_len, - top_cache: Arc::clone(&tape.top_cache), - }, - ); + fn truncate_blob_tape(&mut self, tape: &B, new_len: u64) -> io::Result<()> { + let Some(old_len) = self.blob_tape_len(tape) else { + return Err(io::Error::other("Tape does not exist")); + }; + + if old_len < new_len { + return Err(io::Error::other( + "Cannot truncate a tape to a longer length", + )); + } + + let tape = match self.modified_tapes.get_mut(tape.name()) { + Some(tape) => tape, + None => { + let tape_len = self + .metadata_guard + .get(tape.name()) + .ok_or(io::Error::other("Tape does not exist"))? + .len; + self.modified_tapes + .insert(tape.name(), Box::new(tape.writer(tape_len)?)); + + self.modified_tapes.get_mut(tape.name()).unwrap() + } + }; + + tape.truncate(new_len); + + Ok(()) } } @@ -373,7 +399,15 @@ impl TapesRead for TapesReadTransaction { /// Returns the number of bytes in a blob tape. /// /// Returns `None` if the tape doesn't exist. - fn blob_tape_len(&self, tape: &BlobTape) -> Option { - self.metadata_guard.get(tape.name).copied() + fn blob_tape_len(&self, tape: &B) -> Option { + self.metadata_guard + .get(tape.name()) + .map(|metadata| metadata.len) + } + + fn blob_tape_start(&self, tape: &B) -> Option { + self.metadata_guard + .get(tape.name()) + .map(|metadata| metadata.start) } } diff --git a/src/tapes/cached_tape.rs b/src/tapes/cached_tape.rs new file mode 100644 index 0000000..633e045 --- /dev/null +++ b/src/tapes/cached_tape.rs @@ -0,0 +1,245 @@ +use std::{ + cmp::{Ordering, max}, + io, + sync::Arc, +}; + +use parking_lot::RwLock; + +use crate::{ + Persistence, + metadata::TapeMetadata, + traits::{BlobTape, BlobTapeWriter, OpenConfig}, +}; + +mod ring_buffer; +use ring_buffer::RingBuffer; + +/// Open options for a [`CachedBlobTape`]. +pub struct CachedTapeOpenOptions { + /// The inner tapes' config. + pub inner: B::OpenConfig, + /// The size of the top cache in bytes, this amount of data from the top of the tape will be cached in memory. + pub top_cache_size: u64, +} + +impl OpenConfig for CachedTapeOpenOptions { + fn start_index(&self) -> u64 { + self.inner.start_index() + } +} + +impl> Clone for CachedTapeOpenOptions { + fn clone(&self) -> Self { + CachedTapeOpenOptions { + inner: self.inner.clone(), + top_cache_size: self.top_cache_size, + } + } +} + +/// A wrapper for a [`BlobTape`] that caches bytes from the top of the tape in memory to speed +/// up access and reduce disk I/O. +pub struct CachedBlobTape { + tape: B, + cache: Arc>, +} + +impl BlobTape for CachedBlobTape { + type OpenConfig = CachedTapeOpenOptions; + type Writer = CachedBlobTapeWriter; + + fn name(&self) -> &'static str { + self.tape.name() + } + + fn open( + name: &'static str, + tape_metadata: Option, + current_epoch: u64, + config: Self::OpenConfig, + ) -> io::Result { + let tape = B::open(name, tape_metadata, current_epoch, config.inner)?; + + let metadata = tape_metadata.unwrap_or_default(); + let start = max( + metadata.len.saturating_sub(config.top_cache_size), + metadata.start, + ); + + let mut ring_buffer = RingBuffer::new(config.top_cache_size as usize, 0); + let buf = ring_buffer.reset((metadata.len - start) as usize, start as usize); + + tape.read_bytes(start, buf)?; + + let cache = Arc::new(RwLock::new(ring_buffer)); + + Ok(Self { tape, cache }) + } + + fn read_bytes(&self, offset: u64, buf: &mut [u8]) -> io::Result<()> { + let top_cache = self.cache.read(); + let cached_offset = top_cache.cache_start_idx() as u64; + + let mut last_byte_needed_offset = offset + buf.len() as u64; + + if last_byte_needed_offset > cached_offset { + let read_start = offset.saturating_sub(cached_offset); + let buf_to_fill = &mut buf[(cached_offset.saturating_sub(offset)) as usize..]; + + top_cache.fill(read_start as usize, buf_to_fill); + last_byte_needed_offset -= buf_to_fill.len() as u64; + } + + if last_byte_needed_offset != offset { + self.tape.read_bytes( + offset, + &mut buf[0..(last_byte_needed_offset - offset) as usize], + )?; + } + + Ok(()) + } + + fn writer(&self, len: u64) -> io::Result { + Ok(CachedBlobTapeWriter { + ring_buffer: self.cache.clone(), + tape_writer: self.tape.writer(len)?, + bytes_to_flush: 0, + len, + is_truncation: None, + }) + } + + fn delete(self) -> io::Result<()> { + self.tape.delete() + } +} + +/// A writer that writes to a [`CachedBlobTape`]. +/// +/// The cache is flushed to disk when it is full, or on commit. +/// +/// This should not be used directly. +pub struct CachedBlobTapeWriter { + pub(crate) ring_buffer: Arc>, + pub(crate) bytes_to_flush: usize, + pub(crate) tape_writer: B::Writer, + pub(crate) len: u64, + is_truncation: Option, +} + +impl BlobTapeWriter for CachedBlobTapeWriter { + fn write_bytes(&mut self, buf: &[u8]) -> io::Result { + assert!(self.is_truncation.is_none_or(|x| !x)); + self.is_truncation = Some(false); + + let mut ring_buffer = self.ring_buffer.write(); + ring_buffer.prepare_write(self.len as usize); + let capacity = ring_buffer.capacity(); + + // Writing enough data to completely fill the ring buffer. + if buf.len() >= capacity { + flush::( + &mut self.tape_writer, + &ring_buffer, + self.bytes_to_flush, + Persistence::Buffer, + )?; + self.bytes_to_flush = 0; + + self.tape_writer.write_bytes(buf)?; + + ring_buffer.push(&buf[buf.len() - capacity..], buf.len() - capacity) + } + // Writing enough data to push data that hasn't been flushed to disk yet out of the ring buffer. + else if self.bytes_to_flush + buf.len() > capacity { + // Just flush everything that needs to be flushed to disk to reduce the number of flushes. + flush::( + &mut self.tape_writer, + &ring_buffer, + self.bytes_to_flush, + Persistence::Buffer, + )?; + self.bytes_to_flush = 0; + + self.bytes_to_flush += buf.len(); + ring_buffer.push(buf, 0) + } + // Writing data that won't push data that hasn't been flushed to disk yet out of the ring buffer. + else { + self.bytes_to_flush += buf.len(); + ring_buffer.push(buf, 0) + } + + let old_len = self.len; + self.len += buf.len() as u64; + + Ok(old_len) + } + + fn flush(&mut self, persistence: Persistence) -> io::Result<()> { + if self.is_truncation.is_some_and(|x| x) { + self.ring_buffer.write().truncate(self.len as usize); + return self.tape_writer.flush(persistence); + } + + flush::( + &mut self.tape_writer, + &self.ring_buffer.read(), + self.bytes_to_flush, + persistence, + ) + } + + fn truncate(&mut self, new_len: u64) { + assert!(self.is_truncation.is_none_or(|x| x)); + self.is_truncation = Some(true); + + self.len = new_len; + self.tape_writer.truncate(new_len) + } + + fn revert(&mut self, bytes: usize) { + self.ring_buffer.write().pop(bytes); + } + + fn len(&self) -> u64 { + self.len + } + + fn remove_old_files( + &self, + metadata: TapeMetadata, + current_epoch: u64, + oldest_reader_epoch: u64, + ) -> io::Result<()> { + self.tape_writer + .remove_old_files(metadata, current_epoch, oldest_reader_epoch) + } +} + +fn flush( + tape_writer: &mut B::Writer, + ring_buffer: &RingBuffer, + bytes_to_flush: usize, + persistence: Persistence, +) -> io::Result<()> { + if bytes_to_flush != 0 { + let (fist_slice, second_slice) = ring_buffer.as_slices(); + match bytes_to_flush.cmp(&second_slice.len()) { + Ordering::Less | Ordering::Equal => { + tape_writer.write_bytes(&second_slice[second_slice.len() - bytes_to_flush..])?; + } + Ordering::Greater => { + let first_slice_top_needed = bytes_to_flush - second_slice.len(); + + tape_writer + .write_bytes(&fist_slice[fist_slice.len() - first_slice_top_needed..])?; + tape_writer.write_bytes(second_slice)?; + } + } + } + + tape_writer.flush(persistence) +} diff --git a/src/ring_buffer.rs b/src/tapes/cached_tape/ring_buffer.rs similarity index 53% rename from src/ring_buffer.rs rename to src/tapes/cached_tape/ring_buffer.rs index a318284..bca1dae 100644 --- a/src/ring_buffer.rs +++ b/src/tapes/cached_tape/ring_buffer.rs @@ -1,13 +1,4 @@ -use std::{ - cmp::{Ordering, min}, - fs::File, - io, - sync::Arc, -}; - -use parking_lot::RwLock; - -use crate::Persistence; +use std::cmp::min; /// A simple ring buffer to cache the top of a tape. #[derive(Debug)] @@ -134,7 +125,7 @@ impl RingBuffer { self.len = self.len.min(tape_len.saturating_sub(self.cache_start_idx)); } - fn prepare_write(&mut self, tape_len: usize) { + pub(crate) fn prepare_write(&mut self, tape_len: usize) { debug_assert!(self.len == 0 || self.cache_start_idx + self.len == tape_len); if self.len == 0 { @@ -142,130 +133,3 @@ impl RingBuffer { } } } - -/// A file writer that writes to a [`RingBuffer`] and flushes to disk periodically. -pub struct RingBufferFileWriter { - pub(crate) ring_buffer: Arc>, - pub(crate) bytes_to_flush: usize, - pub(crate) file: Arc, - pub(crate) len: u64, -} - -impl RingBufferFileWriter { - /// Flush the buffer to disk. - pub fn flush(&mut self, persistence: Persistence) -> io::Result<()> { - flush( - &self.file, - &self.ring_buffer.read(), - self.bytes_to_flush, - self.len, - persistence, - ) - } - - /// Write some data to the tape. - pub fn write(&mut self, data: &[u8]) -> io::Result { - let mut ring_buffer = self.ring_buffer.write(); - ring_buffer.prepare_write(self.len as usize); - let capacity = ring_buffer.capacity(); - - // Writing enough data to completely fill the ring buffer. - if data.len() >= capacity { - flush( - &self.file, - &ring_buffer, - self.bytes_to_flush, - self.len, - Persistence::Buffer, - )?; - self.bytes_to_flush = 0; - - write_all_at(&self.file, data, self.len)?; - - ring_buffer.push(&data[data.len() - capacity..], data.len() - capacity) - } - // Writing enough data to push data that hasn't been flushed to disk yet out of the ring buffer. - else if self.bytes_to_flush + data.len() > capacity { - // Just flush everything that needs to be flushed to disk to reduce the number of flushes. - flush( - &self.file, - &ring_buffer, - self.bytes_to_flush, - self.len, - Persistence::Buffer, - )?; - self.bytes_to_flush = 0; - - self.bytes_to_flush += data.len(); - ring_buffer.push(data, 0) - } - // Writing data that won't push data that hasn't been flushed to disk yet out of the ring buffer. - else { - self.bytes_to_flush += data.len(); - ring_buffer.push(data, 0) - } - - let old_len = self.len; - self.len += data.len() as u64; - - Ok(old_len) - } -} - -fn write_all_at(file: &File, buf: &[u8], offset: u64) -> io::Result<()> { - #[cfg(unix)] - { - use std::os::unix::fs::FileExt; - - file.write_all_at(buf, offset) - } - #[cfg(windows)] - { - use std::os::windows::fs::FileExt; - - let n = file.seek_write(buf, offset)?; - if n != buf.len() { - return Err(io::Error::other("Failed to write all bytes to tape")); - } - - Ok(()) - } -} - -fn flush( - file: &File, - ring_buffer: &RingBuffer, - mut bytes_to_flush: usize, - tape_len: u64, - persistence: Persistence, -) -> std::io::Result<()> { - if bytes_to_flush != 0 { - let (fist_slice, second_slice) = ring_buffer.as_slices(); - match bytes_to_flush.cmp(&second_slice.len()) { - Ordering::Less | Ordering::Equal => { - write_all_at( - file, - &second_slice[second_slice.len() - bytes_to_flush..], - tape_len - bytes_to_flush as u64, - )?; - } - Ordering::Greater => { - let first_slice_top_needed = bytes_to_flush - second_slice.len(); - write_all_at( - file, - &fist_slice[fist_slice.len() - first_slice_top_needed..], - tape_len - bytes_to_flush as u64, - )?; - bytes_to_flush -= first_slice_top_needed; - - write_all_at(file, second_slice, tape_len - bytes_to_flush as u64)?; - } - } - } - - match persistence { - Persistence::Buffer => Ok(()), - Persistence::SyncData => file.sync_data(), - Persistence::SyncAll => file.sync_all(), - } -} diff --git a/src/tapes/fixed_sized_iter.rs b/src/tapes/fixed_sized_iter.rs index 1dff160..5b6cef9 100644 --- a/src/tapes/fixed_sized_iter.rs +++ b/src/tapes/fixed_sized_iter.rs @@ -1,9 +1,10 @@ -use std::io; +use std::{cmp::max, io}; -use crate::{FixedSizedTape, TapesRead}; +use crate::{BlobTape, FixedSizedTape, TapesRead}; -pub struct Iter<'a, E, T: ?Sized> { - tape: &'a FixedSizedTape, +/// An iterator over the entries of a fixed-sized tape. +pub struct Iter<'a, B: BlobTape, E, T: ?Sized> { + tape: &'a FixedSizedTape, tx: &'a T, start_index: u64, @@ -13,16 +14,16 @@ pub struct Iter<'a, E, T: ?Sized> { tape_len: u64, } -impl<'a, E: bytemuck::Pod, T: TapesRead + ?Sized> Iter<'a, E, T> { +impl<'a, B: BlobTape, E: bytemuck::Pod, T: TapesRead + ?Sized> Iter<'a, B, E, T> { pub(crate) fn new( - tape: &'a FixedSizedTape, + tape: &'a FixedSizedTape, tx: &'a T, start_index: u64, tape_len: u64, ) -> io::Result { const READ_AHEAD_SIZE: usize = 8 * 1024; - let mut buf = vec![E::zeroed(); READ_AHEAD_SIZE / size_of::()]; + let mut buf = vec![E::zeroed(); max(1, READ_AHEAD_SIZE / size_of::())]; let entries_to_read = buf.len() - (start_index as usize + buf.len()).saturating_sub(tape_len as usize); @@ -42,7 +43,7 @@ impl<'a, E: bytemuck::Pod, T: TapesRead + ?Sized> Iter<'a, E, T> { } } -impl<'a, E: bytemuck::Pod, T: TapesRead + ?Sized> Iterator for Iter<'a, E, T> { +impl<'a, B: BlobTape, E: bytemuck::Pod, T: TapesRead + ?Sized> Iterator for Iter<'a, B, E, T> { type Item = io::Result; fn next(&mut self) -> Option { diff --git a/src/tapes/rolling_tape.rs b/src/tapes/rolling_tape.rs new file mode 100644 index 0000000..3f69817 --- /dev/null +++ b/src/tapes/rolling_tape.rs @@ -0,0 +1,450 @@ +use std::{ + cmp::min, + collections::VecDeque, + fs::File, + fs::OpenOptions, + io, + path::{Path, PathBuf}, + str::FromStr, + sync::Arc, +}; + +use parking_lot::RwLock; + +use crate::{ + Persistence, + io_helpers::{read_exact_at_file, write_all_at}, + metadata::TapeMetadata, + traits::{BlobTape, BlobTapeWriter, OpenConfig}, +}; + +/// Open options for a [`RollingBlobTape`]. +#[derive(Clone)] +pub struct RollingTapeOpenOptions { + /// The max size of each file in the rolling tape. + /// + /// Bigger size means less likely to need to read across multiple files when reading data, but also + /// means a longer time for old data to be deleted. + /// + /// You should keep in mind that if this is set too low and keep lots of data in the tape then lots of + /// files will be created. + pub file_size: u64, + /// The directory to store the rolling tapes. + pub dir: PathBuf, + /// The byte index to start the rolling tapes at, only used when creating a new tape. + /// + /// This can be used to start the tape indexing at a specific index when creating a new tape. + pub start_index: u64, +} + +impl OpenConfig for RollingTapeOpenOptions { + fn start_index(&self) -> u64 { + self.start_index + } +} + +/// A rolling tape, allows removing data FIFO from the tape. +/// +/// This should only be used when you need to remove data FIFO, if not use a regular tape. +pub struct RollingBlobTape { + name: &'static str, + files: Arc>>, + dir: PathBuf, + file_size: u64, +} + +impl BlobTape for RollingBlobTape { + type OpenConfig = RollingTapeOpenOptions; + type Writer = RollingBlobTapeWriter; + + fn name(&self) -> &'static str { + self.name + } + + fn open( + name: &'static str, + tape_metadata: Option, + current_epoch: u64, + config: Self::OpenConfig, + ) -> io::Result { + if config.file_size == 0 { + return Err(io::Error::other("file_size must not be 0.")); + } + + let path = config.dir.join("tapes").join(name); + + if tape_metadata.is_none() { + match std::fs::remove_dir_all(&path) { + Ok(_) => (), + Err(e) if e.kind() == io::ErrorKind::NotFound => (), + Err(e) => return Err(e), + } + } + + let mut files = Vec::new(); + + match std::fs::read_dir(&path) { + Ok(dir) => { + for entry in dir { + let entry = entry?; + let Some(Ok(index)) = entry + .path() + .file_name() + .and_then(|i| i.to_str()) + .map(u64::from_str) + else { + return Err(io::Error::other("File in rolling tapes has invalid name")); + }; + + let rolling_tape_file = RollingTapeFile::open(&path, index)?; + + files.push(rolling_tape_file); + } + } + Err(e) if e.kind() == io::ErrorKind::NotFound => { + if tape_metadata.is_some() { + return Err(io::Error::other( + "tape was in metadata but files was not found.", + )); + } + + std::fs::create_dir_all(path.clone())?; + + let first_file = RollingTapeFile::new( + &path, + offset_to_file_index(config.start_index, config.file_size), + )?; + + return Ok(Self { + name, + files: Arc::new(RwLock::new(VecDeque::from([first_file]))), + dir: path, + file_size: config.file_size, + }); + } + Err(e) => return Err(e), + } + + files.sort_by_key(|f| f.file_index); + + let files = files.into_iter().collect(); + + let this = Self { + name, + files: Arc::new(RwLock::new(files)), + dir: path, + file_size: config.file_size, + }; + + if let Some(tape_metadata) = tape_metadata { + remove_old_files( + &this.files, + &this.dir, + tape_metadata, + current_epoch, + u64::MAX, + this.file_size, + )?; + } + + Ok(this) + } + + fn read_bytes(&self, mut offset: u64, mut buf: &mut [u8]) -> io::Result<()> { + let files = self.files.read(); + + let files = files + .iter() + .filter(|f| { + let file_start = file_index_to_offset(f.file_index, self.file_size); + let file_end = file_start + self.file_size; + + file_start < offset + buf.len() as u64 && file_end > offset + }) + .cloned() + .collect::>(); + + for file in files { + let next_file_start = + file_index_to_offset(file.file_index, self.file_size) + self.file_size; + + let bytes_left_on_file = next_file_start - offset; + + let bytes_to_read = min(buf.len(), bytes_left_on_file.try_into().unwrap()); + + read_exact_at_file( + &file.file, + &mut buf[..bytes_to_read], + offset - file_index_to_offset(file.file_index, self.file_size), + )?; + + offset += u64::try_from(bytes_to_read).unwrap(); + buf = &mut buf[bytes_to_read..]; + } + + Ok(()) + } + + fn writer(&self, len: u64) -> io::Result { + let files = self.files.read(); + + let currently_writing_idx = match files + .iter() + .enumerate() + .find(|(_, tape)| file_index_to_offset(tape.file_index, self.file_size) > len) + .map(|(i, _)| i) + .unwrap_or(files.len()) + .checked_sub(1) + { + Some(i) => i, + None => { + unreachable!("Tape files must contain at least one tape file."); + } + }; + + drop(files); + + Ok(RollingBlobTapeWriter { + files: self.files.clone(), + dir: self.dir.clone(), + file_size: self.file_size, + currently_writing_idx, + first_file_touched: currently_writing_idx, + len, + is_truncation: None, + }) + } + + fn delete(self) -> io::Result<()> { + std::fs::remove_dir_all(self.dir) + } +} + +#[derive(Clone)] +struct RollingTapeFile { + file: Arc, + out_of_range_at_epoch: Option, + file_index: u64, +} + +impl RollingTapeFile { + pub fn new(dir: &Path, file_index: u64) -> io::Result { + let file_path = rolling_file_path(dir, file_index); + + let file = Arc::new( + OpenOptions::new() + .write(true) + .read(true) + .create(true) + .truncate(true) + .open(&file_path)?, + ); + + Ok(Self { + file, + out_of_range_at_epoch: None, + file_index, + }) + } + + pub fn open(dir: &Path, file_index: u64) -> io::Result { + let path = rolling_file_path(dir, file_index); + let file = OpenOptions::new().read(true).write(true).open(&path)?; + + Ok(Self { + file: Arc::new(file), + out_of_range_at_epoch: None, + file_index, + }) + } +} + +/// A writer that writes to a [`RollingBlobTape`]. +/// +/// This should not be used directly. +pub struct RollingBlobTapeWriter { + files: Arc>>, + dir: PathBuf, + file_size: u64, + currently_writing_idx: usize, + first_file_touched: usize, + len: u64, + is_truncation: Option, +} + +impl RollingBlobTapeWriter { + fn make_new_file(&self) -> io::Result { + let next_index = self + .files + .read() + .back() + .map_or_default(|f| f.file_index + 1); + + let file = RollingTapeFile::new(&self.dir, next_index)?; + + self.files.write().push_back(file.clone()); + + Ok(file) + } +} + +impl BlobTapeWriter for RollingBlobTapeWriter { + fn write_bytes(&mut self, mut buf: &[u8]) -> io::Result { + assert!(self.is_truncation.is_none_or(|x| !x)); + self.is_truncation = Some(false); + + let idx = self.len; + + while !buf.is_empty() { + let files = self.files.read(); + let tape_file = match files.get(self.currently_writing_idx) { + Some(tape_file) => { + let tape_file = tape_file.clone(); + drop(files); + tape_file + } + None => { + drop(files); + self.make_new_file()? + } + }; + + let bytes_left_on_first_tape = self.file_size + - (self.len - file_index_to_offset(tape_file.file_index, self.file_size)); + + let bytes_to_write = min(buf.len(), bytes_left_on_first_tape.try_into().unwrap()); + + write_all_at( + &tape_file.file, + &buf[..bytes_to_write], + self.len - file_index_to_offset(tape_file.file_index, self.file_size), + )?; + + self.len += u64::try_from(bytes_to_write).unwrap(); + buf = &buf[bytes_to_write..]; + + if !buf.is_empty() { + self.currently_writing_idx += 1; + } + } + + Ok(idx) + } + + fn flush(&mut self, persistence: Persistence) -> io::Result<()> { + if self.is_truncation.is_some_and(|x| x) { + let mut files = self.files.write(); + + let mut file_index = files.front().map(|f| f.file_index).unwrap_or_default(); + + let mut start_index = file_index_to_offset(file_index, self.file_size); + + while start_index > self.len { + files.push_front(RollingTapeFile::new(&self.dir, file_index - 1)?); + + file_index -= 1; + start_index -= self.file_size; + } + + for file in files.iter_mut() { + if file_index_to_offset(file.file_index, self.file_size) <= self.len + && self.len + < file_index_to_offset(file.file_index, self.file_size) + self.file_size + { + file.out_of_range_at_epoch = None; + } + } + + return Ok(()); + } + + if matches!(persistence, Persistence::Buffer) { + return Ok(()); + } + + for i in self.first_file_touched..(self.currently_writing_idx + 1) { + let file = self.files.read().get(i).unwrap().file.clone(); + + match persistence { + Persistence::Buffer => (), + Persistence::SyncData => file.sync_data()?, + Persistence::SyncAll => file.sync_all()?, + } + } + + Ok(()) + } + + fn truncate(&mut self, new_len: u64) { + assert!(self.is_truncation.is_none_or(|x| x)); + self.is_truncation = Some(true); + + self.len = new_len; + } + + fn len(&self) -> u64 { + self.len + } + + fn remove_old_files( + &self, + metadata: TapeMetadata, + current_epoch: u64, + oldest_reader_epoch: u64, + ) -> io::Result<()> { + remove_old_files( + &self.files, + &self.dir, + metadata, + current_epoch, + oldest_reader_epoch, + self.file_size, + ) + } +} + +fn remove_old_files( + files: &RwLock>, + dir: &Path, + metadata: TapeMetadata, + current_epoch: u64, + oldest_reader_epoch: u64, + file_size: u64, +) -> io::Result<()> { + let mut files = files.write(); + + let mut i = 1; + while let Some(file_2) = files.get(i) { + if file_index_to_offset(file_2.file_index, file_size) <= metadata.start { + files + .get_mut(i - 1) + .unwrap() + .out_of_range_at_epoch + .get_or_insert(current_epoch); + } + + i += 1; + } + + while let Some(file) = files.pop_front_if(|file| { + file.out_of_range_at_epoch + .is_some_and(|e| e < oldest_reader_epoch) + }) { + let path = rolling_file_path(dir, file.file_index); + std::fs::remove_file(&path)?; + } + + Ok(()) +} + +fn offset_to_file_index(offset: u64, file_size: u64) -> u64 { + offset / file_size +} + +fn file_index_to_offset(file_index: u64, file_size: u64) -> u64 { + file_index * file_size +} + +fn rolling_file_path(dir: &Path, file_index: u64) -> PathBuf { + dir.join(file_index.to_string()) +} diff --git a/src/tapes/whole_tape.rs b/src/tapes/whole_tape.rs new file mode 100644 index 0000000..271e136 --- /dev/null +++ b/src/tapes/whole_tape.rs @@ -0,0 +1,140 @@ +use std::{ + fs::{File, OpenOptions}, + io, + path::PathBuf, + sync::Arc, +}; + +use crate::{ + Persistence, + io_helpers::{read_exact_at_file, write_all_at}, + metadata::TapeMetadata, + traits::{BlobTape, BlobTapeWriter, OpenConfig}, +}; + +/// Open options for a [`WholeBlobTape`]. +#[derive(Clone)] +pub struct WholeTapeOpenOptions { + /// The directory to store the tapes. + pub dir: PathBuf, +} + +impl OpenConfig for WholeTapeOpenOptions { + fn start_index(&self) -> u64 { + 0 + } +} + +/// A handle to a whole blob tape. +pub struct WholeBlobTape { + name: &'static str, + file: Arc, + dir: PathBuf, +} + +impl BlobTape for WholeBlobTape { + type OpenConfig = WholeTapeOpenOptions; + type Writer = WholeBlobTapeWriter; + + fn name(&self) -> &'static str { + self.name + } + + fn open( + name: &'static str, + tape_metadata: Option, + _: u64, + config: Self::OpenConfig, + ) -> io::Result { + let dir = config.dir.join("tapes"); + + match OpenOptions::new() + .write(true) + .read(true) + .open(dir.join(name)) + { + Ok(file) => { + let metadata = tape_metadata.unwrap_or_default(); + + if file.metadata()?.len() < metadata.len { + return Err(io::Error::other("Tape file is too small")); + } + + let file = Arc::new(file); + + Ok(WholeBlobTape { name, file, dir }) + } + Err(e) if e.kind() == io::ErrorKind::NotFound => { + if tape_metadata.is_some() { + return Err(io::Error::other( + "tape was in metadata but file was not found.", + )); + } + + let file = Arc::new( + OpenOptions::new() + .write(true) + .read(true) + .create(true) + .truncate(true) + .open(dir.join(name))?, + ); + + Ok(WholeBlobTape { name, file, dir }) + } + Err(e) => Err(e), + } + } + + fn read_bytes(&self, offset: u64, buf: &mut [u8]) -> io::Result<()> { + read_exact_at_file(&self.file, buf, offset) + } + + fn writer(&self, len: u64) -> io::Result { + Ok(WholeBlobTapeWriter { + file: self.file.clone(), + len, + }) + } + + fn delete(self) -> io::Result<()> { + std::fs::remove_file(self.dir.join(self.name)) + } +} + +/// A writer that writes to a [`WholeBlobTape`] +/// +/// This should not be used directly. +pub struct WholeBlobTapeWriter { + file: Arc, + len: u64, +} + +impl BlobTapeWriter for WholeBlobTapeWriter { + fn flush(&mut self, persistence: Persistence) -> io::Result<()> { + match persistence { + Persistence::Buffer => Ok(()), + Persistence::SyncData => self.file.sync_data(), + Persistence::SyncAll => self.file.sync_all(), + } + } + + fn write_bytes(&mut self, buf: &[u8]) -> io::Result { + write_all_at(&self.file, buf, self.len)?; + let idx = self.len; + self.len += buf.len() as u64; + Ok(idx) + } + + fn truncate(&mut self, new_len: u64) { + self.len = new_len; + } + + fn len(&self) -> u64 { + self.len + } + + fn remove_old_files(&self, _: TapeMetadata, _: u64, _: u64) -> io::Result<()> { + Ok(()) + } +} diff --git a/src/traits.rs b/src/traits.rs index d3e911a..17578a1 100644 --- a/src/traits.rs +++ b/src/traits.rs @@ -1,29 +1,109 @@ -use std::fs::File; use std::io; -use crate::{BlobTape, FixedSizedTape}; +use crate::{FixedSizedTape, Persistence, metadata::TapeMetadata}; +/// A trait for a tape of bytes. +/// +/// You should not use any functions or types from this trait directly. You should use the transaction +/// API. +/// +/// This trait is sealed, only this crate provides implementations. +pub trait BlobTape: Sized { + type OpenConfig: OpenConfig; + + type Writer: BlobTapeWriter + 'static; + + fn name(&self) -> &'static str; + + fn open( + name: &'static str, + tape_metadata: Option, + current_epoch: u64, + config: Self::OpenConfig, + ) -> io::Result; + + fn read_bytes(&self, offset: u64, buf: &mut [u8]) -> io::Result<()>; + + fn writer(&self, len: u64) -> io::Result; + + fn delete(self) -> io::Result<()>; +} + +/// A trait for the configuration of a tape. +pub trait OpenConfig { + /// The byte index to start the tape at, only used when creating a new tape. + fn start_index(&self) -> u64; +} + +/// An internal trait for a writer to a tape. +/// +/// A single writer can only append _or_ truncate, you should not mix both in 1 writer instance. +pub trait BlobTapeWriter { + fn write_bytes(&mut self, buf: &[u8]) -> io::Result; + + fn truncate(&mut self, new_len: u64); + + fn flush(&mut self, persistence: Persistence) -> io::Result<()>; + + fn revert(&mut self, _bytes: usize) {} + + fn len(&self) -> u64; + + fn remove_old_files( + &self, + metadata: TapeMetadata, + current_epoch: u64, + oldest_reader_epoch: u64, + ) -> io::Result<()>; +} + +/// A trait for reading from tapes. pub trait TapesRead { - fn blob_tape_len(&self, tape: &BlobTape) -> Option; + /// Returns the length of a [`BlobTape`]. + /// + /// This will not take into account the removed bytes and will be the total bytes written + /// excluding those popped. + fn blob_tape_len(&self, tape: &B) -> Option; + + /// Gets the start index of a tape. + /// + /// This will be `0` for a tape that has not had its start index shifted. + fn blob_tape_start(&self, tape: &B) -> Option; - fn read_bytes(&self, blob_tape: &BlobTape, offset: u64, buf: &mut [u8]) -> io::Result<()> { + /// Fills a mutable buffer with bytes from a tape, starting at the given `offset`. + /// + /// Will return an error if the read goes past the end of the tape. + fn read_bytes( + &self, + blob_tape: &B, + offset: u64, + buf: &mut [u8], + ) -> io::Result<()> { let tape_len = self .blob_tape_len(blob_tape) .ok_or(io::Error::other("Tape not found"))?; - read_bytes(blob_tape, tape_len, offset, buf) + let tape_start = self + .blob_tape_start(blob_tape) + .ok_or(io::Error::other("Tape not found"))?; + + read_bytes(blob_tape, tape_len, tape_start, offset, buf) } - fn fixed_sized_tape_len(&self, tape: &FixedSizedTape) -> Option { + /// Gets the length of a fixed-sized tape in entries. + fn fixed_sized_tape_len( + &self, + tape: &FixedSizedTape, + ) -> Option { self.blob_tape_len(&tape.inner) .map(|bytes| bytes / size_of::() as u64) } /// Reads an entry from a fixed-sized tape. /// - /// Returns `None` if read goes past the end of the tape. - fn read_entry( + /// Returns `None` if the read goes past the end of the tape, or before the start of the tape. + fn read_entry( &self, - fixed_sized_tape: &FixedSizedTape, + fixed_sized_tape: &FixedSizedTape, index: u64, ) -> io::Result> { let mut entry = E::zeroed(); @@ -49,9 +129,9 @@ pub trait TapesRead { /// when accessing the tape file. /// /// If there is an error, the state of the buffer is not guaranteed. - fn read_entries( + fn read_entries( &self, - fixed_sized_tape: &FixedSizedTape, + fixed_sized_tape: &FixedSizedTape, offset: u64, buf: &mut [E], ) -> io::Result<()> { @@ -68,11 +148,11 @@ pub trait TapesRead { /// /// Will return an error if the start is past the end of the tape or on any other I/O error /// when accessing the tape file. - fn iter_from<'b, E: bytemuck::Pod>( + fn iter_from<'b, B: BlobTape, E: bytemuck::Pod>( &'b self, - fixed_sized_tape: &'b FixedSizedTape, + fixed_sized_tape: &'b FixedSizedTape, from: u64, - ) -> io::Result> { + ) -> io::Result> { let tape_len = self .fixed_sized_tape_len(fixed_sized_tape) .ok_or(io::Error::other("Tape not found"))?; @@ -88,127 +168,119 @@ pub trait TapesRead { } } +/// A trait for truncating and popping tapes. pub trait TapesTruncate: TapesRead { - fn truncate_blob_tape(&mut self, tape: &BlobTape, new_len: u64); + /// Truncates a blob tape to `new_len` bytes. + /// + /// Returns an error if the tape does not exist or if `new_len` is longer than the tape. + /// + /// If `new_len` is before the tape's start index, the tape is emptied and its start index + /// moves to `new_len`. + fn truncate_blob_tape(&mut self, tape: &B, new_len: u64) -> io::Result<()>; - fn truncate_fixed_sized_tape( + /// Truncates a fixed-sized tape to `new_len` entries. + /// + /// Returns an error if the tape does not exist or if `new_len` is longer than the tape. + /// + /// If `new_len` is before the tape's start index, the tape is emptied and its start index + /// moves to `new_len`. + fn truncate_fixed_sized_tape( &mut self, - tape: &FixedSizedTape, + tape: &FixedSizedTape, new_len: u64, - ) { - self.truncate_blob_tape(&tape.inner, new_len * size_of::() as u64); + ) -> io::Result<()> { + self.truncate_blob_tape(&tape.inner, new_len * size_of::() as u64) } - fn drop_fixed_sized_tape( + /// Drops the last `numb_to_drop` entries from a fixed-sized tape. + fn drop_fixed_sized_tape( &mut self, - tape: &FixedSizedTape, + tape: &FixedSizedTape, numb_to_drop: u64, - ) { + ) -> io::Result<()> { let Some(len) = self.fixed_sized_tape_len(tape) else { - return; + return Err(io::Error::other("Tape not found")); }; let new_len = len.saturating_sub(numb_to_drop); - self.truncate_fixed_sized_tape(tape, new_len); + self.truncate_fixed_sized_tape(tape, new_len) } - fn pop_fixed_sized_tape( + /// Pops the last entry from a fixed-sized tape. + /// + /// Returns the index and entry of the popped entry, or `None` if the tape does not exist or is + /// empty. + fn pop_fixed_sized_tape( &mut self, - tape: &FixedSizedTape, + tape: &FixedSizedTape, ) -> io::Result> { let Some(len) = self.fixed_sized_tape_len(tape) else { return Ok(None); }; + if len == 0 { + return Ok(None); + } + let Some(entry) = self.read_entry(tape, len - 1)? else { return Ok(None); }; - self.truncate_fixed_sized_tape(tape, len - 1); + self.truncate_fixed_sized_tape(tape, len - 1)?; Ok(Some((len - 1, entry))) } } +/// A trait for appending to tapes. pub trait TapesAppend: TapesRead { - fn append_bytes(&mut self, tape: &BlobTape, buf: &[u8]) -> io::Result; - fn append_entries( + /// Appends bytes to a tape. + /// + /// Returns the index at which the data was written. + fn append_bytes(&mut self, tape: &B, buf: &[u8]) -> io::Result; + /// Appends entries to a fixed-sized tape. + /// + /// Returns the index of the first appended entry. + fn append_entries( &mut self, - fixed_sized_tape: &FixedSizedTape, + fixed_sized_tape: &FixedSizedTape, entries: &[E], ) -> io::Result { self.append_bytes(&fixed_sized_tape.inner, bytemuck::cast_slice(entries)) .map(|len| len / size_of::() as u64) } -} - -fn read_bytes(blob_tape: &BlobTape, tape_len: u64, offset: u64, buf: &mut [u8]) -> io::Result<()> { - if tape_len < offset + buf.len() as u64 { - return Err(io::Error::new( - io::ErrorKind::UnexpectedEof, - "Read past end of tape", - )); - } - - let top_cache = blob_tape.top_cache.read(); - let cached_offset = top_cache.cache_start_idx() as u64; - - let mut last_byte_needed_offset = offset + buf.len() as u64; - if last_byte_needed_offset > cached_offset { - let read_start = offset.saturating_sub(cached_offset); - let buf_to_fill = &mut buf[(cached_offset.saturating_sub(offset)) as usize..]; - - top_cache.fill(read_start as usize, buf_to_fill); - last_byte_needed_offset -= buf_to_fill.len() as u64; - } + /// Shift the start index of a tape. + /// + /// This will do nothing if `new_start` is less than the current start of the tape, and for a whole + /// tape it will not free up disk space. + fn shift_start_idx(&mut self, tape: &B, new_start: u64) -> io::Result<()>; - if last_byte_needed_offset != offset { - read_exact_at( - &blob_tape.file, - &mut buf[0..(last_byte_needed_offset - offset) as usize], - offset, - )?; + /// Shift the start index of a fixed size tape. + /// + /// This will do nothing if `new_start` is less than the current start of the tape, and for a whole + /// tape it will not free up disk space. + fn shift_start_idx_fixed( + &mut self, + fixed_sized_tape: &FixedSizedTape, + new_start: u64, + ) -> io::Result<()> { + self.shift_start_idx(&fixed_sized_tape.inner, new_start * size_of::() as u64) } - - Ok(()) } -pub(crate) fn read_exact_at(file: &File, buf: &mut [u8], offset: u64) -> io::Result<()> { - #[cfg(unix)] - { - use std::os::unix::fs::FileExt; - - file.read_exact_at(buf, offset) +fn read_bytes( + blob_tape: &B, + tape_len: u64, + start_index: u64, + offset: u64, + buf: &mut [u8], +) -> io::Result<()> { + if tape_len < offset + buf.len() as u64 || start_index > offset { + return Err(io::Error::new( + io::ErrorKind::UnexpectedEof, + "Read out of bounds", + )); } - #[cfg(windows)] - { - use std::os::windows::fs::FileExt; - - let mut buf = buf; - let mut offset = offset; - while !buf.is_empty() { - match file.seek_read(buf, offset) { - Ok(0) => { - break; - } - Ok(n) => { - buf = &mut buf[n..]; - offset += n as u64; - } - Err(e) => { - return Err(e); - } - } - } - - if !buf.is_empty() { - Err(io::Error::new( - io::ErrorKind::UnexpectedEof, - "failed to fill the whole buffer", - )) - } else { - Ok(()) - } - } + blob_tape.read_bytes(offset, buf) } diff --git a/tests/coverage_gaps.rs b/tests/coverage_gaps.rs new file mode 100644 index 0000000..7bbaac5 --- /dev/null +++ b/tests/coverage_gaps.rs @@ -0,0 +1,433 @@ +use std::io; + +use tapes::{ + CachedBlobTape, CachedTapeOpenOptions, FixedSizedTape, Persistence, RollingBlobTape, + RollingTapeOpenOptions, Tapes, TapesAppend, TapesRead, WholeBlobTape, WholeTapeOpenOptions, +}; + +/// Writes `len` bytes of a deterministic pattern. +fn pattern(len: usize) -> Vec { + (0..len).map(|i| b'a' + (i % 26) as u8).collect() +} + +fn whole_options(dir: &tempfile::TempDir) -> WholeTapeOpenOptions { + WholeTapeOpenOptions { + dir: dir.path().to_path_buf(), + } +} + +/// Two blob tapes in one database must be independent. +#[test] +fn multiple_blob_tapes_isolated() { + let dir = tempfile::tempdir().unwrap(); + let tapes = Tapes::open(dir.path()).unwrap(); + let options = whole_options(&dir); + + let (a, b) = { + let mut append = tapes.append(); + let a = append + .open_blob_tape::("tape_a", options.clone()) + .unwrap(); + let b = append + .open_blob_tape::("tape_b", options.clone()) + .unwrap(); + append.append_bytes(&a, b"hello").unwrap(); + append.append_bytes(&b, b"world!!").unwrap(); + append.commit(Persistence::Buffer).unwrap(); + (a, b) + }; + + let reader = tapes.reader(); + assert_eq!(reader.blob_tape_len(&a), Some(5)); + assert_eq!(reader.blob_tape_len(&b), Some(7)); + + let mut a_data = [0u8; 5]; + reader.read_bytes(&a, 0, &mut a_data).unwrap(); + assert_eq!(&a_data, b"hello"); + + let mut b_data = [0u8; 7]; + reader.read_bytes(&b, 0, &mut b_data).unwrap(); + assert_eq!(&b_data, b"world!!"); + + // Appending to one must not affect the other. + { + let mut append = tapes.append(); + append.append_bytes(&a, b"!").unwrap(); + append.commit(Persistence::Buffer).unwrap(); + } + + let reader = tapes.reader(); + assert_eq!(reader.blob_tape_len(&a), Some(6)); + assert_eq!(reader.blob_tape_len(&b), Some(7)); +} + +/// Two fixed-sized tapes in one database must be independent. +#[test] +fn multiple_fixed_tapes_isolated() { + let dir = tempfile::tempdir().unwrap(); + let tapes = Tapes::open(dir.path()).unwrap(); + let options = whole_options(&dir); + + let (a, b) = { + let mut append = tapes.append(); + let a = append + .open_fixed_sized_tape::("tape_a", options.clone()) + .unwrap(); + let b = append + .open_fixed_sized_tape::("tape_b", options.clone()) + .unwrap(); + append.append_entries(&a, &[1, 2, 3]).unwrap(); + append.append_entries(&b, &[10, 20]).unwrap(); + append.commit(Persistence::Buffer).unwrap(); + (a, b) + }; + + let reader = tapes.reader(); + assert_eq!(reader.fixed_sized_tape_len(&a), Some(3)); + assert_eq!(reader.fixed_sized_tape_len(&b), Some(2)); + + let mut a_buf = [0u64; 3]; + reader.read_entries(&a, 0, &mut a_buf).unwrap(); + assert_eq!(a_buf, [1, 2, 3]); + + let mut b_buf = [0u64; 2]; + reader.read_entries(&b, 0, &mut b_buf).unwrap(); + assert_eq!(b_buf, [10, 20]); +} + +/// `shift_start_idx_fixed` rolls by entry index (multiplies by `size_of::()`). +#[test] +fn fixed_shift_start_by_entry_index() { + let dir = tempfile::tempdir().unwrap(); + let tapes = Tapes::open(dir.path()).unwrap(); + let options = RollingTapeOpenOptions { + file_size: 8, + dir: dir.path().to_path_buf(), + start_index: 0, + }; + let entries: Vec = vec![10, 20, 30, 40, 50, 60]; + + let tape: FixedSizedTape = { + let mut append = tapes.append(); + let tape = append + .open_fixed_sized_tape::("tape", options.clone()) + .unwrap(); + append.append_entries(&tape, &entries).unwrap(); + append.commit(Persistence::Buffer).unwrap(); + tape + }; + + // A separate blob handle (same tape) for the start check. + let blob: RollingBlobTape = { + let mut append = tapes.append(); + append + .open_blob_tape::("tape", options) + .unwrap() + }; + + // Roll to entry index 2 (byte offset 16), popping entries 0,1. + { + let mut append = tapes.append(); + append.shift_start_idx_fixed(&tape, 2).unwrap(); + append.commit(Persistence::Buffer).unwrap(); + } + + let reader = tapes.reader(); + assert_eq!(reader.blob_tape_start(&blob), Some(16)); + + // Live entries are indices 2..6. + let mut buf = [0u64; 4]; + reader.read_entries(&tape, 2, &mut buf).unwrap(); + assert_eq!(buf, [30, 40, 50, 60]); + + // Reading a popped entry fails cleanly. + let mut err = [0u64; 1]; + let e = reader.read_entries(&tape, 0, &mut err).unwrap_err(); + assert_eq!(e.kind(), io::ErrorKind::UnexpectedEof); +} + +/// Rolling to a start offset that is not a multiple of `file_size` must still +/// read the live region byte-precisely. +#[test] +fn rolling_non_aligned_start() { + let dir = tempfile::tempdir().unwrap(); + let tapes = Tapes::open(dir.path()).unwrap(); + let options = RollingTapeOpenOptions { + file_size: 8, + dir: dir.path().to_path_buf(), + start_index: 0, + }; + let data = pattern(40); + + let tape = { + let mut append = tapes.append(); + let tape = append + .open_blob_tape::("tape", options.clone()) + .unwrap(); + append.append_bytes(&tape, &data).unwrap(); + append.commit(Persistence::Buffer).unwrap(); + tape + }; + + // Roll to start = 10 (not a multiple of file_size=8). + { + let mut append = tapes.append(); + append.shift_start_idx(&tape, 10).unwrap(); + append.commit(Persistence::Buffer).unwrap(); + } + + let reader = tapes.reader(); + assert_eq!(reader.blob_tape_start(&tape), Some(10)); + assert_eq!(reader.blob_tape_len(&tape), Some(40)); + + // Live region is [10, 40). + let mut contents = [0; 30]; + reader.read_bytes(&tape, 10, &mut contents).unwrap(); + assert_eq!(&contents, &data[10..40]); + + // Reading before the start fails cleanly. + let mut err_buf = [0; 4]; + let err = reader.read_bytes(&tape, 0, &mut err_buf).unwrap_err(); + assert_eq!(err.kind(), io::ErrorKind::UnexpectedEof); +} + +/// A `CachedBlobTape` wrapping a `RollingBlobTape` must write and read correctly. +#[test] +fn cached_rolling_combo() { + let dir = tempfile::tempdir().unwrap(); + let tapes = Tapes::open(dir.path()).unwrap(); + let options = CachedTapeOpenOptions:: { + inner: RollingTapeOpenOptions { + file_size: 8, + dir: dir.path().to_path_buf(), + start_index: 0, + }, + top_cache_size: 16, + }; + let data = pattern(40); + + let tape: CachedBlobTape = { + let mut append = tapes.append(); + let tape = append + .open_blob_tape::>("tape", options.clone()) + .unwrap(); + append.append_bytes(&tape, &data).unwrap(); + append.commit(Persistence::Buffer).unwrap(); + tape + }; + + let reader = tapes.reader(); + assert_eq!(reader.blob_tape_len(&tape), Some(40)); + + // Full read (spans the cache boundary at 24). + let mut contents = [0; 40]; + reader.read_bytes(&tape, 0, &mut contents).unwrap(); + assert_eq!(&contents, &data[..]); + + // A sub-range entirely inside the cache. + let mut in_cache = [0; 4]; + reader.read_bytes(&tape, 26, &mut in_cache).unwrap(); + assert_eq!(&in_cache, &data[26..30]); + + // A sub-range that straddles the cache boundary. + let mut straddle = [0; 8]; + reader.read_bytes(&tape, 22, &mut straddle).unwrap(); + assert_eq!(&straddle, &data[22..30]); +} + +/// Multiple concurrent readers must all see the committed data. +#[test] +fn concurrent_readers() { + let dir = tempfile::tempdir().unwrap(); + let data = pattern(200); + + { + let tapes = Tapes::open(dir.path()).unwrap(); + let mut append = tapes.append(); + let tape = append + .open_blob_tape::("tape", whole_options(&dir)) + .unwrap(); + append.append_bytes(&tape, &data).unwrap(); + append.commit(Persistence::SyncData).unwrap(); + } + + let dir_path = dir.path().to_path_buf(); + let data_clone = data.clone(); + let mut handles = Vec::new(); + for _ in 0..4 { + let data = data_clone.clone(); + let dir_path = dir_path.clone(); + handles.push(std::thread::spawn(move || { + let tapes = Tapes::open(&dir_path).unwrap(); + let tape = { + let mut append = tapes.append(); + let tape = append + .open_blob_tape::( + "tape", + WholeTapeOpenOptions { + dir: dir_path.clone(), + }, + ) + .unwrap(); + append.commit(Persistence::Buffer).unwrap(); + tape + }; + let reader = tapes.reader(); + let mut contents = [0; 200]; + reader.read_bytes(&tape, 0, &mut contents).unwrap(); + assert_eq!(&contents, &data[..]); + })); + } + + for h in handles { + h.join().unwrap(); + } +} + +/// Opening a blob tape whose committed length is not a multiple of the entry +/// size as a fixed-sized tape must error. +#[test] +fn fixed_size_mismatch_errors() { + let dir = tempfile::tempdir().unwrap(); + let tapes = Tapes::open(dir.path()).unwrap(); + let options = whole_options(&dir); + + // Create a blob tape with 5 bytes (not a multiple of 8). + { + let mut append = tapes.append(); + let tape = append + .open_blob_tape::("tape", options.clone()) + .unwrap(); + append.append_bytes(&tape, b"abcde").unwrap(); + append.commit(Persistence::Buffer).unwrap(); + } + + let result = { + let mut append = tapes.append(); + append.open_fixed_sized_tape::("tape", options) + }; + assert!(result.is_err(), "expected a size-mismatch error, got Ok"); +} + +/// `iter_from` from exactly the end yields an empty iterator; past the end errors. +#[test] +fn iter_from_bounds() { + let dir = tempfile::tempdir().unwrap(); + let tapes = Tapes::open(dir.path()).unwrap(); + let options = whole_options(&dir); + + let tape: FixedSizedTape = { + let mut append = tapes.append(); + let tape = append + .open_fixed_sized_tape::("tape", options.clone()) + .unwrap(); + append.append_entries(&tape, &[1, 2, 3]).unwrap(); + append.commit(Persistence::Buffer).unwrap(); + tape + }; + + let reader = tapes.reader(); + + // `from == len` → empty iterator, no error. + let collected: Vec = reader + .iter_from(&tape, 3) + .unwrap() + .map(|r| r.unwrap()) + .collect(); + assert!(collected.is_empty()); + + // `from > len` → error. + let result = reader.iter_from(&tape, 4); + assert!(result.is_err(), "expected a past-end error"); +} + +/// A dropped append to a rolling tape must not corrupt the committed data. +/// +/// (The rolling writer's `revert` is a no-op, so this verifies the stale bytes +/// written to disk are never observed.) +#[test] +fn dropped_rolling_append_does_not_corrupt() { + let dir = tempfile::tempdir().unwrap(); + let tapes = Tapes::open(dir.path()).unwrap(); + let options = RollingTapeOpenOptions { + file_size: 8, + dir: dir.path().to_path_buf(), + start_index: 0, + }; + let data = pattern(40); + + let tape = { + let mut append = tapes.append(); + let tape = append + .open_blob_tape::("tape", options.clone()) + .unwrap(); + append.append_bytes(&tape, &data).unwrap(); + append.commit(Persistence::Buffer).unwrap(); + tape + }; + + // A dropped append (no commit) writes stale bytes past the committed length. + { + let mut append = tapes.append(); + let tape = append + .open_blob_tape::("tape", options.clone()) + .unwrap(); + append.append_bytes(&tape, b"STALE").unwrap(); + // `append` is dropped here WITHOUT commit. + } + + // The committed data must be intact. + let reader = tapes.reader(); + assert_eq!(reader.blob_tape_len(&tape), Some(40)); + + let mut contents = [0; 40]; + reader.read_bytes(&tape, 0, &mut contents).unwrap(); + assert_eq!(&contents, &data[..]); +} + +/// Multiple appends to the same tape within one transaction. +#[test] +fn multiple_appends_same_transaction() { + let dir = tempfile::tempdir().unwrap(); + let tapes = Tapes::open(dir.path()).unwrap(); + let options = whole_options(&dir); + + let tape = { + let mut append = tapes.append(); + let tape = append + .open_blob_tape::("tape", options.clone()) + .unwrap(); + append.append_bytes(&tape, b"abc").unwrap(); + append.append_bytes(&tape, b"def").unwrap(); + append.append_bytes(&tape, b"ghi").unwrap(); + append.commit(Persistence::Buffer).unwrap(); + tape + }; + + let reader = tapes.reader(); + assert_eq!(reader.blob_tape_len(&tape), Some(9)); + + let mut contents = [0; 9]; + reader.read_bytes(&tape, 0, &mut contents).unwrap(); + assert_eq!(&contents, b"abcdefghi"); +} + +/// `tape_exists` distinguishes existing from missing tapes. +#[test] +fn tape_exist_distinguishes_missing() { + let dir = tempfile::tempdir().unwrap(); + let tapes = Tapes::open(dir.path()).unwrap(); + + { + let mut append = tapes.append(); + let tape = append + .open_blob_tape::("existing", whole_options(&dir)) + .unwrap(); + append.append_bytes(&tape, b"hi").unwrap(); + append.commit(Persistence::Buffer).unwrap(); + } + + let append = tapes.append(); + assert!(append.tape_exists("existing")); + assert!(!append.tape_exists("nonexistent")); +} diff --git a/tests/delete_tape.rs b/tests/delete_tape.rs new file mode 100644 index 0000000..42e5cd5 --- /dev/null +++ b/tests/delete_tape.rs @@ -0,0 +1,39 @@ +use tapes::{Persistence, Tapes, TapesAppend, TapesRead, WholeBlobTape, WholeTapeOpenOptions}; + +#[test] +fn delete_blob_tape() { + let dir = tempfile::tempdir().unwrap(); + let tapes = Tapes::open(dir.path()).unwrap(); + + let options = WholeTapeOpenOptions { + dir: dir.path().to_path_buf(), + }; + + let mut append = tapes.append(); + let tape = append + .open_blob_tape::("blob", options) + .unwrap(); + append.append_bytes(&tape, b"contents").unwrap(); + append.commit(Persistence::Buffer).unwrap(); + + { + let reader = tapes.reader(); + assert_eq!(reader.blob_tape_len(&tape), Some(8)); + } + + let tape_dir = dir.path().join("tapes").join("blob"); + assert!(tape_dir.exists()); + + { + let mut append = tapes.append(); + append.delete_tape(tape); + append.commit(Persistence::Buffer).unwrap(); + } + + assert!(!tape_dir.exists()); + + { + let append = tapes.append(); + assert!(!append.tape_exists("blob")); + } +} diff --git a/tests/dropped_transaction.rs b/tests/dropped_transaction.rs new file mode 100644 index 0000000..1118b9a --- /dev/null +++ b/tests/dropped_transaction.rs @@ -0,0 +1,140 @@ +use tapes::{ + Persistence, Tapes, TapesAppend, TapesRead, TapesTruncate, WholeBlobTape, WholeTapeOpenOptions, +}; + +#[test] +fn dropped_append_is_overwritten_by_the_next_append() { + let dir = tempfile::tempdir().unwrap(); + let tapes = Tapes::open(dir.path()).unwrap(); + let options = WholeTapeOpenOptions { + dir: dir.path().to_path_buf(), + }; + + let mut append = tapes.append(); + let tape = append + .open_blob_tape::("tape", options) + .unwrap(); + append.append_bytes(&tape, b"abcde").unwrap(); + append.commit(Persistence::Buffer).unwrap(); + + { + let mut append = tapes.append(); + assert_eq!(append.append_bytes(&tape, b"X").unwrap(), 5); + } + + let mut append = tapes.append(); + assert_eq!(append.append_bytes(&tape, b"Y").unwrap(), 5); + append.commit(Persistence::Buffer).unwrap(); + + let reader = tapes.reader(); + let mut contents = [0; 6]; + reader.read_bytes(&tape, 0, &mut contents).unwrap(); + assert_eq!(&contents, b"abcdeY"); +} + +#[test] +fn dropped_large_append_is_not_visible() { + let dir = tempfile::tempdir().unwrap(); + let tapes = Tapes::open(dir.path()).unwrap(); + let options = WholeTapeOpenOptions { + dir: dir.path().to_path_buf(), + }; + + let mut append = tapes.append(); + let tape = append + .open_blob_tape::("tape", options) + .unwrap(); + append.append_bytes(&tape, b"abcdefgh").unwrap(); + append.commit(Persistence::Buffer).unwrap(); + + { + let mut append = tapes.append(); + assert_eq!(append.append_bytes(&tape, b"stale").unwrap(), 8); + } + + let reader = tapes.reader(); + let mut original_contents = [0; 8]; + reader.read_bytes(&tape, 0, &mut original_contents).unwrap(); + assert_eq!(&original_contents, b"abcdefgh"); + drop(reader); + + let mut append = tapes.append(); + assert_eq!(append.append_bytes(&tape, b"XY").unwrap(), 8); + append.commit(Persistence::Buffer).unwrap(); + + let reader = tapes.reader(); + let mut contents = [0; 10]; + reader.read_bytes(&tape, 0, &mut contents).unwrap(); + assert_eq!(&contents, b"abcdefghXY"); +} + +#[test] +fn dropped_truncate_does_not_change_the_tape() { + let dir = tempfile::tempdir().unwrap(); + let tapes = Tapes::open(dir.path()).unwrap(); + let options = WholeTapeOpenOptions { + dir: dir.path().to_path_buf(), + }; + + let mut append = tapes.append(); + let tape = append + .open_blob_tape::("tape", options) + .unwrap(); + append.append_bytes(&tape, b"abcdefgh").unwrap(); + append.commit(Persistence::Buffer).unwrap(); + + { + let mut truncate = tapes.truncate(); + truncate.truncate_blob_tape(&tape, 6).unwrap(); + truncate.truncate_blob_tape(&tape, 4).unwrap(); + + let mut contents = [0; 4]; + truncate.read_bytes(&tape, 0, &mut contents).unwrap(); + assert_eq!(&contents, b"abcd"); + } + + let reader = tapes.reader(); + let mut contents = [0; 8]; + reader.read_bytes(&tape, 0, &mut contents).unwrap(); + assert_eq!(&contents, b"abcdefgh"); +} + +#[test] +fn append_after_committed_truncate_overwrites_the_suffix() { + let dir = tempfile::tempdir().unwrap(); + let tapes = Tapes::open(dir.path()).unwrap(); + let options = WholeTapeOpenOptions { + dir: dir.path().to_path_buf(), + }; + + let mut append = tapes.append(); + let tape = append + .open_blob_tape::("tape", options) + .unwrap(); + append.append_bytes(&tape, b"abcdefgh").unwrap(); + append.commit(Persistence::Buffer).unwrap(); + + { + let mut truncate = tapes.truncate(); + truncate.truncate_blob_tape(&tape, 6).unwrap(); + truncate.truncate_blob_tape(&tape, 4).unwrap(); + + truncate.commit(Persistence::Buffer).unwrap(); + } + + let reader = tapes.reader(); + assert_eq!(reader.blob_tape_len(&tape), Some(4)); + let mut contents = [0; 4]; + reader.read_bytes(&tape, 0, &mut contents).unwrap(); + assert_eq!(&contents, b"abcd"); + drop(reader); + + let mut append = tapes.append(); + assert_eq!(append.append_bytes(&tape, b"XY").unwrap(), 4); + append.commit(Persistence::Buffer).unwrap(); + + let reader = tapes.reader(); + let mut contents = [0; 6]; + reader.read_bytes(&tape, 0, &mut contents).unwrap(); + assert_eq!(&contents, b"abcdXY"); +} diff --git a/tests/rolling_tape.rs b/tests/rolling_tape.rs new file mode 100644 index 0000000..dba85a9 --- /dev/null +++ b/tests/rolling_tape.rs @@ -0,0 +1,486 @@ +use std::io; + +use tapes::{ + FixedSizedTape, Persistence, RollingBlobTape, RollingTapeOpenOptions, Tapes, TapesAppend, + TapesRead, TapesTruncate, +}; + +const NAME: &str = "tape"; + +fn options(dir: &tempfile::TempDir, file_size: u64, start_index: u64) -> RollingTapeOpenOptions { + RollingTapeOpenOptions { + file_size, + dir: dir.path().to_path_buf(), + start_index, + } +} + +fn tape_dir(dir: &tempfile::TempDir) -> std::path::PathBuf { + dir.path().join("tapes").join(NAME) +} + +/// Writes `len` bytes of a deterministic pattern. +fn pattern(len: usize) -> Vec { + (0..len).map(|i| b'a' + (i % 26) as u8).collect() +} + +#[test] +fn rolling_basic_write_read() { + let dir = tempfile::tempdir().unwrap(); + let tapes = Tapes::open(dir.path()).unwrap(); + let options = options(&dir, 8, 0); + let data = pattern(40); + + let tape = { + let mut append = tapes.append(); + let tape = append + .open_blob_tape::(NAME, options) + .unwrap(); + append.append_bytes(&tape, &data).unwrap(); + append.commit(Persistence::Buffer).unwrap(); + tape + }; + + let reader = tapes.reader(); + assert_eq!(reader.blob_tape_len(&tape), Some(40)); + assert_eq!(reader.blob_tape_start(&tape), Some(0)); + let mut contents = [0; 40]; + reader.read_bytes(&tape, 0, &mut contents).unwrap(); + assert_eq!(&contents, &data[..]); +} + +#[test] +fn rolling_read_spans_files() { + let dir = tempfile::tempdir().unwrap(); + let tapes = Tapes::open(dir.path()).unwrap(); + // file_size = 4 → 10 bytes span 3 files. + let options = options(&dir, 4, 0); + let data = pattern(10); + + let tape = { + let mut append = tapes.append(); + let tape = append + .open_blob_tape::(NAME, options) + .unwrap(); + append.append_bytes(&tape, &data).unwrap(); + append.commit(Persistence::Buffer).unwrap(); + tape + }; + + let reader = tapes.reader(); + + // Full read. + let mut contents = [0; 10]; + reader.read_bytes(&tape, 0, &mut contents).unwrap(); + assert_eq!(&contents, &data[..]); + + // A sub-range that spans the file boundaries at 4 and 8. + let mut mid = [0; 6]; + reader.read_bytes(&tape, 3, &mut mid).unwrap(); + assert_eq!(&mid, &data[3..9]); +} + +#[test] +fn rolling_roll_and_read() { + let dir = tempfile::tempdir().unwrap(); + let tapes = Tapes::open(dir.path()).unwrap(); + let options = options(&dir, 8, 0); + let data = pattern(40); + + let tape = { + let mut append = tapes.append(); + let tape = append + .open_blob_tape::(NAME, options) + .unwrap(); + append.append_bytes(&tape, &data).unwrap(); + append.commit(Persistence::Buffer).unwrap(); + tape + }; + + // Roll to start = 16. + { + let mut append = tapes.append(); + append.shift_start_idx(&tape, 16).unwrap(); + append.commit(Persistence::Buffer).unwrap(); + } + + let reader = tapes.reader(); + assert_eq!(reader.blob_tape_len(&tape), Some(40)); + assert_eq!(reader.blob_tape_start(&tape), Some(16)); + + // The live region is [16, 40). + let mut contents = [0; 24]; + reader.read_bytes(&tape, 16, &mut contents).unwrap(); + assert_eq!(&contents, &data[16..40]); + + // Reading before the start must fail cleanly. + let mut err_buf = [0; 4]; + let err = reader.read_bytes(&tape, 0, &mut err_buf).unwrap_err(); + assert_eq!(err.kind(), io::ErrorKind::UnexpectedEof); +} + +#[test] +fn rolling_roll_deletes_old_files() { + let dir = tempfile::tempdir().unwrap(); + let tapes = Tapes::open(dir.path()).unwrap(); + // file_size = 8, write 40 bytes → files 0,1,2,3,4. + let options = options(&dir, 8, 0); + let data = pattern(40); + + let tape = { + let mut append = tapes.append(); + let tape = append + .open_blob_tape::(NAME, options) + .unwrap(); + append.append_bytes(&tape, &data).unwrap(); + append.commit(Persistence::Buffer).unwrap(); + tape + }; + + // Before rolling, all 5 files exist. + let tdir = tape_dir(&dir); + for i in 0..5 { + assert!( + tdir.join(i.to_string()).exists(), + "file {i} should exist before roll" + ); + } + + // Roll to start = 16 → files 0 and 1 ([0,8),[8,16)) are fully before the start. + { + let mut append = tapes.append(); + append.shift_start_idx(&tape, 16).unwrap(); + append.commit(Persistence::Buffer).unwrap(); + } + + // Files 0 and 1 should be gone; 2,3,4 should remain. + assert!(!tdir.join("0").exists(), "file 0 should be deleted"); + assert!(!tdir.join("1").exists(), "file 1 should be deleted"); + assert!(tdir.join("2").exists(), "file 2 should remain"); + assert!(tdir.join("3").exists(), "file 3 should remain"); + assert!(tdir.join("4").exists(), "file 4 should remain"); +} + +#[test] +fn rolling_roll_then_append() { + let dir = tempfile::tempdir().unwrap(); + let tapes = Tapes::open(dir.path()).unwrap(); + let options = options(&dir, 8, 0); + let data = pattern(40); + + let tape = { + let mut append = tapes.append(); + let tape = append + .open_blob_tape::(NAME, options) + .unwrap(); + append.append_bytes(&tape, &data).unwrap(); + append.commit(Persistence::Buffer).unwrap(); + tape + }; + + { + let mut append = tapes.append(); + append.shift_start_idx(&tape, 16).unwrap(); + append.commit(Persistence::Buffer).unwrap(); + } + + // Append 2 more bytes at offset 40. + { + let mut append = tapes.append(); + append.append_bytes(&tape, b"XY").unwrap(); + append.commit(Persistence::Buffer).unwrap(); + } + + let reader = tapes.reader(); + assert_eq!(reader.blob_tape_len(&tape), Some(42)); + assert_eq!(reader.blob_tape_start(&tape), Some(16)); + + // Live region is [16, 42): old [16,40) + new "XY". + let mut contents = [0; 26]; + reader.read_bytes(&tape, 16, &mut contents).unwrap(); + let mut expected = data[16..40].to_vec(); + expected.extend_from_slice(b"XY"); + assert_eq!(&contents, &expected[..]); +} + +#[test] +fn rolling_roll_then_truncate() { + let dir = tempfile::tempdir().unwrap(); + let tapes = Tapes::open(dir.path()).unwrap(); + let options = options(&dir, 8, 0); + let data = pattern(40); + + let tape = { + let mut append = tapes.append(); + let tape = append + .open_blob_tape::(NAME, options) + .unwrap(); + append.append_bytes(&tape, &data).unwrap(); + append.commit(Persistence::Buffer).unwrap(); + tape + }; + + { + let mut append = tapes.append(); + append.shift_start_idx(&tape, 16).unwrap(); + append.commit(Persistence::Buffer).unwrap(); + } + + // Truncate the tail down to len = 24 → live region [16, 24). + { + let mut truncate = tapes.truncate(); + truncate.truncate_blob_tape(&tape, 24).unwrap(); + truncate.commit(Persistence::Buffer).unwrap(); + } + + let reader = tapes.reader(); + assert_eq!(reader.blob_tape_len(&tape), Some(24)); + assert_eq!(reader.blob_tape_start(&tape), Some(16)); + + let mut contents = [0; 8]; + reader.read_bytes(&tape, 16, &mut contents).unwrap(); + assert_eq!(&contents, &data[16..24]); +} + +#[test] +fn rolling_reopen_after_roll() { + let dir = tempfile::tempdir().unwrap(); + let options = options(&dir, 8, 0); + let data = pattern(40); + + { + let tapes = Tapes::open(dir.path()).unwrap(); + let tape = { + let mut append = tapes.append(); + let tape = append + .open_blob_tape::(NAME, options.clone()) + .unwrap(); + append.append_bytes(&tape, &data).unwrap(); + append.commit(Persistence::Buffer).unwrap(); + tape + }; + + let mut append = tapes.append(); + append.shift_start_idx(&tape, 16).unwrap(); + append.commit(Persistence::Buffer).unwrap(); + } + + // Reopen the database and verify the live region survived. + let tapes = Tapes::open(dir.path()).unwrap(); + let tape = { + let mut append = tapes.append(); + append + .open_blob_tape::(NAME, options) + .unwrap() + }; + + let reader = tapes.reader(); + assert_eq!(reader.blob_tape_len(&tape), Some(40)); + assert_eq!(reader.blob_tape_start(&tape), Some(16)); + + let mut contents = [0; 24]; + reader.read_bytes(&tape, 16, &mut contents).unwrap(); + assert_eq!(&contents, &data[16..40]); + + // Old files should still be gone after reopen. + let tdir = tape_dir(&dir); + assert!(!tdir.join("0").exists()); + assert!(!tdir.join("1").exists()); + assert!(tdir.join("2").exists()); +} + +#[test] +fn rolling_delete() { + let dir = tempfile::tempdir().unwrap(); + let tapes = Tapes::open(dir.path()).unwrap(); + let options = options(&dir, 8, 0); + + let tape = { + let mut append = tapes.append(); + let tape = append + .open_blob_tape::(NAME, options) + .unwrap(); + append.append_bytes(&tape, &pattern(20)).unwrap(); + append.commit(Persistence::Buffer).unwrap(); + tape + }; + + assert!(tape_dir(&dir).exists()); + + { + let mut append = tapes.append(); + append.delete_tape(tape); + append.commit(Persistence::Buffer).unwrap(); + } + + assert!(!tape_dir(&dir).exists()); + { + let append = tapes.append(); + assert!(!append.tape_exists(NAME)); + } +} + +#[test] +fn rolling_fixed_write_read() { + let dir = tempfile::tempdir().unwrap(); + let tapes = Tapes::open(dir.path()).unwrap(); + // file_size = 8 → each file holds exactly one u64 entry. + let options = options(&dir, 8, 0); + let entries: Vec = vec![10, 20, 30, 40, 50]; + + let tape: FixedSizedTape = { + let mut append = tapes.append(); + let tape = append + .open_fixed_sized_tape::(NAME, options) + .unwrap(); + append.append_entries(&tape, &entries).unwrap(); + append.commit(Persistence::Buffer).unwrap(); + tape + }; + + let reader = tapes.reader(); + assert_eq!( + reader.fixed_sized_tape_len(&tape), + Some(entries.len() as u64) + ); + + let mut buf = [0u64; 5]; + reader.read_entries(&tape, 0, &mut buf).unwrap(); + assert_eq!(buf, entries.as_slice()); + + // Read a sub-range of entries. + let mut sub = [0u64; 2]; + reader.read_entries(&tape, 2, &mut sub).unwrap(); + assert_eq!(sub, [30, 40]); + + // Iterator over the whole tape. + let collected: Vec = reader + .iter_from(&tape, 0) + .unwrap() + .map(|r| r.unwrap()) + .collect(); + assert_eq!(collected, entries); +} + +#[test] +fn rolling_fixed_roll_read() { + let dir = tempfile::tempdir().unwrap(); + let tapes = Tapes::open(dir.path()).unwrap(); + let options = options(&dir, 8, 0); + let entries: Vec = vec![10, 20, 30, 40, 50, 60]; + + let tape: FixedSizedTape = { + let mut append = tapes.append(); + let tape = append + .open_fixed_sized_tape::(NAME, options.clone()) + .unwrap(); + append.append_entries(&tape, &entries).unwrap(); + append.commit(Persistence::Buffer).unwrap(); + tape + }; + + // A separate blob handle (same tape) for the roll and start check. + let blob: RollingBlobTape = { + let mut append = tapes.append(); + append + .open_blob_tape::(NAME, options) + .unwrap() + }; + + // Roll to start = 16 → entries 0,1 (bytes [0,16)) are popped. + { + let mut append = tapes.append(); + append.shift_start_idx(&blob, 16).unwrap(); + append.commit(Persistence::Buffer).unwrap(); + } + + let reader = tapes.reader(); + assert_eq!(reader.blob_tape_start(&blob), Some(16)); + + // Live entries are indices 2..6. + let mut buf = [0u64; 4]; + reader.read_entries(&tape, 2, &mut buf).unwrap(); + assert_eq!(buf, [30, 40, 50, 60]); + + // Reading popped entries fails cleanly. + let mut err = [0u64; 1]; + let e = reader.read_entries(&tape, 0, &mut err).unwrap_err(); + assert_eq!(e.kind(), io::ErrorKind::UnexpectedEof); +} + +#[test] +fn rolling_truncate_below_start_empties_tape() { + let dir = tempfile::tempdir().unwrap(); + let tapes = Tapes::open(dir.path()).unwrap(); + let options = options(&dir, 8, 0); + let data = pattern(40); + + let tape = { + let mut append = tapes.append(); + let tape = append + .open_blob_tape::(NAME, options.clone()) + .unwrap(); + append.append_bytes(&tape, &data).unwrap(); + append.commit(Persistence::Buffer).unwrap(); + tape + }; + + // Roll to start = 16. + { + let mut append = tapes.append(); + append.shift_start_idx(&tape, 16).unwrap(); + append.commit(Persistence::Buffer).unwrap(); + } + + // Truncating below the start index is allowed and empties the tape, + // moving the start index down to the new length. + { + let mut truncate = tapes.truncate(); + truncate.truncate_blob_tape(&tape, 10).unwrap(); + truncate.commit(Persistence::Buffer).unwrap(); + } + + let reader = tapes.reader(); + assert_eq!(reader.blob_tape_len(&tape), Some(10)); + assert_eq!(reader.blob_tape_start(&tape), Some(10)); + + // Appending after such a truncation continues from the new length. + { + let mut append = tapes.append(); + append.append_bytes(&tape, &pattern(8)).unwrap(); + append.commit(Persistence::Buffer).unwrap(); + } + + let reader = tapes.reader(); + assert_eq!(reader.blob_tape_len(&tape), Some(18)); + assert_eq!(reader.blob_tape_start(&tape), Some(10)); + + let mut contents = vec![0u8; 8]; + reader.read_bytes(&tape, 10, &mut contents).unwrap(); + assert_eq!(contents, pattern(8)); +} + +#[test] +fn rolling_new_tape_start_index() { + let dir = tempfile::tempdir().unwrap(); + let tapes = Tapes::open(dir.path()).unwrap(); + let options = options(&dir, 8, 16); + + let tape = { + let mut append = tapes.append(); + let tape = append + .open_blob_tape::(NAME, options) + .unwrap(); + append.append_bytes(&tape, b"abcdef").unwrap(); + append.commit(Persistence::Buffer).unwrap(); + tape + }; + + let reader = tapes.reader(); + assert_eq!(reader.blob_tape_len(&tape), Some(22)); + assert_eq!(reader.blob_tape_start(&tape), Some(16)); + + let mut contents = [0; 6]; + reader.read_bytes(&tape, 16, &mut contents).unwrap(); + assert_eq!(&contents, b"abcdef"); +} diff --git a/tests/sm_blob.rs b/tests/sm_blob.rs index 20c6c08..78fcceb 100644 --- a/tests/sm_blob.rs +++ b/tests/sm_blob.rs @@ -3,7 +3,10 @@ use proptest::prelude::*; use proptest_state_machine::{ReferenceStateMachine, StateMachineTest, prop_state_machine}; use std::fmt::{Debug, Formatter}; -use tapes::{BlobTape, Persistence, TapeOpenOptions, Tapes, TapesAppend, TapesRead, TapesTruncate}; +use tapes::{ + CachedBlobTape, CachedTapeOpenOptions, Persistence, Tapes, TapesAppend, TapesRead, + TapesTruncate, WholeBlobTape, WholeTapeOpenOptions, +}; #[derive(Clone, Debug)] enum TapeTransition { @@ -16,14 +19,12 @@ struct TapesState { dir: tempfile::TempDir, tapes: Tapes, - tape: BlobTape, + tape: CachedBlobTape, } impl Debug for TapesState { fn fmt(&self, f: &mut Formatter<'_>) -> std::fmt::Result { - f.debug_struct("tapes") - .field("cache", &self.tape.top_cache.read().as_slices()) - .finish() + f.debug_struct("tapes").finish() } } @@ -91,9 +92,11 @@ impl StateMachineTest for TapesState { let tape = append .open_blob_tape( "tape", - &TapeOpenOptions { + CachedTapeOpenOptions { + inner: WholeTapeOpenOptions { + dir: dir.path().to_path_buf(), + }, top_cache_size: ref_state.cache_size, - dir: dir.path().to_path_buf(), }, ) .unwrap(); @@ -121,7 +124,7 @@ impl StateMachineTest for TapesState { TapeTransition::Truncate(new_len) => { let mut truncate = state.tapes.truncate(); - truncate.truncate_blob_tape(&state.tape, new_len); + truncate.truncate_blob_tape(&state.tape, new_len).unwrap(); truncate.commit(Persistence::Buffer).unwrap(); } @@ -132,9 +135,11 @@ impl StateMachineTest for TapesState { let tape = append .open_blob_tape( "tape", - &TapeOpenOptions { + CachedTapeOpenOptions { + inner: WholeTapeOpenOptions { + dir: state.dir.path().to_path_buf(), + }, top_cache_size: ref_state.cache_size, - dir: state.dir.path().to_path_buf(), }, ) .unwrap(); diff --git a/tests/sm_fixed.rs b/tests/sm_fixed.rs index 8470822..c0bb537 100644 --- a/tests/sm_fixed.rs +++ b/tests/sm_fixed.rs @@ -4,7 +4,8 @@ use proptest_state_machine::{ReferenceStateMachine, StateMachineTest, prop_state use std::fmt::{Debug, Formatter}; use tapes::{ - FixedSizedTape, Persistence, TapeOpenOptions, Tapes, TapesAppend, TapesRead, TapesTruncate, + CachedBlobTape, CachedTapeOpenOptions, FixedSizedTape, Persistence, Tapes, TapesAppend, + TapesRead, TapesTruncate, WholeBlobTape, WholeTapeOpenOptions, }; #[derive(Clone, Debug)] @@ -18,7 +19,7 @@ struct TapesState { dir: tempfile::TempDir, tapes: Tapes, - tape: FixedSizedTape, + tape: FixedSizedTape>, } impl Debug for TapesState { @@ -91,11 +92,13 @@ impl StateMachineTest for TapesState { let mut append = tapes.append(); let tape = append - .open_fixed_sized_tape( + .open_fixed_sized_tape::>( "tape", - &TapeOpenOptions { + CachedTapeOpenOptions { + inner: WholeTapeOpenOptions { + dir: dir.path().to_path_buf(), + }, top_cache_size: ref_state.cache_size, - dir: dir.path().to_path_buf(), }, ) .unwrap(); @@ -126,7 +129,7 @@ impl StateMachineTest for TapesState { let low = amt / 2; let high = amt - low; - truncate.drop_fixed_sized_tape(&state.tape, low); + truncate.drop_fixed_sized_tape(&state.tape, low).unwrap(); for _ in 0..high { truncate.pop_fixed_sized_tape(&state.tape).unwrap(); @@ -139,11 +142,13 @@ impl StateMachineTest for TapesState { let mut append = tapes.append(); let tape = append - .open_fixed_sized_tape( + .open_fixed_sized_tape::>( "tape", - &TapeOpenOptions { + CachedTapeOpenOptions { + inner: WholeTapeOpenOptions { + dir: state.dir.path().to_path_buf(), + }, top_cache_size: ref_state.cache_size, - dir: state.dir.path().to_path_buf(), }, ) .unwrap(); diff --git a/tests/snapshot_reads.rs b/tests/snapshot_reads.rs new file mode 100644 index 0000000..0119d86 --- /dev/null +++ b/tests/snapshot_reads.rs @@ -0,0 +1,148 @@ +use std::{sync::mpsc, thread, time::Duration}; + +use tapes::{ + CachedBlobTape, CachedTapeOpenOptions, Persistence, Tapes, TapesAppend, TapesRead, + TapesTruncate, WholeBlobTape, WholeTapeOpenOptions, +}; + +#[test] +fn reader_keeps_its_snapshot_after_truncate_commits() { + let dir = tempfile::tempdir().unwrap(); + let tapes = Tapes::open(dir.path()).unwrap(); + let options = CachedTapeOpenOptions { + inner: WholeTapeOpenOptions { + dir: dir.path().to_path_buf(), + }, + top_cache_size: 4, + }; + + let mut append = tapes.append(); + let tape = append + .open_blob_tape::>("tape", options) + .unwrap(); + append.append_bytes(&tape, b"abcdefgh").unwrap(); + append.commit(Persistence::Buffer).unwrap(); + + let old_reader = tapes.reader(); + + let mut truncate = tapes.truncate(); + truncate.truncate_blob_tape(&tape, 6).unwrap(); + truncate.commit(Persistence::Buffer).unwrap(); + + let new_reader = tapes.reader(); + assert_eq!(new_reader.blob_tape_len(&tape), Some(6)); + let mut new_contents = [0; 6]; + new_reader.read_bytes(&tape, 0, &mut new_contents).unwrap(); + assert_eq!(&new_contents, b"abcdef"); + + assert_eq!(old_reader.blob_tape_len(&tape), Some(8)); + let mut old_contents = [0; 8]; + old_reader.read_bytes(&tape, 0, &mut old_contents).unwrap(); + assert_eq!(&old_contents, b"abcdefgh"); +} + +#[test] +fn reader_keeps_old_bytes_after_a_dropped_append_and_truncate() { + let dir = tempfile::tempdir().unwrap(); + let tapes = Tapes::open(dir.path()).unwrap(); + let options = CachedTapeOpenOptions { + inner: WholeTapeOpenOptions { + dir: dir.path().to_path_buf(), + }, + top_cache_size: 8, + }; + + let mut append = tapes.append(); + let tape = append + .open_blob_tape::>("tape", options) + .unwrap(); + append.append_bytes(&tape, b"abcdefgh").unwrap(); + append.commit(Persistence::Buffer).unwrap(); + + { + let mut append = tapes.append(); + append.append_bytes(&tape, b"stale").unwrap(); + } + + let old_reader = tapes.reader(); + + let mut truncate = tapes.truncate(); + truncate.truncate_blob_tape(&tape, 4).unwrap(); + truncate.commit(Persistence::Buffer).unwrap(); + + let mut old_contents = [0; 8]; + old_reader.read_bytes(&tape, 0, &mut old_contents).unwrap(); + assert_eq!(&old_contents, b"abcdefgh"); + drop(old_reader); + + let mut append = tapes.append(); + append.append_bytes(&tape, b"XY").unwrap(); + append.commit(Persistence::Buffer).unwrap(); + + let reader = tapes.reader(); + let mut contents = [0; 6]; + reader.read_bytes(&tape, 0, &mut contents).unwrap(); + assert_eq!(&contents, b"abcdXY"); +} + +#[test] +fn append_after_truncate_waits_for_the_old_reader() { + let dir = tempfile::tempdir().unwrap(); + let tapes = Tapes::open(dir.path()).unwrap(); + let options = CachedTapeOpenOptions { + inner: WholeTapeOpenOptions { + dir: dir.path().to_path_buf(), + }, + top_cache_size: 4, + }; + + let mut append = tapes.append(); + let tape = append + .open_blob_tape::>("tape", options) + .unwrap(); + append.append_bytes(&tape, b"abcdefgh").unwrap(); + append.commit(Persistence::Buffer).unwrap(); + + let old_reader = tapes.reader(); + + let mut truncate = tapes.truncate(); + truncate.truncate_blob_tape(&tape, 6).unwrap(); + truncate.commit(Persistence::Buffer).unwrap(); + + thread::scope(|scope| { + let (attempting_tx, attempting_rx) = mpsc::channel(); + let (finished_tx, finished_rx) = mpsc::channel(); + let tapes = &tapes; + let tape = &tape; + + scope.spawn(move || { + attempting_tx.send(()).unwrap(); + + let mut append = tapes.append(); + append.append_bytes(tape, b"XY").unwrap(); + append.commit(Persistence::Buffer).unwrap(); + + finished_tx.send(()).unwrap(); + }); + + attempting_rx.recv().unwrap(); + let finished_while_reader_was_alive = + finished_rx.recv_timeout(Duration::from_millis(100)).is_ok(); + + let mut old_contents = [0; 8]; + old_reader.read_bytes(tape, 0, &mut old_contents).unwrap(); + assert_eq!(&old_contents, b"abcdefgh"); + drop(old_reader); + + if !finished_while_reader_was_alive { + finished_rx.recv_timeout(Duration::from_secs(1)).unwrap(); + } + + assert!(!finished_while_reader_was_alive); + }); + + let reader = tapes.reader(); + let mut contents = [0; 8]; + reader.read_bytes(&tape, 0, &mut contents).unwrap(); + assert_eq!(&contents, b"abcdefXY"); +}