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

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
19 changes: 13 additions & 6 deletions src/tapes.rs
Original file line number Diff line number Diff line change
Expand Up @@ -97,18 +97,25 @@ impl TapesAppendTransaction {
name: &'static str,
options: B::OpenConfig,
) -> io::Result<FixedSizedTape<E, B>> {
let inner = self.open_blob_tape(name, options)?;
let metadata = self.metadata_guard.get(name).copied();
let start = metadata.map_or(options.start_index(), |metadata| metadata.start);
let len = metadata.map_or(start, |metadata| metadata.len);

let entry_size = size_of::<E>() as u64;
if !start.is_multiple_of(entry_size) {
return Err(io::Error::other(
"Tape start is not a multiple of entry size",
));
}

if self
.metadata_guard
.get(name)
.is_some_and(|metadata| !(metadata.len as usize).is_multiple_of(size_of::<E>()))
{
if !len.is_multiple_of(entry_size) {
return Err(io::Error::other(
"Tape size is not a multiple of entry size",
));
}

let inner = self.open_blob_tape(name, options)?;

Ok(FixedSizedTape {
inner,
phantom_data: PhantomData,
Expand Down
8 changes: 7 additions & 1 deletion src/tapes/cached_tape.rs
Original file line number Diff line number Diff line change
Expand Up @@ -59,9 +59,13 @@ impl<B: BlobTape + 'static> BlobTape for CachedBlobTape<B> {
current_epoch: u64,
config: Self::OpenConfig,
) -> io::Result<Self> {
let start_index = config.inner.start_index();
let tape = B::open(name, tape_metadata, current_epoch, config.inner)?;

let metadata = tape_metadata.unwrap_or_default();
let metadata = tape_metadata.unwrap_or(TapeMetadata {
start: start_index,
len: start_index,
});
let start = max(
metadata.len.saturating_sub(config.top_cache_size),
metadata.start,
Expand Down Expand Up @@ -91,6 +95,8 @@ impl<B: BlobTape + 'static> BlobTape for CachedBlobTape<B> {
last_byte_needed_offset -= buf_to_fill.len() as u64;
}

drop(top_cache);

if last_byte_needed_offset != offset {
self.tape.read_bytes(
offset,
Expand Down
224 changes: 110 additions & 114 deletions src/tapes/rolling_tape.rs
Original file line number Diff line number Diff line change
Expand Up @@ -71,6 +71,10 @@ impl BlobTape for RollingBlobTape {
return Err(io::Error::other("file_size must not be 0."));
}

if name == "metadata" {
return Err(io::Error::other("The tape name `metadata` is reserved."));
}

let path = config.dir.join("tapes").join(name);

if tape_metadata.is_none() {
Expand All @@ -96,6 +100,12 @@ impl BlobTape for RollingBlobTape {
return Err(io::Error::other("File in rolling tapes has invalid name"));
};

if tape_metadata
.is_some_and(|m| index > offset_to_file_index(m.len, config.file_size))
{
continue;
}

let rolling_tape_file = RollingTapeFile::open(&path, index)?;

files.push(rolling_tape_file);
Expand All @@ -108,19 +118,7 @@ impl BlobTape for RollingBlobTape {
));
}

std::fs::create_dir_all(path.clone())?;

let first_file = RollingTapeFile::new(
&path,
offset_to_file_index(config.start_index, config.file_size),
)?;

return Ok(Self {
name,
files: Arc::new(RwLock::new(VecDeque::from([first_file]))),
dir: path,
file_size: config.file_size,
});
std::fs::create_dir_all(&path)?;
}
Err(e) => return Err(e),
}
Expand Down Expand Up @@ -151,32 +149,18 @@ impl BlobTape for RollingBlobTape {
}

fn read_bytes(&self, mut offset: u64, mut buf: &mut [u8]) -> io::Result<()> {
let files = self.files.read();

let files = files
.iter()
.filter(|f| {
let file_start = file_index_to_offset(f.file_index, self.file_size);
let file_end = file_start + self.file_size;

file_start < offset + buf.len() as u64 && file_end > offset
})
.cloned()
.collect::<Vec<_>>();
while !buf.is_empty() {
let file_index = offset_to_file_index(offset, self.file_size);

for file in files {
let next_file_start =
file_index_to_offset(file.file_index, self.file_size) + self.file_size;
let file = file_at(&self.files.read(), file_index)
.ok_or_else(|| io::Error::other("Tape file not found"))?;

let bytes_left_on_file = next_file_start - offset;
let offset_in_file = offset - file_index_to_offset(file_index, self.file_size);
let bytes_left_on_file = self.file_size - offset_in_file;

let bytes_to_read = min(buf.len(), bytes_left_on_file.try_into().unwrap());

read_exact_at_file(
&file.file,
&mut buf[..bytes_to_read],
offset - file_index_to_offset(file.file_index, self.file_size),
)?;
read_exact_at_file(&file, &mut buf[..bytes_to_read], offset_in_file)?;

offset += u64::try_from(bytes_to_read).unwrap();
buf = &mut buf[bytes_to_read..];
Expand All @@ -186,30 +170,29 @@ impl BlobTape for RollingBlobTape {
}

fn writer(&self, len: u64) -> io::Result<Self::Writer> {
let files = self.files.read();

let currently_writing_idx = match files
.iter()
.enumerate()
.find(|(_, tape)| file_index_to_offset(tape.file_index, self.file_size) > len)
.map(|(i, _)| i)
.unwrap_or(files.len())
.checked_sub(1)
{
Some(i) => i,
None => {
unreachable!("Tape files must contain at least one tape file.");
let current_file_idx = offset_to_file_index(len, self.file_size);
let first_file_touched = {
let files = self.files.read();
let slot = slot_for(&files, current_file_idx);
let current_exists = files
.get(slot)
.is_some_and(|f| f.file_index == current_file_idx);

if !current_exists && let Some(prev) = slot.checked_sub(1) {
// The previous file could have unsynced writes.
files[prev].file_index
} else {
current_file_idx
}
};
Comment on lines +174 to +187

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

This confused me, couldn't we just do first_file_touched: current_file_idx without all this?

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

This ensures that if the current file is missing, we sync the previous one as it might not have been synced yet

Made it easier to reason about

@Boog900 Boog900 Sep 12, 2026 •

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Couldn't there be more than 1 previous file that we haven't synced yet? If we write 2 files but only buffer their writes then this will only sync one of the 2 right?

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

I'll merge this anyway as it is harmless


drop(files);

Ok(RollingBlobTapeWriter {
files: self.files.clone(),
dir: self.dir.clone(),
file_size: self.file_size,
currently_writing_idx,
first_file_touched: currently_writing_idx,
current_file_idx,
first_file_touched,
current_file: None,
len,
is_truncation: None,
})
Expand All @@ -236,7 +219,7 @@ impl RollingTapeFile {
.write(true)
.read(true)
.create(true)
.truncate(true)
.truncate(false)
.open(&file_path)?,
);

Expand Down Expand Up @@ -266,23 +249,26 @@ pub struct RollingBlobTapeWriter {
files: Arc<RwLock<VecDeque<RollingTapeFile>>>,
dir: PathBuf,
file_size: u64,
currently_writing_idx: usize,
first_file_touched: usize,
current_file_idx: u64,
first_file_touched: u64,
current_file: Option<Arc<File>>,
len: u64,
is_truncation: Option<bool>,
}

impl RollingBlobTapeWriter {
fn make_new_file(&self) -> io::Result<RollingTapeFile> {
let next_index = self
.files
.read()
.back()
.map_or_default(|f| f.file_index + 1);

let file = RollingTapeFile::new(&self.dir, next_index)?;

self.files.write().push_back(file.clone());
fn make_new_file(&self, file_index: u64) -> io::Result<RollingTapeFile> {
let file = RollingTapeFile::new(&self.dir, file_index)?;

let mut files = self.files.write();
let slot = slot_for(&files, file_index);
debug_assert!(
files.get(slot).is_none_or(|f| f.file_index != file_index),
"rolling tape file {file_index} already exists"
);
Comment thread
redsh4de marked this conversation as resolved.
// We can clean up files here that have fallen off from the top due to a truncation here.
files.truncate(slot);
files.push_back(file.clone());

Ok(file)
}
Expand All @@ -296,36 +282,30 @@ impl BlobTapeWriter for RollingBlobTapeWriter {
let idx = self.len;

while !buf.is_empty() {
let files = self.files.read();
let tape_file = match files.get(self.currently_writing_idx) {
Some(tape_file) => {
let tape_file = tape_file.clone();
drop(files);
tape_file
}
None => {
drop(files);
self.make_new_file()?
}
};
let file_index = offset_to_file_index(self.len, self.file_size);

if self.current_file.is_none() || file_index != self.current_file_idx {
let file = file_at(&self.files.read(), file_index);
self.current_file = Some(match file {
Some(file) => file,
None => self.make_new_file(file_index)?.file,
});
self.current_file_idx = file_index;
}

let bytes_left_on_first_tape = self.file_size
- (self.len - file_index_to_offset(tape_file.file_index, self.file_size));
let offset_in_file = self.len - file_index_to_offset(file_index, self.file_size);
let bytes_left_on_first_tape = self.file_size - offset_in_file;

let bytes_to_write = min(buf.len(), bytes_left_on_first_tape.try_into().unwrap());

write_all_at(
&tape_file.file,
self.current_file.as_ref().unwrap(),
&buf[..bytes_to_write],
self.len - file_index_to_offset(tape_file.file_index, self.file_size),
offset_in_file,
)?;

self.len += u64::try_from(bytes_to_write).unwrap();
buf = &buf[bytes_to_write..];

if !buf.is_empty() {
self.currently_writing_idx += 1;
}
}

Ok(idx)
Expand All @@ -334,25 +314,13 @@ impl BlobTapeWriter for RollingBlobTapeWriter {
fn flush(&mut self, persistence: Persistence) -> io::Result<()> {
if self.is_truncation.is_some_and(|x| x) {
let mut files = self.files.write();
let first = slot_for(&files, offset_to_file_index(self.len, self.file_size));

let mut file_index = files.front().map(|f| f.file_index).unwrap_or_default();

let mut start_index = file_index_to_offset(file_index, self.file_size);

while start_index > self.len {
files.push_front(RollingTapeFile::new(&self.dir, file_index - 1)?);

file_index -= 1;
start_index -= self.file_size;
}

for file in files.iter_mut() {
if file_index_to_offset(file.file_index, self.file_size) <= self.len
&& self.len
< file_index_to_offset(file.file_index, self.file_size) + self.file_size
{
file.out_of_range_at_epoch = None;
}
for file in files
.range_mut(first..)
.take_while(|file| file.out_of_range_at_epoch.is_some())
{
file.out_of_range_at_epoch = None;
}

return Ok(());
Expand All @@ -362,8 +330,10 @@ impl BlobTapeWriter for RollingBlobTapeWriter {
return Ok(());
}

for i in self.first_file_touched..(self.currently_writing_idx + 1) {
let file = self.files.read().get(i).unwrap().file.clone();
for file_index in self.first_file_touched..=self.current_file_idx {
let Some(file) = file_at(&self.files.read(), file_index) else {
continue;
};

match persistence {
Persistence::Buffer => (),
Expand Down Expand Up @@ -413,17 +383,14 @@ fn remove_old_files(
) -> io::Result<()> {
let mut files = files.write();

let mut i = 1;
while let Some(file_2) = files.get(i) {
if file_index_to_offset(file_2.file_index, file_size) <= metadata.start {
files
.get_mut(i - 1)
.unwrap()
.out_of_range_at_epoch
.get_or_insert(current_epoch);
}
let end = slot_for(&files, offset_to_file_index(metadata.start, file_size));

i += 1;
for file in files
.range_mut(..end)
.rev()
.take_while(|file| file.out_of_range_at_epoch.is_none())
{
file.out_of_range_at_epoch = Some(current_epoch);
}

while let Some(file) = files.pop_front_if(|file| {
Expand All @@ -437,6 +404,35 @@ fn remove_old_files(
Ok(())
}

fn slot_for(files: &VecDeque<RollingTapeFile>, file_index: u64) -> usize {
let (Some(first), Some(last)) = (files.front(), files.back()) else {
return 0;
};

if file_index <= first.file_index {
return 0;
}

if file_index > last.file_index {
return files.len();
}

if let Ok(slot) = usize::try_from(file_index - first.file_index)
&& files.get(slot).is_some_and(|f| f.file_index == file_index)
{
return slot;
}

files.partition_point(|f| f.file_index < file_index)
}

fn file_at(files: &VecDeque<RollingTapeFile>, file_index: u64) -> Option<Arc<File>> {
files
.get(slot_for(files, file_index))
.filter(|file| file.file_index == file_index)
.map(|file| Arc::clone(&file.file))
}

fn offset_to_file_index(offset: u64, file_size: u64) -> u64 {
offset / file_size
}
Expand Down
Loading