fix(tranquil-store): arc-counted cache handles, reader-eviction race

Lewis: May this revision serve well! <lu5a@proton.me>
This commit is contained in:
Lewis
2026-04-21 22:04:24 +03:00
parent 6d2d3b4be4
commit 00c9eb732f
17 changed files with 619 additions and 300 deletions
@@ -68,8 +68,8 @@ pub(super) fn compact_on_writer_thread<S: StorageIO>(
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<S: StorageIO>(
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<S: StorageIO>(
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())?;
@@ -553,7 +553,8 @@ fn initialize_active_state<S: StorageIO>(
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<S: StorageIO>(
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<S: StorageIO>(
shutdown_checkpoint(index, epoch, &ctx.hint_positions);
}
struct RotationState {
struct RotationState<S: StorageIO> {
file_id: DataFileId,
fd: FileId,
handle: Arc<super::manager::CachedHandle<S>>,
hint_fd: FileId,
}
@@ -1178,7 +1180,7 @@ const VERIFY_RETRY_ATTEMPTS: u32 = 4;
fn rollback_batch<S: StorageIO>(
manager: &DataFileManager<S>,
state: &ActiveState,
rotations: &[RotationState],
rotations: &[RotationState<S>],
) {
let _ = manager.io().truncate(state.fd, state.position.raw());
let _ = manager.io().sync(state.fd);
@@ -1187,7 +1189,7 @@ fn rollback_batch<S: StorageIO>(
.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<S: StorageIO>(
let mut all_decrements: Vec<[u8; CID_SIZE]> = Vec::new();
let mut current_hint_fd = state.hint_fd;
let mut rotations: Vec<RotationState> = Vec::new();
let mut rotations: Vec<RotationState<S>> = Vec::new();
let mut data_writer =
DataFileWriter::resume(manager.io(), state.fd, state.file_id, state.position);
@@ -1222,12 +1224,14 @@ fn process_batch<S: StorageIO>(
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<S: StorageIO>(
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<S: StorageIO>(
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);
+109 -92
View File
@@ -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<S: StorageIO> {
fd: FileId,
io: Arc<S>,
writable: bool,
}
impl<S: StorageIO> CachedHandle<S> {
pub fn fd(&self) -> FileId {
self.fd
}
pub fn is_writable(&self) -> bool {
self.writable
}
}
impl<S: StorageIO> Drop for CachedHandle<S> {
fn drop(&mut self) {
let _ = self.io.close(self.fd);
}
}
pub struct DataFileManager<S: StorageIO> {
io: S,
io: Arc<S>,
data_dir: PathBuf,
max_file_size: u64,
handles: RwLock<HashMap<DataFileId, CachedHandle>>,
handles: RwLock<HashMap<DataFileId, Arc<CachedHandle<S>>>>,
}
impl<S: StorageIO> DataFileManager<S> {
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<S: StorageIO> DataFileManager<S> {
}
pub fn io(&self) -> &S {
&self.io
self.io.as_ref()
}
pub fn data_dir(&self) -> &Path {
@@ -56,76 +74,79 @@ impl<S: StorageIO> DataFileManager<S> {
.join(format!("{file_id}.{DATA_FILE_EXTENSION}"))
}
pub fn open_for_append(&self, file_id: DataFileId) -> io::Result<FileId> {
pub fn open_for_append(&self, file_id: DataFileId) -> io::Result<Arc<CachedHandle<S>>> {
{
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<FileId> {
pub fn open_for_read(&self, file_id: DataFileId) -> io::Result<Arc<CachedHandle<S>>> {
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<CachedHandle<S>>)> {
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<CachedHandle<S>>) {
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<S: StorageIO> DataFileManager<S> {
}
pub fn list_files(&self) -> io::Result<Vec<DataFileId>> {
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<S: StorageIO> DataFileManager<S> {
}
}
impl<S: StorageIO> Drop for DataFileManager<S> {
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();
+1 -1
View File
@@ -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};
@@ -104,11 +104,11 @@ impl<S: StorageIO> BlockStoreReader<S> {
});
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<S: StorageIO> BlockStoreReader<S> {
}
fn read_block_at(&self, location: BlockLocation) -> Result<Bytes, ReadError> {
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(
+133 -130
View File
@@ -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<SegmentId> {
(ext == SEGMENT_FILE_EXTENSION).then(|| stem.parse::<u32>().ok().map(SegmentId::new))?
}
struct CachedSegmentHandle {
pub struct CachedSegmentHandle<S: StorageIO> {
fd: FileId,
sealed: bool,
io: Arc<S>,
sealed: AtomicBool,
writable: bool,
}
impl<S: StorageIO> CachedSegmentHandle<S> {
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<S: StorageIO> Drop for CachedSegmentHandle<S> {
fn drop(&mut self) {
let _ = self.io.close(self.fd);
}
}
pub struct SegmentManager<S: StorageIO> {
io: S,
io: Arc<S>,
segments_dir: PathBuf,
max_segment_size: u64,
handles: RwLock<HashMap<SegmentId, CachedSegmentHandle>>,
handles: RwLock<HashMap<SegmentId, Arc<CachedSegmentHandle<S>>>>,
retention_epoch: AtomicU64,
}
@@ -51,7 +73,7 @@ impl<S: StorageIO> SegmentManager<S> {
);
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<S: StorageIO> SegmentManager<S> {
}
pub fn io(&self) -> &S {
&self.io
self.io.as_ref()
}
pub fn segments_dir(&self) -> &Path {
@@ -92,52 +114,51 @@ impl<S: StorageIO> SegmentManager<S> {
Ok(ids)
}
pub fn open_for_read(&self, id: SegmentId) -> io::Result<FileId> {
pub fn open_for_read(&self, id: SegmentId) -> io::Result<Arc<CachedSegmentHandle<S>>> {
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<FileId> {
pub fn open_for_append(&self, id: SegmentId) -> io::Result<Arc<CachedSegmentHandle<S>>> {
{
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<S: StorageIO> SegmentManager<S> {
}
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<S: StorageIO> SegmentManager<S> {
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<CachedSegmentHandle<S>>)> {
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<CachedSegmentHandle<S>>) {
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<S: StorageIO> SegmentManager<S> {
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<S: StorageIO> SegmentManager<S> {
}
pub fn shutdown(&self) {
self.handles.write().drain().for_each(|(_, handle)| {
let _ = self.io.close(handle.fd);
});
}
}
impl<S: StorageIO> Drop for SegmentManager<S> {
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);
}
+2 -2
View File
@@ -411,8 +411,8 @@ impl<S: StorageIO + 'static> EventLog<S> {
pub fn disk_usage(&self) -> io::Result<u64> {
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))
})
}
+9 -9
View File
@@ -145,10 +145,10 @@ impl<S: StorageIO> EventLogReader<S> {
}
fn rebuild_index(&self, segment_id: SegmentId) -> io::Result<SegmentIndex> {
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<S: StorageIO> EventLogReader<S> {
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<S: StorageIO> EventLogReader<S> {
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<S: StorageIO> EventLogReader<S> {
}
fn rebuild_sidecar(&self, segment_id: SegmentId) -> io::Result<SidecarIndex> {
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)
}
+17 -15
View File
@@ -66,9 +66,9 @@ impl<S: StorageIO> EventLogWriter<S> {
index_interval: usize,
max_payload: u32,
) -> io::Result<Self> {
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<S: StorageIO> EventLogWriter<S> {
index_interval: usize,
max_payload: u32,
) -> io::Result<Self> {
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<S: StorageIO> EventLogWriter<S> {
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::<S>(
self.manager.io(),
new_fd,
new_handle.fd(),
new_id,
self.next_seq,
self.max_payload,
@@ -322,11 +323,12 @@ impl<S: StorageIO> EventLogWriter<S> {
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<S: StorageIO> EventLogWriter<S> {
}
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<S: StorageIO>(
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()
+3 -3
View File
@@ -595,7 +595,7 @@ fn eventlog_snapshot<S: StorageIO + Send + Sync + 'static>(
let mut segment_last_ts: Vec<(SegmentId, u64)> = Vec::new();
segments.iter().for_each(|&id| {
let per_segment: Vec<ValidEvent> = 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<S: StorageIO + Send + Sync + 'static>(
manager: &SegmentManager<S>,
id: SegmentId,
) -> std::io::Result<Option<u64>> {
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()))
}
@@ -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<u8> = (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))
@@ -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<SimulatedIO> = 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()
);
}
@@ -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();
@@ -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<Option<PathBuf>>,
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<FailSpec>,
fd_to_path: Mutex<HashMap<FileId, PathBuf>>,
}
impl FailingIO {
fn new(spec: Arc<FailSpec>) -> 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<FileId> {
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<usize> {
self.inner.read_at(fd, offset, buf)
}
fn write_at(&self, fd: FileId, offset: u64, buf: &[u8]) -> io::Result<usize> {
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<u64> {
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<Vec<PathBuf>> {
self.inner.list_dir(path)
}
fn mmap_file(&self, fd: FileId) -> io::Result<MappedFile> {
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::<FailingIO>::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<SimulatedIO> = 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()
);
}
@@ -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;
+10 -9
View File
@@ -39,7 +39,7 @@ fn read_all_events(mgr: &SegmentManager<SimulatedIO>, seed: u64) -> Vec<ValidEve
.flat_map(|&seg_id| {
let fd = mgr
.open_for_read(seg_id)
.unwrap_or_else(|e| panic!("seed {seed}: open_for_read({seg_id}) failed: {e}"));
.unwrap_or_else(|e| panic!("seed {seed}: open_for_read({seg_id}) failed: {e}")).fd();
SegmentReader::open(mgr.io(), fd, MAX_EVENT_PAYLOAD)
.unwrap_or_else(|e| {
panic!(
@@ -256,7 +256,7 @@ fn segment_deletion_does_not_corrupt_neighbors() {
let seg2_fd = mgr
.open_for_read(SegmentId::new(2))
.unwrap_or_else(|e| panic!("seed {seed}: open_for_read(2) failed: {e}"));
.unwrap_or_else(|e| panic!("seed {seed}: open_for_read(2) failed: {e}")).fd();
let seg2_events = SegmentReader::open(mgr.io(), seg2_fd, MAX_EVENT_PAYLOAD)
.unwrap_or_else(|e| {
panic!("seed {seed}: SegmentReader::open(2, MAX_EVENT_PAYLOAD) failed: {e}")
@@ -271,7 +271,7 @@ fn segment_deletion_does_not_corrupt_neighbors() {
let seg3_fd = mgr
.open_for_read(SegmentId::new(3))
.unwrap_or_else(|e| panic!("seed {seed}: open_for_read(3) failed: {e}"));
.unwrap_or_else(|e| panic!("seed {seed}: open_for_read(3) failed: {e}")).fd();
let seg3_events = SegmentReader::open(mgr.io(), seg3_fd, MAX_EVENT_PAYLOAD)
.unwrap_or_else(|e| {
panic!("seed {seed}: SegmentReader::open(3, MAX_EVENT_PAYLOAD) failed: {e}")
@@ -388,7 +388,7 @@ fn fsync_ordering_unsynced_events_never_durable() {
let fd = mgr
.open_for_read(SegmentId::new(1))
.unwrap_or_else(|e| panic!("seed {seed}: open_for_read(1) failed: {e}"));
.unwrap_or_else(|e| panic!("seed {seed}: open_for_read(1) failed: {e}")).fd();
let recovered = SegmentReader::open(mgr.io(), fd, MAX_EVENT_PAYLOAD)
.unwrap_or_else(|e| {
panic!("seed {seed}: SegmentReader::open(1, MAX_EVENT_PAYLOAD) failed: {e}")
@@ -458,7 +458,7 @@ fn fsync_ordering_proof_sync_before_blockstore_ack() {
let fd = mgr
.open_for_read(SegmentId::new(1))
.unwrap_or_else(|e| panic!("seed {seed}: open_for_read(1) failed: {e}"));
.unwrap_or_else(|e| panic!("seed {seed}: open_for_read(1) failed: {e}")).fd();
let recovered = SegmentReader::open(mgr.io(), fd, MAX_EVENT_PAYLOAD)
.unwrap_or_else(|e| {
panic!("seed {seed}: SegmentReader::open(1, MAX_EVENT_PAYLOAD) failed: {e}")
@@ -579,7 +579,7 @@ fn group_sync_crash_mid_sync_partial_fsync() {
let fd = mgr
.open_for_read(SegmentId::new(1))
.unwrap_or_else(|e| panic!("seed {seed}: open_for_read(1) failed: {e}"));
.unwrap_or_else(|e| panic!("seed {seed}: open_for_read(1) failed: {e}")).fd();
let events = SegmentReader::open(mgr.io(), fd, MAX_EVENT_PAYLOAD)
.unwrap_or_else(|e| {
panic!("seed {seed}: SegmentReader::open(1, MAX_EVENT_PAYLOAD) failed: {e}")
@@ -652,7 +652,7 @@ fn group_sync_no_double_sync_no_skipped_events() {
let fd = mgr
.open_for_read(SegmentId::new(1))
.unwrap_or_else(|e| panic!("seed {seed}: open_for_read(1) failed: {e}"));
.unwrap_or_else(|e| panic!("seed {seed}: open_for_read(1) failed: {e}")).fd();
let events = SegmentReader::open(mgr.io(), fd, MAX_EVENT_PAYLOAD)
.unwrap_or_else(|e| {
panic!("seed {seed}: SegmentReader::open(1, MAX_EVENT_PAYLOAD) failed: {e}")
@@ -938,7 +938,7 @@ fn aggressive_faults_group_sync_recovery() {
let pristine_fd = pristine_mgr
.open_for_read(SegmentId::new(1))
.unwrap_or_else(|e| panic!("seed {seed}: pristine open_for_read(1) failed: {e}"));
.unwrap_or_else(|e| panic!("seed {seed}: pristine open_for_read(1) failed: {e}")).fd();
let pristine_events = SegmentReader::open(
pristine_mgr.io(),
pristine_fd,
@@ -985,9 +985,10 @@ fn aggressive_faults_group_sync_recovery() {
return;
}
let Ok(fd) = mgr.open_for_read(SegmentId::new(1)) else {
let Ok(handle) = mgr.open_for_read(SegmentId::new(1)) else {
return;
};
let fd = handle.fd();
let Ok(reader) = SegmentReader::open(mgr.io(), fd, MAX_EVENT_PAYLOAD) else {
return;
};
@@ -55,16 +55,18 @@ fn rollback_rotation_does_not_leave_orphan_data_file() {
{
let io: Arc<RealIO> = 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();