use std::ops::Range;
use std::path::{Path, PathBuf};
use crate::common::generic_consts::{AccessPattern, Random};
use crate::common::mmap::{Advice, AdviceSetting};
use crate::common::universal_io::{
CachedReadFs, IsNotFound, OpenOptions, Populate, ReadPipeline, ReadRange, UniversalAppend,
UniversalIoError, UniversalRead, UniversalReadFs, UniversalWriteFileOps, UserData,
};
use crate::blobstore::Result;
use crate::blobstore::blobstore::Flusher;
use crate::blobstore::error::BlobstoreError;
use crate::blobstore::tracker::{OptionalPointer, PointOffset, ValuePointer};
const FILE_NAME: &str = "log_tracker.dat";
const ENTRY_SIZE: u64 = size_of::<OptionalPointer>() as u64;
#[derive(Debug)]
pub(crate) struct AppendOnlyTracker<S> {
path: PathBuf,
file: S,
persisted_count: PointOffset,
pending: Vec<OptionalPointer>,
}
impl<S: UniversalRead> AppendOnlyTracker<S> {
fn tracker_file_name(dir: &Path) -> PathBuf {
dir.join(FILE_NAME)
}
fn open_options(populate: Populate, writeable: bool) -> OpenOptions {
OpenOptions {
writeable,
need_sequential: false,
populate,
advice: AdviceSetting::Advice(Advice::Random),
}
}
pub fn preopen<Fs: CachedReadFs<File = S>>(
fs: &Fs,
dir: &Path,
populate: Populate,
) -> Result<()> {
fs.schedule_prefetch(
&Self::tracker_file_name(dir),
Some(Self::open_options(populate, false)),
None,
)?;
Ok(())
}
fn open_file<Fs: UniversalReadFs<File = S>>(
fs: &Fs,
path: &Path,
populate: Populate,
writeable: bool,
) -> Result<S> {
fs.open(
path,
Self::open_options(populate, writeable),
Default::default(),
)
.map_err(|err| {
if err.is_not_found() {
BlobstoreError::service_error(format!(
"Append-only tracker file does not exist: {}",
path.display(),
))
} else {
BlobstoreError::from(err)
}
})
}
pub fn open_read_only<Fs: UniversalReadFs<File = S>>(
fs: &Fs,
dir: &Path,
populate: Populate,
) -> Result<Self> {
let path = Self::tracker_file_name(dir);
let file = Self::open_file(fs, &path, populate, false)?;
let len = file.len::<u8>()?;
Ok(Self {
path,
file,
persisted_count: count_from_len(len)?,
pending: Vec::new(),
})
}
pub fn files(&self) -> Vec<PathBuf> {
vec![self.path.clone()]
}
pub fn populate(&self) -> Result<()> {
self.file.populate().map_err(Into::into)
}
pub fn clear_cache(&self) -> Result<()> {
self.file.clear_ram_cache().map_err(Into::into)
}
pub fn pointer_count(&self) -> PointOffset {
self.persisted_count + self.pending.len() as PointOffset
}
pub fn get<P: AccessPattern>(&self, point_offset: PointOffset) -> Result<Option<ValuePointer>> {
if point_offset >= self.pointer_count() {
return Ok(None);
}
if point_offset >= self.persisted_count {
let pending_index = (point_offset - self.persisted_count) as usize;
return Ok(self.pending[pending_index].to_option());
}
let range = ReadRange::one(u64::from(point_offset) * ENTRY_SIZE);
let pointer = self.file.read::<_, OptionalPointer>(range, P::default())?[0];
Ok(pointer.to_option())
}
pub fn get_range<P: AccessPattern>(
&self,
point_offsets: Range<PointOffset>,
) -> Result<Vec<Option<ValuePointer>>> {
let mut pointers = Vec::with_capacity(point_offsets.len());
let start = point_offsets.start;
let end = point_offsets.end.min(self.pointer_count());
let persisted_end = end.min(self.persisted_count);
if start < persisted_end {
let range = ReadRange {
byte_offset: u64::from(start) * ENTRY_SIZE,
length: u64::from(persisted_end - start),
};
let entries = self.file.read::<_, OptionalPointer>(range, P::default())?;
pointers.extend(entries.iter().map(|entry| entry.to_option()));
}
for point_offset in start.max(self.persisted_count)..end {
let pending_index = (point_offset - self.persisted_count) as usize;
pointers.push(self.pending[pending_index].to_option());
}
pointers.resize(point_offsets.len(), None);
Ok(pointers)
}
pub fn iter<U, I>(&self, point_offsets: I) -> Result<Iter<'_, U, I, S>>
where
U: UserData,
I: Iterator<Item = (U, PointOffset)>,
{
Ok(Iter {
point_offsets,
tracker: self,
pipeline: S::ReadPipeline::new()?,
})
}
pub fn reload_count(&mut self) -> Result<PendingReload> {
debug_assert!(
self.pending.is_empty(),
"live reload must only be used on read-only instances",
);
self.file.reopen()?;
let len = self.file.len::<u8>()?;
let new_count = count_from_len(len)?;
if new_count < self.persisted_count {
return Err(BlobstoreError::service_error(format!(
"live reload cannot decrease mapping count, possible data loss: old count {}, new count {new_count}",
self.persisted_count,
)));
}
Ok(PendingReload { count: new_count })
}
pub fn commit_reload(&mut self, reload: PendingReload) {
let PendingReload { count } = reload;
debug_assert!(
count >= self.persisted_count,
"a committed reload must not decrease the mapping count",
);
self.persisted_count = count;
}
}
#[must_use = "an observed reload only becomes visible once committed"]
#[derive(Debug, Copy, Clone)]
pub(crate) struct PendingReload {
count: PointOffset,
}
#[cfg(test)]
impl PendingReload {
pub fn count(self) -> PointOffset {
self.count
}
}
impl<S: UniversalAppend> AppendOnlyTracker<S> {
pub fn new(fs: &S::Fs, dir: &Path) -> Result<Self> {
let path = Self::tracker_file_name(dir);
fs.create(&path, 0)?;
let file = fs.open(
&path,
Self::open_options(Populate::No, true),
Default::default(),
)?;
Ok(Self {
path,
file,
persisted_count: 0,
pending: Vec::new(),
})
}
pub fn open_writable(fs: &S::Fs, dir: &Path, populate: Populate) -> Result<Self> {
let path = Self::tracker_file_name(dir);
let mut file = Self::open_file(fs, &path, populate, true)?;
let len = file.len::<u8>()?;
let aligned_len = len - (len % ENTRY_SIZE);
if aligned_len != len {
drop(file);
fs.create(&path, aligned_len as usize)?;
file = Self::open_file(fs, &path, populate, true)?;
}
Ok(Self {
path,
file,
persisted_count: count_from_len(aligned_len)?,
pending: Vec::new(),
})
}
pub fn set(&mut self, point_offset: PointOffset, pointer: ValuePointer) -> Result<()> {
let next = self.pointer_count();
if point_offset < next {
return Err(BlobstoreError::unsupported_operation(format!(
"cannot set mapping for point offset {point_offset}, the tracker is append-only \
and requires monotonically increasing point offsets, the next allowed point \
offset is {next}",
)));
}
let pending_index = (point_offset - self.persisted_count) as usize;
self.pending.resize(pending_index, OptionalPointer::none());
self.pending.push(OptionalPointer::some(pointer));
Ok(())
}
pub fn write_pending(&mut self, target: PointOffset) -> Result<()> {
if target <= self.persisted_count {
return Ok(());
}
let count = (target - self.persisted_count) as usize;
debug_assert!(
count <= self.pending.len(),
"flush target exceeds pending mappings",
);
let count = count.min(self.pending.len());
let offset = u64::from(self.persisted_count) * ENTRY_SIZE;
let end = offset + count as u64 * ENTRY_SIZE;
match self.file.append(offset, &self.pending[..count]) {
Ok(()) => {}
Err(UniversalIoError::AppendOffsetConflict { .. }) => {
self.file.reopen()?;
let len = self.file.len::<u8>()?;
if len != end {
return Err(BlobstoreError::service_error(format!(
"append-only tracker file {} was modified outside this writer: it ends \
at byte {len}, expected {offset} before or {end} after the append",
self.path.display(),
)));
}
}
Err(err) => return Err(err.into()),
}
self.pending.drain(..count);
self.persisted_count += count as PointOffset;
Ok(())
}
pub fn flusher(&self) -> Flusher {
let flusher = self.file.flusher();
Box::new(move || flusher().map_err(BlobstoreError::from))
}
}
pub(crate) struct Iter<'a, U, I, S>
where
U: UserData,
I: Iterator<Item = (U, PointOffset)>,
S: UniversalRead,
{
point_offsets: I,
tracker: &'a AppendOnlyTracker<S>,
pipeline: S::ReadPipeline<'a, U>,
}
impl<'a, U, I, S> Iterator for Iter<'a, U, I, S>
where
U: UserData,
I: Iterator<Item = (U, PointOffset)>,
S: UniversalRead,
{
type Item = Result<(U, Option<ValuePointer>)>;
fn next(&mut self) -> Option<Self::Item> {
while self.pipeline.can_schedule()
&& let Some((user_data, point_offset)) = self.point_offsets.next()
{
if point_offset >= self.tracker.persisted_count {
let pending_index = (point_offset - self.tracker.persisted_count) as usize;
let pointer = self
.tracker
.pending
.get(pending_index)
.and_then(|entry| entry.to_option());
return Some(Ok((user_data, pointer)));
}
let start = u64::from(point_offset) * ENTRY_SIZE;
let result = self.pipeline.schedule::<Random>(
user_data,
&self.tracker.file,
start..start + ENTRY_SIZE,
align_of::<OptionalPointer>(),
);
if let Err(err) = result {
return Some(Err(err.into()));
}
}
let result = self.pipeline.wait_bytemuck::<OptionalPointer>();
let (user_data, entry) = match result {
Ok(entry) => entry?,
Err(err) => return Some(Err(err.into())),
};
let &[entry] = entry.as_ref() else {
unreachable!();
};
Some(Ok((user_data, entry.to_option())))
}
}
fn count_from_len(len: u64) -> Result<PointOffset> {
PointOffset::try_from(len / ENTRY_SIZE).map_err(|_| {
BlobstoreError::service_error(format!(
"append-only tracker file of {len} bytes holds more mappings than supported",
))
})
}
#[cfg(test)]
mod tests {
use std::io::Write as _;
use crate::common::generic_consts::Random;
use crate::common::universal_io::{MmapFile, MmapFs};
use fs_err as fs;
use tempfile::TempDir;
use super::*;
fn empty_tracker() -> (TempDir, AppendOnlyTracker<MmapFile>) {
let dir = TempDir::new().unwrap();
let tracker = AppendOnlyTracker::new(&MmapFs, dir.path()).unwrap();
(dir, tracker)
}
fn pointer(n: u32) -> ValuePointer {
ValuePointer::new(0, n * 2, n * 3 + 1)
}
fn file_len(tracker: &AppendOnlyTracker<MmapFile>) -> u64 {
fs::metadata(&tracker.path).unwrap().len()
}
#[test]
fn test_new_tracker_is_empty() {
let (_dir, tracker) = empty_tracker();
assert_eq!(tracker.pointer_count(), 0);
assert_eq!(tracker.get::<Random>(0).unwrap(), None);
assert_eq!(file_len(&tracker), 0);
}
#[test]
fn test_open_missing_tracker_fails() {
let dir = TempDir::new().unwrap();
assert!(
AppendOnlyTracker::<MmapFile>::open_writable(&MmapFs, dir.path(), Populate::No)
.is_err()
);
assert!(
AppendOnlyTracker::<MmapFile>::open_read_only(&MmapFs, dir.path(), Populate::No)
.is_err()
);
}
#[test]
fn test_set_and_get_pending() {
let (_dir, mut tracker) = empty_tracker();
for n in 0..5 {
tracker.set(n, pointer(n)).unwrap();
}
assert_eq!(tracker.pointer_count(), 5);
for n in 0..5 {
assert_eq!(tracker.get::<Random>(n).unwrap(), Some(pointer(n)));
}
assert_eq!(tracker.get::<Random>(5).unwrap(), None);
assert_eq!(file_len(&tracker), 0);
}
#[test]
fn test_set_rejects_non_monotonic_point_offsets() {
let (_dir, mut tracker) = empty_tracker();
tracker.set(0, pointer(0)).unwrap();
assert!(tracker.set(0, pointer(0)).is_err());
tracker.set(5, pointer(5)).unwrap();
assert!(tracker.set(3, pointer(3)).is_err());
assert!(tracker.set(5, pointer(5)).is_err());
tracker.set(6, pointer(6)).unwrap();
assert_eq!(tracker.pointer_count(), 7);
}
#[test]
fn test_skipped_point_offsets_read_as_none() {
let (_dir, mut tracker) = empty_tracker();
tracker.set(0, pointer(0)).unwrap();
tracker.set(4, pointer(4)).unwrap();
assert_eq!(tracker.pointer_count(), 5);
assert_eq!(tracker.get::<Random>(0).unwrap(), Some(pointer(0)));
for n in 1..4 {
assert_eq!(tracker.get::<Random>(n).unwrap(), None);
}
assert_eq!(tracker.get::<Random>(4).unwrap(), Some(pointer(4)));
tracker.write_pending(tracker.pointer_count()).unwrap();
assert_eq!(file_len(&tracker), 5 * ENTRY_SIZE);
assert_eq!(tracker.get::<Random>(2).unwrap(), None);
}
#[test]
fn test_write_pending_and_reopen() {
let dir = TempDir::new().unwrap();
let mut tracker = AppendOnlyTracker::<MmapFile>::new(&MmapFs, dir.path()).unwrap();
for n in 0..5 {
tracker.set(n, pointer(n)).unwrap();
}
tracker.write_pending(tracker.pointer_count()).unwrap();
tracker.flusher()().unwrap();
assert_eq!(file_len(&tracker), 5 * ENTRY_SIZE);
drop(tracker);
let tracker =
AppendOnlyTracker::<MmapFile>::open_writable(&MmapFs, dir.path(), Populate::No)
.unwrap();
assert_eq!(tracker.pointer_count(), 5);
for n in 0..5 {
assert_eq!(tracker.get::<Random>(n).unwrap(), Some(pointer(n)));
}
assert_eq!(
tracker.get_range::<Random>(0..7).unwrap(),
(0..5)
.map(|n| Some(pointer(n)))
.chain([None, None])
.collect::<Vec<_>>(),
);
}
#[test]
fn test_partial_write_pending() {
let (_dir, mut tracker) = empty_tracker();
for n in 0..5 {
tracker.set(n, pointer(n)).unwrap();
}
tracker.write_pending(3).unwrap();
assert_eq!(file_len(&tracker), 3 * ENTRY_SIZE);
assert_eq!(tracker.pointer_count(), 5);
for n in 0..5 {
assert_eq!(tracker.get::<Random>(n).unwrap(), Some(pointer(n)));
}
tracker.write_pending(5).unwrap();
assert_eq!(file_len(&tracker), 5 * ENTRY_SIZE);
for n in 0..5 {
assert_eq!(tracker.get::<Random>(n).unwrap(), Some(pointer(n)));
}
}
#[test]
fn test_stale_flush_is_noop() {
let (_dir, mut tracker) = empty_tracker();
for n in 0..3 {
tracker.set(n, pointer(n)).unwrap();
}
tracker.write_pending(3).unwrap();
assert_eq!(file_len(&tracker), 3 * ENTRY_SIZE);
tracker.write_pending(2).unwrap();
tracker.write_pending(3).unwrap();
assert_eq!(file_len(&tracker), 3 * ENTRY_SIZE);
assert_eq!(tracker.pointer_count(), 3);
for n in 0..3 {
assert_eq!(tracker.get::<Random>(n).unwrap(), Some(pointer(n)));
}
}
#[test]
fn test_torn_write_is_ignored_and_truncated() {
let dir = TempDir::new().unwrap();
let mut tracker = AppendOnlyTracker::<MmapFile>::new(&MmapFs, dir.path()).unwrap();
for n in 0..5 {
tracker.set(n, pointer(n)).unwrap();
}
tracker.write_pending(tracker.pointer_count()).unwrap();
let path = tracker.path.clone();
drop(tracker);
let mut file = fs::OpenOptions::new().append(true).open(&path).unwrap();
file.write_all(&[0xAA; 7]).unwrap();
drop(file);
assert_eq!(fs::metadata(&path).unwrap().len(), 5 * ENTRY_SIZE + 7);
let tracker =
AppendOnlyTracker::<MmapFile>::open_read_only(&MmapFs, dir.path(), Populate::No)
.unwrap();
assert_eq!(tracker.pointer_count(), 5);
assert_eq!(tracker.get::<Random>(4).unwrap(), Some(pointer(4)));
assert_eq!(file_len(&tracker), 5 * ENTRY_SIZE + 7);
drop(tracker);
let tracker =
AppendOnlyTracker::<MmapFile>::open_writable(&MmapFs, dir.path(), Populate::No)
.unwrap();
assert_eq!(tracker.pointer_count(), 5);
assert_eq!(tracker.get::<Random>(4).unwrap(), Some(pointer(4)));
assert_eq!(file_len(&tracker), 5 * ENTRY_SIZE);
}
#[test]
fn test_write_pending_adopts_lost_append_on_conflict() {
let (_dir, mut tracker) = empty_tracker();
tracker.set(0, pointer(0)).unwrap();
let entry = OptionalPointer::some(pointer(0));
fs::write(&tracker.path, bytemuck::bytes_of(&entry)).unwrap();
tracker.write_pending(1).unwrap();
assert_eq!(tracker.persisted_count, 1);
assert!(tracker.pending.is_empty());
assert_eq!(
file_len(&tracker),
ENTRY_SIZE,
"the mapping must not be appended twice",
);
assert_eq!(tracker.get::<Random>(0).unwrap(), Some(pointer(0)));
}
#[test]
fn test_write_pending_rejects_foreign_growth() {
let (_dir, mut tracker) = empty_tracker();
tracker.set(0, pointer(0)).unwrap();
fs::write(&tracker.path, [9; 7]).unwrap();
let err = tracker.write_pending(1).unwrap_err();
assert!(matches!(err, BlobstoreError::ServiceError { .. }));
assert_eq!(tracker.persisted_count, 0);
assert_eq!(tracker.pending.len(), 1);
}
#[test]
fn test_live_reload() {
let dir = TempDir::new().unwrap();
let mut writer = AppendOnlyTracker::<MmapFile>::new(&MmapFs, dir.path()).unwrap();
for n in 0..3 {
writer.set(n, pointer(n)).unwrap();
}
writer.write_pending(writer.pointer_count()).unwrap();
let mut reader =
AppendOnlyTracker::<MmapFile>::open_read_only(&MmapFs, dir.path(), Populate::No)
.unwrap();
assert_eq!(reader.pointer_count(), 3);
let reload = reader.reload_count().unwrap();
assert_eq!(reload.count(), 3);
reader.commit_reload(reload);
for n in 3..6 {
writer.set(n, pointer(n)).unwrap();
}
writer.write_pending(writer.pointer_count()).unwrap();
let reload = reader.reload_count().unwrap();
assert_eq!(reload.count(), 6);
assert_eq!(reader.pointer_count(), 3);
assert_eq!(reader.get::<Random>(3).unwrap(), None);
reader.commit_reload(reload);
assert_eq!(reader.pointer_count(), 6);
for n in 0..6 {
assert_eq!(reader.get::<Random>(n).unwrap(), Some(pointer(n)));
}
let reload = reader.reload_count().unwrap();
assert_eq!(reload.count(), 6);
reader.commit_reload(reload);
}
#[test]
fn test_iter_spans_persisted_and_pending() {
let (_dir, mut tracker) = empty_tracker();
for n in 0..3 {
tracker.set(n, pointer(n)).unwrap();
}
tracker.write_pending(3).unwrap();
tracker.set(3, pointer(3)).unwrap();
tracker.set(6, pointer(6)).unwrap();
let requested = [2, 0, 6, 4, 3, 9];
let mut collected = tracker
.iter(requested.iter().map(|&offset| (offset, offset)))
.unwrap()
.map(|result| result.unwrap())
.collect::<Vec<_>>();
collected.sort_by_key(|(point_offset, _)| *point_offset);
assert_eq!(
collected,
vec![
(0, Some(pointer(0))),
(2, Some(pointer(2))),
(3, Some(pointer(3))),
(4, None),
(6, Some(pointer(6))),
(9, None),
],
);
}
#[test]
fn test_get_range_spans_persisted_and_pending() {
let (_dir, mut tracker) = empty_tracker();
for n in 0..3 {
tracker.set(n, pointer(n)).unwrap();
}
tracker.write_pending(3).unwrap();
tracker.set(3, pointer(3)).unwrap();
tracker.set(6, pointer(6)).unwrap();
assert_eq!(
tracker.get_range::<Random>(0..9).unwrap(),
vec![
Some(pointer(0)),
Some(pointer(1)),
Some(pointer(2)),
Some(pointer(3)),
None,
None,
Some(pointer(6)),
None,
None,
],
);
assert_eq!(
tracker.get_range::<Random>(2..4).unwrap(),
vec![Some(pointer(2)), Some(pointer(3)),]
);
assert_eq!(tracker.get_range::<Random>(7..9).unwrap(), vec![None, None]);
#[allow(clippy::reversed_empty_ranges)]
let empty = tracker.get_range::<Random>(3..3).unwrap();
assert!(empty.is_empty());
}
}