diff --git a/crates/tranquil-store/src/blockstore/compaction.rs b/crates/tranquil-store/src/blockstore/compaction.rs index 9cfc145..406e6d7 100644 --- a/crates/tranquil-store/src/blockstore/compaction.rs +++ b/crates/tranquil-store/src/blockstore/compaction.rs @@ -68,8 +68,8 @@ pub(super) fn compact_on_writer_thread( return Err(CompactionError::ActiveFileCannotBeCompacted); } - let source_fd = manager.open_for_read(source_file_id)?; - let source_size = manager.io().file_size(source_fd)?; + let source_handle = manager.open_for_read(source_file_id)?; + let source_size = manager.io().file_size(source_handle.fd())?; let new_file_id = file_ids.allocate(); @@ -77,7 +77,7 @@ pub(super) fn compact_on_writer_thread( manager, index, source_file_id, - source_fd, + source_handle.fd(), new_file_id, current_epoch, grace_period_ms, @@ -148,8 +148,8 @@ fn stream_compact( let mut reader = DataFileReader::open(manager.io(), source_fd)?; let now = crate::wall_clock_ms(); - let new_fd = manager.open_for_append(new_file_id)?; - let mut writer = DataFileWriter::new(manager.io(), new_fd, new_file_id)?; + let new_handle = manager.open_for_append(new_file_id)?; + let mut writer = DataFileWriter::new(manager.io(), new_handle.fd(), new_file_id)?; let hint_path = hint_file_path(manager.data_dir(), new_file_id); let hint_fd = manager.io().open(&hint_path, OpenOptions::read_write())?; diff --git a/crates/tranquil-store/src/blockstore/group_commit.rs b/crates/tranquil-store/src/blockstore/group_commit.rs index 331866d..8967202 100644 --- a/crates/tranquil-store/src/blockstore/group_commit.rs +++ b/crates/tranquil-store/src/blockstore/group_commit.rs @@ -553,7 +553,8 @@ fn initialize_active_state( match cursor { Some(wc) => { - let fd = manager.open_for_append(wc.file_id)?; + let handle = manager.open_for_append(wc.file_id)?; + let fd = handle.fd(); let file_size = manager.io().file_size(fd)?; if file_size < wc.offset.raw() { @@ -588,7 +589,8 @@ fn initialize_active_state( None => { let file_id = file_ids.allocate(); - let fd = manager.open_for_append(file_id)?; + let handle = manager.open_for_append(file_id)?; + let fd = handle.fd(); let writer = DataFileWriter::new(manager.io(), fd, file_id)?; writer.sync()?; let position = writer.position(); @@ -1074,9 +1076,9 @@ fn drain_and_process_remaining( shutdown_checkpoint(index, epoch, &ctx.hint_positions); } -struct RotationState { +struct RotationState { file_id: DataFileId, - fd: FileId, + handle: Arc>, hint_fd: FileId, } @@ -1178,7 +1180,7 @@ const VERIFY_RETRY_ATTEMPTS: u32 = 4; fn rollback_batch( manager: &DataFileManager, state: &ActiveState, - rotations: &[RotationState], + rotations: &[RotationState], ) { let _ = manager.io().truncate(state.fd, state.position.raw()); let _ = manager.io().sync(state.fd); @@ -1187,7 +1189,7 @@ fn rollback_batch( .truncate(state.hint_fd, state.hint_position.raw()); let _ = manager.io().sync(state.hint_fd); rotations.iter().for_each(|rot| { - manager.rollback_rotation(rot.file_id, rot.fd); + manager.rollback_rotation(rot.file_id); let _ = manager.io().close(rot.hint_fd); let _ = manager .io() @@ -1210,7 +1212,7 @@ fn process_batch( let mut all_decrements: Vec<[u8; CID_SIZE]> = Vec::new(); let mut current_hint_fd = state.hint_fd; - let mut rotations: Vec = Vec::new(); + let mut rotations: Vec> = Vec::new(); let mut data_writer = DataFileWriter::resume(manager.io(), state.fd, state.file_id, state.position); @@ -1222,12 +1224,14 @@ fn process_batch( hint_writer.sync().map_err(CommitError::from)?; let next_id = ctx.file_ids.allocate(); - let next_fd = manager.open_for_append(next_id)?; + let next_handle = manager.open_for_append(next_id)?; + let next_fd = next_handle.fd(); tracing::info!( from = %data_writer.file_id(), to = %next_id, - "data file rotation (batch boundary)" + trigger = "batch_boundary", + "data file rotation" ); data_writer = DataFileWriter::new(manager.io(), next_fd, next_id)?; @@ -1243,7 +1247,7 @@ fn process_batch( hint_writer = HintFileWriter::new(manager.io(), new_hint_fd); rotations.push(RotationState { file_id: next_id, - fd: next_fd, + handle: next_handle, hint_fd: new_hint_fd, }); } @@ -1348,7 +1352,7 @@ fn process_batch( let last_idx = rotations.len() - 1; rotations.iter().enumerate().for_each(|(i, rot)| { if i == last_idx { - manager.commit_rotation(rot.file_id, rot.fd); + manager.commit_rotation(rot.file_id, &rot.handle); ctx.active_files.register(ctx.shard_id, rot.file_id); } else { let _ = manager.io().close(rot.hint_fd); diff --git a/crates/tranquil-store/src/blockstore/manager.rs b/crates/tranquil-store/src/blockstore/manager.rs index 77b9b93..f6434d2 100644 --- a/crates/tranquil-store/src/blockstore/manager.rs +++ b/crates/tranquil-store/src/blockstore/manager.rs @@ -1,6 +1,7 @@ use std::collections::HashMap; use std::io; use std::path::{Path, PathBuf}; +use std::sync::Arc; use parking_lot::RwLock; @@ -13,22 +14,39 @@ pub const DEFAULT_MAX_FILE_SIZE: u64 = 256 * 1024 * 1024; pub(crate) const DATA_FILE_EXTENSION: &str = "tqb"; -struct CachedHandle { +pub struct CachedHandle { fd: FileId, + io: Arc, writable: bool, } +impl CachedHandle { + pub fn fd(&self) -> FileId { + self.fd + } + + pub fn is_writable(&self) -> bool { + self.writable + } +} + +impl Drop for CachedHandle { + fn drop(&mut self) { + let _ = self.io.close(self.fd); + } +} + pub struct DataFileManager { - io: S, + io: Arc, data_dir: PathBuf, max_file_size: u64, - handles: RwLock>, + handles: RwLock>>>, } impl DataFileManager { pub fn new(io: S, data_dir: PathBuf, max_file_size: u64) -> Self { Self { - io, + io: Arc::new(io), data_dir, max_file_size, handles: RwLock::new(HashMap::new()), @@ -40,7 +58,7 @@ impl DataFileManager { } pub fn io(&self) -> &S { - &self.io + self.io.as_ref() } pub fn data_dir(&self) -> &Path { @@ -56,76 +74,79 @@ impl DataFileManager { .join(format!("{file_id}.{DATA_FILE_EXTENSION}")) } - pub fn open_for_append(&self, file_id: DataFileId) -> io::Result { + pub fn open_for_append(&self, file_id: DataFileId) -> io::Result>> { { let cache = self.handles.read(); if let Some(entry) = cache.get(&file_id) && entry.writable { - return Ok(entry.fd); + return Ok(Arc::clone(entry)); } } let path = self.data_file_path(file_id); let fd = self.io.open(&path, OpenOptions::read_write())?; let mut cache = self.handles.write(); - match cache.get(&file_id) { + match cache.get(&file_id).cloned() { Some(entry) if entry.writable => { let _ = self.io.close(fd); - Ok(entry.fd) + Ok(entry) } - Some(entry) => { - let old_fd = entry.fd; - cache.insert(file_id, CachedHandle { fd, writable: true }); - let _ = self.io.close(old_fd); - Ok(fd) - } - None => { - cache.insert(file_id, CachedHandle { fd, writable: true }); - Ok(fd) + _ => { + let handle = Arc::new(CachedHandle { + fd, + io: Arc::clone(&self.io), + writable: true, + }); + cache.insert(file_id, Arc::clone(&handle)); + Ok(handle) } } } - pub fn open_for_read(&self, file_id: DataFileId) -> io::Result { + pub fn open_for_read(&self, file_id: DataFileId) -> io::Result>> { if let Some(entry) = self.handles.read().get(&file_id) { - return Ok(entry.fd); + return Ok(Arc::clone(entry)); } let path = self.data_file_path(file_id); let fd = self.io.open(&path, OpenOptions::read_only_existing())?; let mut cache = self.handles.write(); - match cache.get(&file_id) { + match cache.get(&file_id).cloned() { Some(entry) => { let _ = self.io.close(fd); - Ok(entry.fd) + Ok(entry) } None => { - cache.insert( - file_id, - CachedHandle { - fd, - writable: false, - }, - ); - Ok(fd) + let handle = Arc::new(CachedHandle { + fd, + io: Arc::clone(&self.io), + writable: false, + }); + cache.insert(file_id, Arc::clone(&handle)); + Ok(handle) } } } - pub fn prepare_rotation(&self, current: DataFileId) -> io::Result<(DataFileId, FileId)> { + pub fn prepare_rotation( + &self, + current: DataFileId, + ) -> io::Result<(DataFileId, Arc>)> { let next = current.next(); let path = self.data_file_path(next); let fd = self.io.open(&path, OpenOptions::read_write())?; - Ok((next, fd)) + let handle = Arc::new(CachedHandle { + fd, + io: Arc::clone(&self.io), + writable: true, + }); + Ok((next, handle)) } - pub fn commit_rotation(&self, file_id: DataFileId, fd: FileId) { - self.handles - .write() - .insert(file_id, CachedHandle { fd, writable: true }); + pub fn commit_rotation(&self, file_id: DataFileId, handle: &Arc>) { + self.handles.write().insert(file_id, Arc::clone(handle)); } - pub fn rollback_rotation(&self, file_id: DataFileId, fd: FileId) { - let _ = self.io.close(fd); + pub fn rollback_rotation(&self, file_id: DataFileId) { self.handles.write().remove(&file_id); let _ = self.io.delete(&self.data_file_path(file_id)); } @@ -135,14 +156,11 @@ impl DataFileManager { } pub fn list_files(&self) -> io::Result> { - list_files_by_extension(&self.io, &self.data_dir, DATA_FILE_EXTENSION) + list_files_by_extension(&*self.io, &self.data_dir, DATA_FILE_EXTENSION) } pub fn evict_handle(&self, file_id: DataFileId) { - let removed = self.handles.write().remove(&file_id); - if let Some(entry) = removed { - let _ = self.io.close(entry.fd); - } + self.handles.write().remove(&file_id); } pub fn delete_data_file(&self, file_id: DataFileId) -> io::Result<()> { @@ -152,14 +170,6 @@ impl DataFileManager { } } -impl Drop for DataFileManager { - fn drop(&mut self) { - self.handles.write().drain().for_each(|(_, entry)| { - let _ = self.io.close(entry.fd); - }); - } -} - #[cfg(test)] mod tests { use super::*; @@ -178,8 +188,8 @@ mod tests { #[test] fn open_for_append_creates_file() { let mgr = setup_manager(1024); - let fd = mgr.open_for_append(DataFileId::new(0)).unwrap(); - assert_eq!(mgr.io().file_size(fd).unwrap(), 0); + let handle = mgr.open_for_append(DataFileId::new(0)).unwrap(); + assert_eq!(mgr.io().file_size(handle.fd()).unwrap(), 0); } #[test] @@ -191,45 +201,52 @@ mod tests { #[test] fn handle_cache_returns_same_fd() { let mgr = setup_manager(1024); - let fd1 = mgr.open_for_append(DataFileId::new(0)).unwrap(); - let fd2 = mgr.open_for_append(DataFileId::new(0)).unwrap(); - assert_eq!(fd1, fd2); + let h1 = mgr.open_for_append(DataFileId::new(0)).unwrap(); + let h2 = mgr.open_for_append(DataFileId::new(0)).unwrap(); + assert_eq!(h1.fd(), h2.fd()); } #[test] fn open_for_read_uses_cache_from_append() { let mgr = setup_manager(1024); - let fd_write = mgr.open_for_append(DataFileId::new(0)).unwrap(); - let fd_read = mgr.open_for_read(DataFileId::new(0)).unwrap(); - assert_eq!(fd_write, fd_read); + let h_write = mgr.open_for_append(DataFileId::new(0)).unwrap(); + let h_read = mgr.open_for_read(DataFileId::new(0)).unwrap(); + assert_eq!(h_write.fd(), h_read.fd()); } #[test] fn rotation_lifecycle_prepare_commit() { let mgr = setup_manager(1024); - let _fd0 = mgr.open_for_append(DataFileId::new(0)).unwrap(); - let (next_id, next_fd) = mgr.prepare_rotation(DataFileId::new(0)).unwrap(); + let _h0 = mgr.open_for_append(DataFileId::new(0)).unwrap(); + let (next_id, next_handle) = mgr.prepare_rotation(DataFileId::new(0)).unwrap(); assert_eq!(next_id, DataFileId::new(1)); - assert_eq!(mgr.io().file_size(next_fd).unwrap(), 0); + assert_eq!(mgr.io().file_size(next_handle.fd()).unwrap(), 0); mgr.io().sync_dir(mgr.data_dir()).unwrap(); - mgr.commit_rotation(next_id, next_fd); - assert_eq!(mgr.open_for_read(next_id).unwrap(), next_fd); + mgr.commit_rotation(next_id, &next_handle); + assert_eq!( + mgr.open_for_read(next_id).unwrap().fd(), + next_handle.fd() + ); } #[test] fn rotation_rollback_cleans_handle_and_deletes_file() { let mgr = setup_manager(1024); - let _fd0 = mgr.open_for_append(DataFileId::new(0)).unwrap(); - let (next_id, next_fd) = mgr.prepare_rotation(DataFileId::new(0)).unwrap(); - mgr.commit_rotation(next_id, next_fd); + let _h0 = mgr.open_for_append(DataFileId::new(0)).unwrap(); + let (next_id, next_handle) = mgr.prepare_rotation(DataFileId::new(0)).unwrap(); + mgr.commit_rotation(next_id, &next_handle); - assert_eq!(mgr.open_for_read(next_id).unwrap(), next_fd); - mgr.rollback_rotation(next_id, next_fd); + assert_eq!( + mgr.open_for_read(next_id).unwrap().fd(), + next_handle.fd() + ); + drop(next_handle); + mgr.rollback_rotation(next_id); let reopen = mgr.open_for_read(next_id); assert!( reopen.is_err_and(|e| e.kind() == io::ErrorKind::NotFound), - "rollback must delete the data file so recovery cannot resurrect uncommitted bytes" + "rollback_rotation must delete the data file so recovery cannot resurrect uncommitted bytes" ); } @@ -245,8 +262,8 @@ mod tests { #[test] fn list_files_finds_data_files() { let mgr = setup_manager(1024); - let _fd0 = mgr.open_for_append(DataFileId::new(0)).unwrap(); - let _fd3 = mgr.open_for_append(DataFileId::new(3)).unwrap(); + let _h0 = mgr.open_for_append(DataFileId::new(0)).unwrap(); + let _h3 = mgr.open_for_append(DataFileId::new(3)).unwrap(); let files = mgr.list_files().unwrap(); assert_eq!(files, vec![DataFileId::new(0), DataFileId::new(3)]); @@ -255,7 +272,7 @@ mod tests { #[test] fn list_files_ignores_non_data_files() { let mgr = setup_manager(1024); - let _fd0 = mgr.open_for_append(DataFileId::new(0)).unwrap(); + let _h0 = mgr.open_for_append(DataFileId::new(0)).unwrap(); mgr.io() .open(Path::new("/data/notes.txt"), OpenOptions::read_write()) .unwrap(); @@ -280,32 +297,32 @@ mod tests { #[test] fn rotate_and_write_across_files() { let mgr = setup_manager(1024); - let fd0 = mgr.open_for_append(DataFileId::new(0)).unwrap(); - let mut writer0 = DataFileWriter::new(mgr.io(), fd0, DataFileId::new(0)).unwrap(); + let h0 = mgr.open_for_append(DataFileId::new(0)).unwrap(); + let mut writer0 = DataFileWriter::new(mgr.io(), h0.fd(), DataFileId::new(0)).unwrap(); let _ = writer0 .append_block(&test_cid(1), b"first file data") .unwrap(); writer0.sync().unwrap(); - let (id1, fd1) = mgr.prepare_rotation(DataFileId::new(0)).unwrap(); + let (id1, h1) = mgr.prepare_rotation(DataFileId::new(0)).unwrap(); mgr.io().sync_dir(mgr.data_dir()).unwrap(); - mgr.commit_rotation(id1, fd1); - let mut writer1 = DataFileWriter::new(mgr.io(), fd1, id1).unwrap(); + mgr.commit_rotation(id1, &h1); + let mut writer1 = DataFileWriter::new(mgr.io(), h1.fd(), id1).unwrap(); let _ = writer1 .append_block(&test_cid(2), b"second file data") .unwrap(); writer1.sync().unwrap(); - let fd0_read = mgr.open_for_read(DataFileId::new(0)).unwrap(); - let blocks0 = DataFileReader::open(mgr.io(), fd0_read) + let h0_read = mgr.open_for_read(DataFileId::new(0)).unwrap(); + let blocks0 = DataFileReader::open(mgr.io(), h0_read.fd()) .unwrap() .valid_blocks() .unwrap(); assert_eq!(blocks0.len(), 1); assert_eq!(blocks0[0].2, b"first file data"); - let fd1_read = mgr.open_for_read(id1).unwrap(); - let blocks1 = DataFileReader::open(mgr.io(), fd1_read) + let h1_read = mgr.open_for_read(id1).unwrap(); + let blocks1 = DataFileReader::open(mgr.io(), h1_read.fd()) .unwrap() .valid_blocks() .unwrap(); @@ -316,11 +333,11 @@ mod tests { #[test] fn read_cache_hit_from_writable_entry() { let mgr = setup_manager(1024); - let fd_write = mgr.open_for_append(DataFileId::new(0)).unwrap(); - DataFileWriter::new(mgr.io(), fd_write, DataFileId::new(0)).unwrap(); + let h_write = mgr.open_for_append(DataFileId::new(0)).unwrap(); + DataFileWriter::new(mgr.io(), h_write.fd(), DataFileId::new(0)).unwrap(); - let fd_read = mgr.open_for_read(DataFileId::new(0)).unwrap(); - assert_eq!(fd_write, fd_read); + let h_read = mgr.open_for_read(DataFileId::new(0)).unwrap(); + assert_eq!(h_write.fd(), h_read.fd()); } #[test] @@ -339,15 +356,15 @@ mod tests { mgr.io().sync_dir(mgr.data_dir()).unwrap(); mgr.io().close(raw_fd).unwrap(); - let fd_read = mgr.open_for_read(DataFileId::new(0)).unwrap(); - let _reader = DataFileReader::open(mgr.io(), fd_read).unwrap(); + let h_read = mgr.open_for_read(DataFileId::new(0)).unwrap(); + let _reader = DataFileReader::open(mgr.io(), h_read.fd()).unwrap(); - let fd_append = mgr.open_for_append(DataFileId::new(0)).unwrap(); - assert_ne!(fd_read, fd_append); + let h_append = mgr.open_for_append(DataFileId::new(0)).unwrap(); + assert_ne!(h_read.fd(), h_append.fd()); let mut writer = DataFileWriter::resume( mgr.io(), - fd_append, + h_append.fd(), DataFileId::new(0), BlockOffset::new(BLOCK_HEADER_SIZE as u64), ); @@ -356,7 +373,7 @@ mod tests { .unwrap(); writer.sync().unwrap(); - let blocks = DataFileReader::open(mgr.io(), fd_append) + let blocks = DataFileReader::open(mgr.io(), h_append.fd()) .unwrap() .valid_blocks() .unwrap(); diff --git a/crates/tranquil-store/src/blockstore/mod.rs b/crates/tranquil-store/src/blockstore/mod.rs index a594901..2bc9c03 100644 --- a/crates/tranquil-store/src/blockstore/mod.rs +++ b/crates/tranquil-store/src/blockstore/mod.rs @@ -24,7 +24,7 @@ pub use hint::{ HINT_FILE_EXTENSION, HINT_RECORD_SIZE, HintFileReader, HintFileWriter, HintIndex, ReadHintRecord, RebuildError, decode_hint_record, hint_file_path, scan_hints_to_memory, }; -pub use manager::{DEFAULT_MAX_FILE_SIZE, DataFileManager}; +pub use manager::{CachedHandle, DEFAULT_MAX_FILE_SIZE, DataFileManager}; pub use reader::{BlockStoreReader, ReadError}; pub use store::QuiesceGuard; pub use store::{BlockStoreConfig, DEFAULT_SHARD_COUNT, TranquilBlockStore}; diff --git a/crates/tranquil-store/src/blockstore/reader.rs b/crates/tranquil-store/src/blockstore/reader.rs index f26f499..ee53844 100644 --- a/crates/tranquil-store/src/blockstore/reader.rs +++ b/crates/tranquil-store/src/blockstore/reader.rs @@ -104,11 +104,11 @@ impl BlockStoreReader { }); by_file.into_iter().try_for_each(|(file_id, mut entries)| { - let fd = self.manager.open_for_read(file_id)?; - let file_size = self.manager.io().file_size(fd)?; + let handle = self.manager.open_for_read(file_id)?; + let file_size = self.manager.io().file_size(handle.fd())?; entries.sort_by_key(|(_, loc)| loc.offset); entries.into_iter().try_for_each(|(orig_idx, loc)| { - let data = self.decode_and_validate(fd, file_size, loc)?; + let data = self.decode_and_validate(handle.fd(), file_size, loc)?; results[orig_idx] = Some(data); Ok::<_, ReadError>(()) }) @@ -116,9 +116,9 @@ impl BlockStoreReader { } fn read_block_at(&self, location: BlockLocation) -> Result { - let fd = self.manager.open_for_read(location.file_id)?; - let file_size = self.manager.io().file_size(fd)?; - self.decode_and_validate(fd, file_size, location) + let handle = self.manager.open_for_read(location.file_id)?; + let file_size = self.manager.io().file_size(handle.fd())?; + self.decode_and_validate(handle.fd(), file_size, location) } fn decode_and_validate( diff --git a/crates/tranquil-store/src/eventlog/manager.rs b/crates/tranquil-store/src/eventlog/manager.rs index b4cc432..1666173 100644 --- a/crates/tranquil-store/src/eventlog/manager.rs +++ b/crates/tranquil-store/src/eventlog/manager.rs @@ -1,7 +1,8 @@ use std::collections::HashMap; use std::io; use std::path::{Path, PathBuf}; -use std::sync::atomic::{AtomicU64, Ordering}; +use std::sync::Arc; +use std::sync::atomic::{AtomicBool, AtomicU64, Ordering}; use parking_lot::RwLock; @@ -25,17 +26,38 @@ pub fn parse_segment_id(path: &Path) -> Option { (ext == SEGMENT_FILE_EXTENSION).then(|| stem.parse::().ok().map(SegmentId::new))? } -struct CachedSegmentHandle { +pub struct CachedSegmentHandle { fd: FileId, - sealed: bool, + io: Arc, + sealed: AtomicBool, writable: bool, } +impl CachedSegmentHandle { + pub fn fd(&self) -> FileId { + self.fd + } + + pub fn is_sealed(&self) -> bool { + self.sealed.load(Ordering::Acquire) + } + + pub fn is_writable(&self) -> bool { + self.writable + } +} + +impl Drop for CachedSegmentHandle { + fn drop(&mut self) { + let _ = self.io.close(self.fd); + } +} + pub struct SegmentManager { - io: S, + io: Arc, segments_dir: PathBuf, max_segment_size: u64, - handles: RwLock>, + handles: RwLock>>>, retention_epoch: AtomicU64, } @@ -51,7 +73,7 @@ impl SegmentManager { ); io.mkdir(&segments_dir)?; Ok(Self { - io, + io: Arc::new(io), segments_dir, max_segment_size, handles: RwLock::new(HashMap::new()), @@ -60,7 +82,7 @@ impl SegmentManager { } pub fn io(&self) -> &S { - &self.io + self.io.as_ref() } pub fn segments_dir(&self) -> &Path { @@ -92,52 +114,51 @@ impl SegmentManager { Ok(ids) } - pub fn open_for_read(&self, id: SegmentId) -> io::Result { + pub fn open_for_read(&self, id: SegmentId) -> io::Result>> { if let Some(entry) = self.handles.read().get(&id) { - return Ok(entry.fd); + return Ok(Arc::clone(entry)); } let path = self.segment_path(id); let fd = self.io.open(&path, OpenOptions::read_only_existing())?; let mut cache = self.handles.write(); - match cache.get(&id) { + match cache.get(&id).cloned() { Some(entry) => { let _ = self.io.close(fd); - Ok(entry.fd) + Ok(entry) } None => { - cache.insert( - id, - CachedSegmentHandle { - fd, - sealed: false, - writable: false, - }, - ); - Ok(fd) + let handle = Arc::new(CachedSegmentHandle { + fd, + io: Arc::clone(&self.io), + sealed: AtomicBool::new(false), + writable: false, + }); + cache.insert(id, Arc::clone(&handle)); + Ok(handle) } } } - pub fn open_for_append(&self, id: SegmentId) -> io::Result { + pub fn open_for_append(&self, id: SegmentId) -> io::Result>> { { let cache = self.handles.read(); if let Some(entry) = cache.get(&id) { - if entry.sealed { + if entry.is_sealed() { return Err(io::Error::new( io::ErrorKind::InvalidInput, format!("cannot append to sealed segment {id}"), )); } if entry.writable { - return Ok(entry.fd); + return Ok(Arc::clone(entry)); } } } let path = self.segment_path(id); let fd = self.io.open(&path, OpenOptions::read_write())?; let mut cache = self.handles.write(); - match cache.get(&id) { - Some(entry) if entry.sealed => { + match cache.get(&id).cloned() { + Some(entry) if entry.is_sealed() => { let _ = self.io.close(fd); Err(io::Error::new( io::ErrorKind::InvalidInput, @@ -146,31 +167,17 @@ impl SegmentManager { } Some(entry) if entry.writable => { let _ = self.io.close(fd); - Ok(entry.fd) + Ok(entry) } - Some(entry) => { - let old_fd = entry.fd; - cache.insert( - id, - CachedSegmentHandle { - fd, - sealed: false, - writable: true, - }, - ); - let _ = self.io.close(old_fd); - Ok(fd) - } - None => { - cache.insert( - id, - CachedSegmentHandle { - fd, - sealed: false, - writable: true, - }, - ); - Ok(fd) + _ => { + let handle = Arc::new(CachedSegmentHandle { + fd, + io: Arc::clone(&self.io), + sealed: AtomicBool::new(false), + writable: true, + }); + cache.insert(id, Arc::clone(&handle)); + Ok(handle) } } } @@ -179,37 +186,39 @@ impl SegmentManager { position.raw() >= self.max_segment_size } - pub fn prepare_rotation(&self, current_id: SegmentId) -> io::Result<(SegmentId, FileId)> { + pub fn prepare_rotation( + &self, + current_id: SegmentId, + ) -> io::Result<(SegmentId, Arc>)> { let next = current_id.next(); let path = self.segment_path(next); let fd = self.io.open(&path, OpenOptions::read_write())?; self.io.truncate(fd, 0)?; self.io.sync_dir(&self.segments_dir)?; - Ok((next, fd)) + let handle = Arc::new(CachedSegmentHandle { + fd, + io: Arc::clone(&self.io), + sealed: AtomicBool::new(false), + writable: true, + }); + Ok((next, handle)) } - pub fn commit_rotation(&self, new_id: SegmentId, fd: FileId) { - self.handles.write().insert( - new_id, - CachedSegmentHandle { - fd, - sealed: false, - writable: true, - }, - ); + pub fn commit_rotation(&self, new_id: SegmentId, handle: &Arc>) { + self.handles.write().insert(new_id, Arc::clone(handle)); } pub fn seal_segment(&self, id: SegmentId, index: &SegmentIndex) -> io::Result<()> { let path = self.index_path(id); - index.save(&self.io, &path)?; - let mut cache = self.handles.write(); - let entry = cache.get_mut(&id).ok_or_else(|| { + index.save(self.io.as_ref(), &path)?; + let cache = self.handles.read(); + let entry = cache.get(&id).ok_or_else(|| { io::Error::new( io::ErrorKind::InvalidInput, format!("seal_segment: segment {id} not in handle cache"), ) })?; - entry.sealed = true; + entry.sealed.store(true, Ordering::Release); Ok(()) } @@ -217,22 +226,16 @@ impl SegmentManager { self.handles .read() .get(&id) - .is_some_and(|entry| entry.sealed) + .is_some_and(|entry| entry.is_sealed()) } - pub fn rollback_rotation(&self, new_id: SegmentId, fd: FileId) { - let _ = self.io.close(fd); + pub fn rollback_rotation(&self, new_id: SegmentId) { self.handles.write().remove(&new_id); let _ = self.io.delete(&self.segment_path(new_id)); } pub fn delete_segment(&self, id: SegmentId) -> io::Result<()> { - { - let mut cache = self.handles.write(); - if let Some(entry) = cache.remove(&id) { - let _ = self.io.close(entry.fd); - } - } + self.handles.write().remove(&id); [self.index_path(id), self.sidecar_path(id)] .iter() .try_for_each(|path| match self.io.delete(path) { @@ -255,15 +258,7 @@ impl SegmentManager { } pub fn shutdown(&self) { - self.handles.write().drain().for_each(|(_, handle)| { - let _ = self.io.close(handle.fd); - }); - } -} - -impl Drop for SegmentManager { - fn drop(&mut self) { - self.shutdown(); + self.handles.write().clear(); } } @@ -329,7 +324,7 @@ mod tests { #[test] fn open_for_append_creates_file() { let mgr = setup_manager(1024); - let fd = mgr.open_for_append(SegmentId::new(1)).unwrap(); + let fd = mgr.open_for_append(SegmentId::new(1)).unwrap().fd(); assert_eq!(mgr.io().file_size(fd).unwrap(), 0); } @@ -342,24 +337,24 @@ mod tests { #[test] fn handle_cache_returns_same_fd() { let mgr = setup_manager(1024); - let fd1 = mgr.open_for_append(SegmentId::new(1)).unwrap(); - let fd2 = mgr.open_for_append(SegmentId::new(1)).unwrap(); + let fd1 = mgr.open_for_append(SegmentId::new(1)).unwrap().fd(); + let fd2 = mgr.open_for_append(SegmentId::new(1)).unwrap().fd(); assert_eq!(fd1, fd2); } #[test] fn open_for_read_uses_cache_from_append() { let mgr = setup_manager(1024); - let fd_write = mgr.open_for_append(SegmentId::new(1)).unwrap(); - let fd_read = mgr.open_for_read(SegmentId::new(1)).unwrap(); + let fd_write = mgr.open_for_append(SegmentId::new(1)).unwrap().fd(); + let fd_read = mgr.open_for_read(SegmentId::new(1)).unwrap().fd(); assert_eq!(fd_write, fd_read); } #[test] fn list_segments_finds_segment_files() { let mgr = setup_manager(1024); - mgr.open_for_append(SegmentId::new(1)).unwrap(); - mgr.open_for_append(SegmentId::new(3)).unwrap(); + mgr.open_for_append(SegmentId::new(1)).unwrap().fd(); + mgr.open_for_append(SegmentId::new(3)).unwrap().fd(); let segments = mgr.list_segments().unwrap(); assert_eq!(segments, vec![SegmentId::new(1), SegmentId::new(3)]); @@ -368,7 +363,7 @@ mod tests { #[test] fn list_segments_ignores_non_segment_files() { let mgr = setup_manager(1024); - mgr.open_for_append(SegmentId::new(1)).unwrap(); + mgr.open_for_append(SegmentId::new(1)).unwrap().fd(); mgr.io() .open(Path::new("/segments/notes.txt"), OpenOptions::read_write()) .unwrap(); @@ -380,7 +375,7 @@ mod tests { #[test] fn list_segments_ignores_index_files() { let mgr = setup_manager(1024); - mgr.open_for_append(SegmentId::new(1)).unwrap(); + mgr.open_for_append(SegmentId::new(1)).unwrap().fd(); mgr.io() .open( Path::new("/segments/00000001.tqi"), @@ -395,9 +390,9 @@ mod tests { #[test] fn list_segments_sorted_ascending() { let mgr = setup_manager(1024); - mgr.open_for_append(SegmentId::new(5)).unwrap(); - mgr.open_for_append(SegmentId::new(1)).unwrap(); - mgr.open_for_append(SegmentId::new(3)).unwrap(); + mgr.open_for_append(SegmentId::new(5)).unwrap().fd(); + mgr.open_for_append(SegmentId::new(1)).unwrap().fd(); + mgr.open_for_append(SegmentId::new(3)).unwrap().fd(); let segments = mgr.list_segments().unwrap(); assert_eq!( @@ -418,23 +413,30 @@ mod tests { #[test] fn rotation_lifecycle_prepare_commit() { let mgr = setup_manager(1024); - let _fd0 = mgr.open_for_append(SegmentId::new(1)).unwrap(); - let (next_id, next_fd) = mgr.prepare_rotation(SegmentId::new(1)).unwrap(); + let _h0 = mgr.open_for_append(SegmentId::new(1)).unwrap().fd(); + let (next_id, next_handle) = mgr.prepare_rotation(SegmentId::new(1)).unwrap(); assert_eq!(next_id, SegmentId::new(2)); - assert_eq!(mgr.io().file_size(next_fd).unwrap(), 0); - mgr.commit_rotation(next_id, next_fd); - assert_eq!(mgr.open_for_read(next_id).unwrap(), next_fd); + assert_eq!(mgr.io().file_size(next_handle.fd()).unwrap(), 0); + mgr.commit_rotation(next_id, &next_handle); + assert_eq!( + mgr.open_for_read(next_id).unwrap().fd(), + next_handle.fd() + ); } #[test] fn rotation_rollback_cleans_up() { let mgr = setup_manager(1024); - let _fd0 = mgr.open_for_append(SegmentId::new(1)).unwrap(); - let (next_id, next_fd) = mgr.prepare_rotation(SegmentId::new(1)).unwrap(); - mgr.commit_rotation(next_id, next_fd); + let _h0 = mgr.open_for_append(SegmentId::new(1)).unwrap().fd(); + let (next_id, next_handle) = mgr.prepare_rotation(SegmentId::new(1)).unwrap(); + mgr.commit_rotation(next_id, &next_handle); - assert_eq!(mgr.open_for_read(next_id).unwrap(), next_fd); - mgr.rollback_rotation(next_id, next_fd); + assert_eq!( + mgr.open_for_read(next_id).unwrap().fd(), + next_handle.fd() + ); + drop(next_handle); + mgr.rollback_rotation(next_id); let segments = mgr.list_segments().unwrap(); assert_eq!(segments, vec![SegmentId::new(1)]); @@ -443,7 +445,7 @@ mod tests { #[test] fn seal_segment_persists_index_and_marks_sealed() { let mgr = setup_manager(64 * 1024); - let fd = mgr.open_for_append(SegmentId::new(1)).unwrap(); + let fd = mgr.open_for_append(SegmentId::new(1)).unwrap().fd(); let mut writer = SegmentWriter::new( mgr.io(), fd, @@ -476,7 +478,7 @@ mod tests { #[test] fn delete_segment_removes_files_and_handle() { let mgr = setup_manager(64 * 1024); - let fd = mgr.open_for_append(SegmentId::new(1)).unwrap(); + let fd = mgr.open_for_append(SegmentId::new(1)).unwrap().fd(); let mut writer = SegmentWriter::new( mgr.io(), fd, @@ -507,9 +509,9 @@ mod tests { let mgr = setup_manager(1024); assert_eq!(mgr.oldest_segment().unwrap(), None); - mgr.open_for_append(SegmentId::new(3)).unwrap(); - mgr.open_for_append(SegmentId::new(1)).unwrap(); - mgr.open_for_append(SegmentId::new(5)).unwrap(); + mgr.open_for_append(SegmentId::new(3)).unwrap().fd(); + mgr.open_for_append(SegmentId::new(1)).unwrap().fd(); + mgr.open_for_append(SegmentId::new(5)).unwrap().fd(); assert_eq!(mgr.oldest_segment().unwrap(), Some(SegmentId::new(1))); } @@ -524,7 +526,7 @@ mod tests { fn rotate_and_write_across_segments() { let mgr = setup_manager(1024); - let fd1 = mgr.open_for_append(SegmentId::new(1)).unwrap(); + let fd1 = mgr.open_for_append(SegmentId::new(1)).unwrap().fd(); let mut writer1 = SegmentWriter::new( mgr.io(), fd1, @@ -538,8 +540,9 @@ mod tests { .unwrap(); writer1.sync(mgr.io()).unwrap(); - let (id2, fd2) = mgr.prepare_rotation(SegmentId::new(1)).unwrap(); - mgr.commit_rotation(id2, fd2); + let (id2, handle2) = mgr.prepare_rotation(SegmentId::new(1)).unwrap(); + let fd2 = handle2.fd(); + mgr.commit_rotation(id2, &handle2); let mut writer2 = SegmentWriter::new(mgr.io(), fd2, id2, EventSequence::new(2), MAX_EVENT_PAYLOAD) @@ -549,7 +552,7 @@ mod tests { .unwrap(); writer2.sync(mgr.io()).unwrap(); - let fd1_read = mgr.open_for_read(SegmentId::new(1)).unwrap(); + let fd1_read = mgr.open_for_read(SegmentId::new(1)).unwrap().fd(); let events1 = crate::eventlog::SegmentReader::open(mgr.io(), fd1_read, MAX_EVENT_PAYLOAD) .unwrap() .valid_prefix() @@ -557,7 +560,7 @@ mod tests { assert_eq!(events1.len(), 1); assert_eq!(events1[0].payload, b"first segment"); - let fd2_read = mgr.open_for_read(id2).unwrap(); + let fd2_read = mgr.open_for_read(id2).unwrap().fd(); let events2 = crate::eventlog::SegmentReader::open(mgr.io(), fd2_read, MAX_EVENT_PAYLOAD) .unwrap() .valid_prefix() @@ -569,7 +572,7 @@ mod tests { #[test] fn seal_then_append_errors() { let mgr = setup_manager(64 * 1024); - let fd = mgr.open_for_append(SegmentId::new(1)).unwrap(); + let fd = mgr.open_for_append(SegmentId::new(1)).unwrap().fd(); SegmentWriter::new( mgr.io(), fd, @@ -596,9 +599,9 @@ mod tests { #[test] fn multiple_deletions_increment_epoch() { let mgr = setup_manager(1024); - mgr.open_for_append(SegmentId::new(1)).unwrap(); - mgr.open_for_append(SegmentId::new(2)).unwrap(); - mgr.open_for_append(SegmentId::new(3)).unwrap(); + mgr.open_for_append(SegmentId::new(1)).unwrap().fd(); + mgr.open_for_append(SegmentId::new(2)).unwrap().fd(); + mgr.open_for_append(SegmentId::new(3)).unwrap().fd(); assert_eq!(mgr.retention_epoch(), 0); mgr.delete_segment(SegmentId::new(1)).unwrap(); @@ -610,7 +613,7 @@ mod tests { #[test] fn open_for_read_does_not_infer_sealed_from_index_file() { let mgr = setup_manager(64 * 1024); - let fd = mgr.open_for_append(SegmentId::new(1)).unwrap(); + let fd = mgr.open_for_append(SegmentId::new(1)).unwrap().fd(); let mut writer = SegmentWriter::new( mgr.io(), fd, @@ -630,26 +633,26 @@ mod tests { mgr.handles.write().remove(&SegmentId::new(1)); - let _read_fd = mgr.open_for_read(SegmentId::new(1)).unwrap(); + let _read_fd = mgr.open_for_read(SegmentId::new(1)).unwrap().fd(); assert!(!mgr.is_sealed(SegmentId::new(1))); } #[test] fn open_for_read_unsealed_allows_append() { let mgr = setup_manager(1024); - let _fd = mgr.open_for_append(SegmentId::new(1)).unwrap(); + let _fd = mgr.open_for_append(SegmentId::new(1)).unwrap().fd(); mgr.handles.write().remove(&SegmentId::new(1)); - let _read_fd = mgr.open_for_read(SegmentId::new(1)).unwrap(); + let _read_fd = mgr.open_for_read(SegmentId::new(1)).unwrap().fd(); assert!(!mgr.is_sealed(SegmentId::new(1))); } #[test] fn shutdown_clears_handles() { let mgr = setup_manager(1024); - mgr.open_for_append(SegmentId::new(1)).unwrap(); - mgr.open_for_append(SegmentId::new(2)).unwrap(); + mgr.open_for_append(SegmentId::new(1)).unwrap().fd(); + mgr.open_for_append(SegmentId::new(2)).unwrap().fd(); mgr.shutdown(); assert!(mgr.handles.read().is_empty()); @@ -665,7 +668,7 @@ mod tests { #[test] fn prepare_rotation_truncates_stale_file() { let mgr = setup_manager(1024); - let _fd0 = mgr.open_for_append(SegmentId::new(1)).unwrap(); + let _fd0 = mgr.open_for_append(SegmentId::new(1)).unwrap().fd(); let stale_path = mgr.segment_path(SegmentId::new(2)); let stale_fd = mgr @@ -677,23 +680,23 @@ mod tests { assert_eq!(mgr.io().file_size(stale_fd).unwrap(), 4096); mgr.io().close(stale_fd).unwrap(); - let (next_id, next_fd) = mgr.prepare_rotation(SegmentId::new(1)).unwrap(); + let (next_id, next_handle) = mgr.prepare_rotation(SegmentId::new(1)).unwrap(); assert_eq!(next_id, SegmentId::new(2)); - assert_eq!(mgr.io().file_size(next_fd).unwrap(), 0); + assert_eq!(mgr.io().file_size(next_handle.fd()).unwrap(), 0); } #[test] fn open_for_append_upgrades_read_only_handle() { let mgr = setup_manager(1024); - let fd_append = mgr.open_for_append(SegmentId::new(1)).unwrap(); + let fd_append = mgr.open_for_append(SegmentId::new(1)).unwrap().fd(); mgr.handles.write().remove(&SegmentId::new(1)); - let fd_read = mgr.open_for_read(SegmentId::new(1)).unwrap(); + let fd_read = mgr.open_for_read(SegmentId::new(1)).unwrap().fd(); assert_ne!(fd_read, fd_append); assert!(!mgr.handles.read().get(&SegmentId::new(1)).unwrap().writable); - let fd_upgraded = mgr.open_for_append(SegmentId::new(1)).unwrap(); + let fd_upgraded = mgr.open_for_append(SegmentId::new(1)).unwrap().fd(); assert_ne!(fd_upgraded, fd_read); assert!(mgr.handles.read().get(&SegmentId::new(1)).unwrap().writable); } diff --git a/crates/tranquil-store/src/eventlog/mod.rs b/crates/tranquil-store/src/eventlog/mod.rs index d23869b..fa08a46 100644 --- a/crates/tranquil-store/src/eventlog/mod.rs +++ b/crates/tranquil-store/src/eventlog/mod.rs @@ -411,8 +411,8 @@ impl EventLog { pub fn disk_usage(&self) -> io::Result { let segments = self.manager.list_segments()?; segments.iter().try_fold(0u64, |acc, &id| { - let fd = self.manager.open_for_read(id)?; - let size = self.manager.io().file_size(fd)?; + let handle = self.manager.open_for_read(id)?; + let size = self.manager.io().file_size(handle.fd())?; Ok(acc.saturating_add(size)) }) } diff --git a/crates/tranquil-store/src/eventlog/reader.rs b/crates/tranquil-store/src/eventlog/reader.rs index 7890e34..9b40eeb 100644 --- a/crates/tranquil-store/src/eventlog/reader.rs +++ b/crates/tranquil-store/src/eventlog/reader.rs @@ -145,10 +145,10 @@ impl EventLogReader { } fn rebuild_index(&self, segment_id: SegmentId) -> io::Result { - let fd = self.manager.open_for_read(segment_id)?; + let handle = self.manager.open_for_read(segment_id)?; let (idx, _) = rebuild_from_segment( self.manager.io(), - fd, + handle.fd(), DEFAULT_INDEX_INTERVAL, self.max_payload, )?; @@ -243,8 +243,8 @@ impl EventLogReader { return Ok(Arc::clone(m)); } - let fd = self.manager.open_for_read(segment_id)?; - let mapped = self.manager.io().mmap_file(fd)?; + let handle = self.manager.open_for_read(segment_id)?; + let mapped = self.manager.io().mmap_file(handle.fd())?; let arc = Arc::new(mapped); self.mmaps.write().insert(segment_id, Arc::clone(&arc)); Ok(arc) @@ -269,10 +269,10 @@ impl EventLogReader { predicate, ) } else { - let fd = self.manager.open_for_read(segment_id)?; - let file_size = self.manager.io().file_size(fd)?; + let handle = self.manager.open_for_read(segment_id)?; + let file_size = self.manager.io().file_size(handle.fd())?; self.scan_direct( - fd, + handle.fd(), file_size, start_offset, start_seq, @@ -513,8 +513,8 @@ impl EventLogReader { } fn rebuild_sidecar(&self, segment_id: SegmentId) -> io::Result { - let fd = self.manager.open_for_read(segment_id)?; - let sidecar = build_sidecar_from_segment(self.manager.io(), fd, self.max_payload)?; + let handle = self.manager.open_for_read(segment_id)?; + let sidecar = build_sidecar_from_segment(self.manager.io(), handle.fd(), self.max_payload)?; let _ = sidecar.save(self.manager.io(), &self.manager.sidecar_path(segment_id)); Ok(sidecar) } diff --git a/crates/tranquil-store/src/eventlog/writer.rs b/crates/tranquil-store/src/eventlog/writer.rs index 288d01c..b08df0a 100644 --- a/crates/tranquil-store/src/eventlog/writer.rs +++ b/crates/tranquil-store/src/eventlog/writer.rs @@ -66,9 +66,9 @@ impl EventLogWriter { index_interval: usize, max_payload: u32, ) -> io::Result { - let fd = manager.open_for_append(segment_id)?; - manager.io().truncate(fd, 0)?; - let writer = SegmentWriter::new(manager.io(), fd, segment_id, next_seq, max_payload)?; + let handle = manager.open_for_append(segment_id)?; + manager.io().truncate(handle.fd(), 0)?; + let writer = SegmentWriter::new(manager.io(), handle.fd(), segment_id, next_seq, max_payload)?; writer.sync(manager.io())?; manager.io().sync_dir(manager.segments_dir())?; @@ -93,7 +93,8 @@ impl EventLogWriter { index_interval: usize, max_payload: u32, ) -> io::Result { - let fd = manager.open_for_append(active_id)?; + let handle = manager.open_for_append(active_id)?; + let fd = handle.fd(); let (index, last_seq_in_active) = match rebuild_from_segment( manager.io(), @@ -308,11 +309,11 @@ impl EventLogWriter { Err(e) => warn!(segment = %old_id, error = %e, "non-fatal sidecar build failure"), } - let (new_id, new_fd) = self.manager.prepare_rotation(old_id)?; + let (new_id, new_handle) = self.manager.prepare_rotation(old_id)?; match SegmentWriter::new::( self.manager.io(), - new_fd, + new_handle.fd(), new_id, self.next_seq, self.max_payload, @@ -322,11 +323,12 @@ impl EventLogWriter { self.active_index = SegmentIndex::new(); self.event_count_in_segment = 0; self.last_event_offset = None; - self.manager.commit_rotation(new_id, new_fd); + self.manager.commit_rotation(new_id, &new_handle); Ok(Some(old_id)) } Err(e) => { - self.manager.rollback_rotation(new_id, new_fd); + drop(new_handle); + self.manager.rollback_rotation(new_id); Err(e) } } @@ -361,8 +363,8 @@ impl EventLogWriter { } fn build_sidecar_for_segment(&self, segment_id: SegmentId) -> io::Result<()> { - let fd = self.manager.open_for_read(segment_id)?; - let sidecar = build_sidecar_from_segment(self.manager.io(), fd, self.max_payload)?; + let handle = self.manager.open_for_read(segment_id)?; + let sidecar = build_sidecar_from_segment(self.manager.io(), handle.fd(), self.max_payload)?; let path = self.manager.sidecar_path(segment_id); sidecar.save(self.manager.io(), &path) } @@ -397,9 +399,9 @@ fn find_last_seq_from_segments( Ok(Some(idx)) => Ok(idx.last_seq()), Err(e) if e.kind() != io::ErrorKind::InvalidData => Err(e), _ => { - let fd = manager.open_for_read(seg_id)?; + let handle = manager.open_for_read(seg_id)?; let (_, last_seq) = - rebuild_from_segment(manager.io(), fd, DEFAULT_INDEX_INTERVAL, max_payload)?; + rebuild_from_segment(manager.io(), handle.fd(), DEFAULT_INDEX_INTERVAL, max_payload)?; Ok(last_seq) } } @@ -551,7 +553,7 @@ mod tests { assert_eq!(writer.synced_seq(), EventSequence::new(5)); assert_eq!(writer.active_segment_id(), SegmentId::new(1)); - let fd = mgr.open_for_read(SegmentId::new(1)).unwrap(); + let fd = mgr.open_for_read(SegmentId::new(1)).unwrap().fd(); let events = SegmentReader::open(mgr.io(), fd, MAX_EVENT_PAYLOAD) .unwrap() .valid_prefix() @@ -796,7 +798,7 @@ mod tests { .unwrap(); assert_eq!(writer.next_seq, EventSequence::new(3)); - let fd = mgr.open_for_read(SegmentId::new(1)).unwrap(); + let fd = mgr.open_for_read(SegmentId::new(1)).unwrap().fd(); let events = SegmentReader::open(mgr.io(), fd, MAX_EVENT_PAYLOAD) .unwrap() .valid_prefix() @@ -1003,7 +1005,7 @@ mod tests { assert_eq!(seq, EventSequence::new(4)); writer.sync().unwrap(); - let fd = mgr.open_for_read(SegmentId::new(1)).unwrap(); + let fd = mgr.open_for_read(SegmentId::new(1)).unwrap().fd(); let events = SegmentReader::open(mgr.io(), fd, MAX_EVENT_PAYLOAD) .unwrap() .valid_prefix() diff --git a/crates/tranquil-store/src/gauntlet/runner.rs b/crates/tranquil-store/src/gauntlet/runner.rs index 24d1ff6..3f1400b 100644 --- a/crates/tranquil-store/src/gauntlet/runner.rs +++ b/crates/tranquil-store/src/gauntlet/runner.rs @@ -595,7 +595,7 @@ fn eventlog_snapshot( let mut segment_last_ts: Vec<(SegmentId, u64)> = Vec::new(); segments.iter().for_each(|&id| { let per_segment: Vec = match s.manager.open_for_read(id) { - Ok(fd) => match SegmentReader::open(s.manager.io(), fd, MAX_EVENT_PAYLOAD) { + Ok(handle) => match SegmentReader::open(s.manager.io(), handle.fd(), MAX_EVENT_PAYLOAD) { Ok(reader) => reader.valid_prefix().unwrap_or_default(), Err(_) => Vec::new(), }, @@ -1108,8 +1108,8 @@ fn segment_last_timestamp( manager: &SegmentManager, id: SegmentId, ) -> std::io::Result> { - let fd = manager.open_for_read(id)?; - let reader = SegmentReader::open(manager.io(), fd, MAX_EVENT_PAYLOAD)?; + let handle = manager.open_for_read(id)?; + let reader = SegmentReader::open(manager.io(), handle.fd(), MAX_EVENT_PAYLOAD)?; let events = reader.valid_prefix()?; Ok(events.last().map(|e: &ValidEvent| e.timestamp.raw())) } diff --git a/crates/tranquil-store/tests/eventlog_crash.rs b/crates/tranquil-store/tests/eventlog_crash.rs index f910e61..eb87011 100644 --- a/crates/tranquil-store/tests/eventlog_crash.rs +++ b/crates/tranquil-store/tests/eventlog_crash.rs @@ -52,7 +52,7 @@ fn synced_events_survive_crash() { "seed {seed}: expected all synced events to survive" ); - let fd = mgr.open_for_read(SegmentId::new(1)).unwrap(); + let fd = mgr.open_for_read(SegmentId::new(1)).unwrap().fd(); let events = SegmentReader::open(mgr.io(), fd, MAX_EVENT_PAYLOAD) .unwrap() .valid_prefix() @@ -181,7 +181,7 @@ fn partial_event_truncated_on_recovery() { mgr.io().sync_dir(Path::new("/segments")).unwrap(); } - let fd = mgr.open_for_read(SegmentId::new(1)).unwrap(); + let fd = mgr.open_for_read(SegmentId::new(1)).unwrap().fd(); let file_size = mgr.io().file_size(fd).unwrap(); let partial_bytes = ((seed % 20) + 1) as usize; let junk: Vec = (0..partial_bytes) @@ -252,7 +252,7 @@ fn cross_segment_recovery() { ); sealed_segments[..sealed_count].iter().for_each(|&seg_id| { - let fd = mgr.open_for_read(seg_id).unwrap(); + let fd = mgr.open_for_read(seg_id).unwrap().fd(); let events = SegmentReader::open(mgr.io(), fd, MAX_EVENT_PAYLOAD) .unwrap() .valid_prefix() @@ -342,7 +342,7 @@ fn large_sealed_segment_index_rebuild_latency() { let index_path = mgr.index_path(SegmentId::new(1)); let _ = mgr.io().delete(&index_path); - let fd = mgr.open_for_read(SegmentId::new(1)).unwrap(); + let fd = mgr.open_for_read(SegmentId::new(1)).unwrap().fd(); let start = std::time::Instant::now(); let (index, last_seq) = rebuild_from_segment(mgr.io(), fd, 256, MAX_EVENT_PAYLOAD).unwrap(); @@ -417,7 +417,7 @@ fn pristine_comparison_under_faults() { } pristine_mgr.shutdown(); - let pristine_fd = pristine_mgr.open_for_read(SegmentId::new(1)).unwrap(); + let pristine_fd = pristine_mgr.open_for_read(SegmentId::new(1)).unwrap().fd(); let pristine_events = SegmentReader::open(pristine_mgr.io(), pristine_fd, MAX_EVENT_PAYLOAD) .unwrap() @@ -465,7 +465,7 @@ fn pristine_comparison_under_faults() { return Ok(None); } - let fd = faulty_clone.open_for_read(SegmentId::new(1))?; + let fd = faulty_clone.open_for_read(SegmentId::new(1))?.fd(); let events = SegmentReader::open(faulty_clone.io(), fd, MAX_EVENT_PAYLOAD)? .valid_prefix()?; Ok(Some(events)) @@ -592,7 +592,7 @@ fn pristine_comparison_parameterized_faults() { } pristine_mgr.shutdown(); - let pristine_fd = pristine_mgr.open_for_read(SegmentId::new(1)).unwrap(); + let pristine_fd = pristine_mgr.open_for_read(SegmentId::new(1)).unwrap().fd(); let pristine_events = SegmentReader::open(pristine_mgr.io(), pristine_fd, MAX_EVENT_PAYLOAD) .unwrap() .valid_prefix() @@ -637,7 +637,7 @@ fn pristine_comparison_parameterized_faults() { return Ok(None); } - let fd = faulty_clone.open_for_read(SegmentId::new(1))?; + let fd = faulty_clone.open_for_read(SegmentId::new(1))?.fd(); let events = SegmentReader::open(faulty_clone.io(), fd, MAX_EVENT_PAYLOAD)? .valid_prefix()?; Ok(Some(events)) diff --git a/crates/tranquil-store/tests/eventlog_manager_race.rs b/crates/tranquil-store/tests/eventlog_manager_race.rs new file mode 100644 index 0000000..5d58261 --- /dev/null +++ b/crates/tranquil-store/tests/eventlog_manager_race.rs @@ -0,0 +1,53 @@ +use std::path::PathBuf; +use std::sync::Arc; + +use tranquil_store::SimulatedIO; +use tranquil_store::StorageIO; +use tranquil_store::eventlog::{SegmentId, SegmentManager}; + +#[test] +fn concurrent_reader_survives_evict_on_segment_delete() { + let sim: Arc = Arc::new(SimulatedIO::pristine(0x1eed7a11)); + let segments_dir = PathBuf::from("/segments"); + + let manager = Arc::new( + SegmentManager::new(Arc::clone(&sim), segments_dir.clone(), 1 << 20).unwrap(), + ); + + let seg_id = SegmentId::new(1); + + let write_handle = manager.open_for_append(seg_id).unwrap(); + sim.write_at(write_handle.fd(), 0, b"arbitrary seed bytes for the segment") + .unwrap(); + sim.sync(write_handle.fd()).unwrap(); + sim.sync_dir(&segments_dir).unwrap(); + drop(write_handle); + + let ready_to_evict = Arc::new(std::sync::Barrier::new(2)); + let evict_done = Arc::new(std::sync::Barrier::new(2)); + + let reader_manager = Arc::clone(&manager); + let reader_io = Arc::clone(&sim); + let reader_ready = Arc::clone(&ready_to_evict); + let reader_done = Arc::clone(&evict_done); + + let reader = std::thread::spawn(move || { + let read_handle = reader_manager.open_for_read(seg_id).unwrap(); + reader_ready.wait(); + reader_done.wait(); + reader_io.file_size(read_handle.fd()) + }); + + ready_to_evict.wait(); + manager.delete_segment(seg_id).unwrap(); + evict_done.wait(); + + let read_result = reader.join().unwrap(); + assert!( + read_result.is_ok(), + "read against a FileId obtained before delete_segment must still succeed; \ + SegmentManager's delete_segment / rollback_rotation close the fd while a reader holds it. \ + error: {:?}", + read_result.err() + ); +} diff --git a/crates/tranquil-store/tests/eventlog_properties.rs b/crates/tranquil-store/tests/eventlog_properties.rs index 84b72fa..b341215 100644 --- a/crates/tranquil-store/tests/eventlog_properties.rs +++ b/crates/tranquil-store/tests/eventlog_properties.rs @@ -611,9 +611,9 @@ fn fsync_ordering_blocks_before_events() { SegmentManager::new(Arc::clone(&sim), PathBuf::from("/segments"), 64 * 1024).unwrap(), ); - let block_fd = block_mgr.open_for_append(DataFileId::new(0)).unwrap(); + let block_handle = block_mgr.open_for_append(DataFileId::new(0)).unwrap(); let mut block_writer = - DataFileWriter::new(block_mgr.io(), block_fd, DataFileId::new(0)).unwrap(); + DataFileWriter::new(block_mgr.io(), block_handle.fd(), DataFileId::new(0)).unwrap(); let cid = test_cid(1); let _ = block_writer.append_block(&cid, &[0xAA; 128]).unwrap(); block_writer.sync().unwrap(); diff --git a/crates/tranquil-store/tests/rotation_robustness.rs b/crates/tranquil-store/tests/rotation_robustness.rs new file mode 100644 index 0000000..f60b1ed --- /dev/null +++ b/crates/tranquil-store/tests/rotation_robustness.rs @@ -0,0 +1,234 @@ +mod common; + +use std::collections::HashMap; +use std::io; +use std::path::{Path, PathBuf}; +use std::sync::{Arc, Mutex}; +use std::sync::atomic::{AtomicBool, Ordering}; + +use tranquil_store::blockstore::{ + BlockStoreConfig, DataFileId, DataFileManager, DataFileWriter, GroupCommitConfig, + TranquilBlockStore, +}; +use tranquil_store::{FileId, MappedFile, OpenOptions, RealIO, SimulatedIO, StorageIO}; + +use common::{test_cid, with_runtime}; + +struct FailSpec { + target_path: Mutex>, + armed: AtomicBool, + tripped: AtomicBool, +} + +impl FailSpec { + fn new() -> Self { + Self { + target_path: Mutex::new(None), + armed: AtomicBool::new(false), + tripped: AtomicBool::new(false), + } + } + + fn arm_fail_first_sync_on(&self, path: &Path) { + *self.target_path.lock().unwrap() = Some(path.to_path_buf()); + self.tripped.store(false, Ordering::SeqCst); + self.armed.store(true, Ordering::SeqCst); + } + + fn fired(&self) -> bool { + self.tripped.load(Ordering::SeqCst) + } +} + +struct FailingIO { + inner: RealIO, + spec: Arc, + fd_to_path: Mutex>, +} + +impl FailingIO { + fn new(spec: Arc) -> Self { + Self { + inner: RealIO::new(), + spec, + fd_to_path: Mutex::new(HashMap::new()), + } + } +} + +impl StorageIO for FailingIO { + fn open(&self, path: &Path, opts: OpenOptions) -> io::Result { + let fd = self.inner.open(path, opts)?; + self.fd_to_path.lock().unwrap().insert(fd, path.to_path_buf()); + Ok(fd) + } + + fn close(&self, fd: FileId) -> io::Result<()> { + self.fd_to_path.lock().unwrap().remove(&fd); + self.inner.close(fd) + } + + fn read_at(&self, fd: FileId, offset: u64, buf: &mut [u8]) -> io::Result { + self.inner.read_at(fd, offset, buf) + } + + fn write_at(&self, fd: FileId, offset: u64, buf: &[u8]) -> io::Result { + self.inner.write_at(fd, offset, buf) + } + + fn sync(&self, fd: FileId) -> io::Result<()> { + let should_fail = self.spec.armed.load(Ordering::SeqCst) + && !self.spec.tripped.load(Ordering::SeqCst) + && match ( + self.fd_to_path.lock().unwrap().get(&fd).cloned(), + self.spec.target_path.lock().unwrap().clone(), + ) { + (Some(fd_path), Some(target)) => fd_path == target, + _ => false, + }; + match should_fail { + true => { + self.spec.tripped.store(true, Ordering::SeqCst); + Err(io::Error::other("injected sync failure on target path")) + } + false => self.inner.sync(fd), + } + } + + fn file_size(&self, fd: FileId) -> io::Result { + self.inner.file_size(fd) + } + + fn truncate(&self, fd: FileId, size: u64) -> io::Result<()> { + self.inner.truncate(fd, size) + } + + fn rename(&self, from: &Path, to: &Path) -> io::Result<()> { + self.inner.rename(from, to) + } + + fn delete(&self, path: &Path) -> io::Result<()> { + self.inner.delete(path) + } + + fn mkdir(&self, path: &Path) -> io::Result<()> { + self.inner.mkdir(path) + } + + fn sync_dir(&self, path: &Path) -> io::Result<()> { + self.inner.sync_dir(path) + } + + fn list_dir(&self, path: &Path) -> io::Result> { + self.inner.list_dir(path) + } + + fn mmap_file(&self, fd: FileId) -> io::Result { + self.inner.mmap_file(fd) + } +} + +#[test] +fn post_rotation_sync_failure_deletes_new_rotation_files() { + with_runtime(|| { + let dir = tempfile::TempDir::new().unwrap(); + let data_dir = dir.path().join("data"); + let index_dir = dir.path().join("index"); + + let spec = Arc::new(FailSpec::new()); + let spec_for_factory = Arc::clone(&spec); + + let config = BlockStoreConfig { + data_dir: data_dir.clone(), + index_dir, + max_file_size: 256, + group_commit: GroupCommitConfig::default(), + shard_count: 1, + }; + + let store = TranquilBlockStore::::open_with_io(config, move || { + FailingIO::new(Arc::clone(&spec_for_factory)) + }) + .unwrap(); + + store + .put_blocks_blocking(vec![(test_cid(1), vec![0xAA; 300])]) + .expect("priming put succeeds"); + + let rotated_data_path = data_dir.join("000002.tqb"); + let rotated_hint_path = data_dir.join("000002.tqh"); + + spec.arm_fail_first_sync_on(&rotated_data_path); + + let result = store.put_blocks_blocking(vec![(test_cid(2), vec![0xBB; 300])]); + assert!( + result.is_err(), + "put_blocks_blocking should surface the injected post-rotation sync failure" + ); + assert!( + spec.fired(), + "injector never observed a sync on the rotated data file; timing changed" + ); + + assert!( + !rotated_data_path.exists(), + "rotation rollback must delete the new data file after a post-write sync failure; \ + leaked file at {rotated_data_path:?}" + ); + assert!( + !rotated_hint_path.exists(), + "rotation rollback must delete the new hint file after a post-write sync failure; \ + leaked file at {rotated_hint_path:?}" + ); + }); +} + +#[test] +fn concurrent_reader_survives_evict_handle() { + let sim: Arc = Arc::new(SimulatedIO::pristine(0x13579bdf)); + let data_dir = Path::new("/data"); + sim.mkdir(data_dir).unwrap(); + sim.sync_dir(data_dir).unwrap(); + + let manager = Arc::new(DataFileManager::new( + Arc::clone(&sim), + data_dir.to_path_buf(), + 1 << 20, + )); + + let file_id = DataFileId::new(0); + let write_handle = manager.open_for_append(file_id).unwrap(); + { + let mut writer = DataFileWriter::new(&*sim, write_handle.fd(), file_id).unwrap(); + let _ = writer.append_block(&test_cid(1), &vec![0x11; 128]).unwrap(); + writer.sync().unwrap(); + } + drop(write_handle); + + let ready_to_evict = Arc::new(std::sync::Barrier::new(2)); + let evict_done = Arc::new(std::sync::Barrier::new(2)); + + let reader_manager = Arc::clone(&manager); + let reader_io = Arc::clone(&sim); + let reader_ready = Arc::clone(&ready_to_evict); + let reader_done = Arc::clone(&evict_done); + + let reader = std::thread::spawn(move || { + let read_handle = reader_manager.open_for_read(file_id).unwrap(); + reader_ready.wait(); + reader_done.wait(); + reader_io.file_size(read_handle.fd()) + }); + + ready_to_evict.wait(); + manager.evict_handle(file_id); + evict_done.wait(); + + let read_result = reader.join().unwrap(); + assert!( + read_result.is_ok(), + "read against a FileId obtained before evict_handle must still succeed; \ + evict_handle closed the underlying fd while the reader held it. error: {:?}", + read_result.err() + ); +} diff --git a/crates/tranquil-store/tests/sim_blockstore.rs b/crates/tranquil-store/tests/sim_blockstore.rs index 89fedb8..d35e72e 100644 --- a/crates/tranquil-store/tests/sim_blockstore.rs +++ b/crates/tranquil-store/tests/sim_blockstore.rs @@ -35,7 +35,8 @@ impl SimHarness { Arc::clone(&self.sim), self.data_dir.to_path_buf(), ); - let fd = manager.open_for_append(file_id).unwrap(); + let handle = manager.open_for_append(file_id).unwrap(); + let fd = handle.fd(); let file_size = self.sim.file_size(fd).unwrap(); match file_size { 0 => { @@ -504,9 +505,10 @@ fn sim_aggressive_faults_data_integrity() { let manager = DataFileManager::with_default_max_size(Arc::clone(&sim), data_dir.to_path_buf()); - let Ok(fd) = manager.open_for_append(file_id) else { + let Ok(handle) = manager.open_for_append(file_id) else { return; }; + let fd = handle.fd(); let writer_result = DataFileWriter::new(&*sim, fd, file_id); let Ok(writer) = writer_result else { return }; @@ -515,7 +517,8 @@ fn sim_aggressive_faults_data_integrity() { return; }; let start_pos = writer.position(); - let _ = sim.close(fd); + drop(writer); + drop(handle); let mut rng = Rng::new(seed); let block_count = (rng.range_u32(15) + 5) as u16; diff --git a/crates/tranquil-store/tests/sim_eventlog.rs b/crates/tranquil-store/tests/sim_eventlog.rs index 0637a8f..51e771b 100644 --- a/crates/tranquil-store/tests/sim_eventlog.rs +++ b/crates/tranquil-store/tests/sim_eventlog.rs @@ -39,7 +39,7 @@ fn read_all_events(mgr: &SegmentManager, seed: u64) -> Vec = Arc::new(RealIO::new()); let manager = DataFileManager::new(Arc::clone(&io), data_dir.clone(), 4096); - let (next_id, next_fd) = manager.prepare_rotation(DataFileId::new(0)).unwrap(); - manager.commit_rotation(next_id, next_fd); + let (next_id, next_handle) = manager.prepare_rotation(DataFileId::new(0)).unwrap(); + manager.commit_rotation(next_id, &next_handle); - let mut writer = DataFileWriter::new(&*io, next_fd, next_id).unwrap(); + let mut writer = DataFileWriter::new(&*io, next_handle.fd(), next_id).unwrap(); let _ = writer.append_block(&orphan_cid, &vec![0xAB; 256]).unwrap(); writer.sync().unwrap(); io.sync_dir(&data_dir).unwrap(); let _ = io.delete(&hint_file_path(&data_dir, next_id)); - manager.rollback_rotation(next_id, next_fd); + drop(writer); + drop(next_handle); + manager.rollback_rotation(next_id); } let store = TranquilBlockStore::open(config).unwrap();