pub(crate) mod append_only;
pub mod iter;
pub mod read_only;
#[cfg(test)]
mod tests;
use std::path::{Path, PathBuf};
use ahash::{AHashMap, AHashSet};
use crate::common::generic_consts::Random;
use crate::common::mmap::{Advice, AdviceSetting, create_and_ensure_length};
use crate::common::universal_io::{
CachedReadFs, OpenOptions, Populate, ReadRange, UniversalIoError, UniversalRead,
UniversalReadFs, UniversalWrite, UserData,
};
use smallvec::SmallVec;
pub use self::iter::{Iter, PointerItem};
pub use self::read_only::ReadOnlyTracker;
use crate::blobstore::Result;
use crate::blobstore::error::BlobstoreError;
pub type PointOffset = u32;
pub type BlockOffset = u32;
pub type PageId = u32;
fn tracker_open_options(populate: Populate, writeable: bool) -> OpenOptions {
OpenOptions {
writeable,
need_sequential: false,
populate,
advice: AdviceSetting::Advice(Advice::Random),
}
}
#[derive(Debug, Copy, Clone, bytemuck::Pod, bytemuck::Zeroable)]
#[repr(C)]
pub(crate) struct OptionalPointer {
discriminant: u32,
value: ValuePointer,
}
impl From<Option<ValuePointer>> for OptionalPointer {
fn from(value: Option<ValuePointer>) -> Self {
match value {
Some(value) => Self::some(value),
None => Self::none(),
}
}
}
impl OptionalPointer {
const OPTIONAL_NONE: u32 = 0;
const OPTIONAL_SOME: u32 = 1;
pub fn none() -> Self {
Self {
discriminant: Self::OPTIONAL_NONE,
value: ValuePointer::new(0, 0, 0),
}
}
pub const fn some(value: ValuePointer) -> Self {
Self {
discriminant: Self::OPTIONAL_SOME,
value,
}
}
pub fn to_option(self) -> Option<ValuePointer> {
if self.discriminant == Self::OPTIONAL_NONE {
None
} else {
Some(self.value)
}
}
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash, bytemuck::Pod, bytemuck::Zeroable)]
#[repr(C)]
pub struct ValuePointer {
pub page_id: PageId,
pub block_offset: BlockOffset,
pub length: u32,
}
impl ValuePointer {
pub fn new(page_id: PageId, block_offset: BlockOffset, length: u32) -> Self {
Self {
page_id,
block_offset,
length,
}
}
}
pub trait TrackerRead<S: UniversalRead> {
fn max_point_offset(&self) -> Result<PointOffset>;
fn get(&self, point_offset: PointOffset) -> Result<Option<ValuePointer>>;
fn iter<U, I>(&self, point_offsets: I) -> Result<Iter<'_, U, I, S>>
where
U: UserData,
I: Iterator<Item = (U, PointOffset)>;
}
fn read_slot<S: UniversalRead>(
storage: &S,
point_offset: PointOffset,
) -> Result<Option<ValuePointer>> {
let start_offset =
size_of::<TrackerHeader>() + point_offset as usize * size_of::<OptionalPointer>();
let end_offset = start_offset + size_of::<OptionalPointer>();
let storage_len = storage.len::<u8>()?;
if end_offset as u64 > storage_len {
return Ok(None);
}
let opt = storage.read::<_, OptionalPointer>(ReadRange::one(start_offset as u64), Random)?[0];
Ok(opt.to_option())
}
#[derive(Debug, Default, Clone, PartialEq)]
pub(crate) struct PointerUpdates {
current: Option<ValuePointer>,
to_free: SmallVec<[ValuePointer; 1]>,
}
impl PointerUpdates {
fn set(&mut self, pointer: ValuePointer) {
if self.current == Some(pointer) {
debug_assert!(false, "we should not set the same point twice");
return;
}
if let Some(old_pointer) = self.current.replace(pointer) {
self.to_free.push(old_pointer);
debug_assert_eq!(
self.to_free.iter().copied().collect::<AHashSet<_>>().len(),
self.to_free.len(),
"should not have duplicate pointers to free",
);
}
debug_assert!(
!self.to_free.contains(&pointer),
"old list cannot contain pointer we just set",
);
}
fn unset(&mut self, pointer: ValuePointer) {
let old_pointer = self.current.take();
debug_assert!(
old_pointer.is_none_or(|p| p == pointer),
"new unset pointer should match with current one, if any",
);
if let Some(old_pointer) = old_pointer
&& old_pointer != pointer
{
self.to_free.push(old_pointer);
}
self.to_free.push(pointer);
debug_assert_eq!(
self.to_free.iter().copied().collect::<AHashSet<_>>().len(),
self.to_free.len(),
"should not have duplicate pointers to free",
);
}
fn is_empty(&self) -> bool {
self.current.is_none() && self.to_free.is_empty()
}
fn drain_persisted(&mut self, persisted: &Self) -> bool {
debug_assert!(!self.is_empty(), "must have at least one pointer");
debug_assert!(
!persisted.is_empty(),
"persisted must have at least one pointer",
);
if self == persisted {
*self = Self::default();
return true;
}
let Self {
current: previous_current,
to_free: freed,
} = persisted;
if let (Some(current), Some(previous_current)) = (self.current, *previous_current)
&& current == previous_current
{
self.current.take();
}
self.to_free.retain(|pointer| !freed.contains(pointer));
self.is_empty()
}
}
#[derive(Debug, Default, Clone, Copy, bytemuck::Pod, bytemuck::Zeroable)]
#[repr(C)]
struct TrackerHeader {
next_pointer_offset: u32,
}
#[derive(Debug)]
pub struct Tracker<S> {
path: PathBuf,
header: TrackerHeader,
storage: S,
pub(super) pending_updates: AHashMap<PointOffset, PointerUpdates>,
next_pointer_offset: PointOffset,
}
impl<S> Tracker<S> {
const FILE_NAME: &'static str = "tracker.dat";
fn tracker_file_name(path: &Path) -> PathBuf {
path.join(Self::FILE_NAME)
}
pub fn files(&self) -> Vec<PathBuf> {
vec![self.path.clone()]
}
pub fn pointer_count(&self) -> u32 {
self.next_pointer_offset
}
}
impl<S: UniversalRead> Tracker<S> {
pub fn preopen<Fs: CachedReadFs<File = S>>(
fs: &Fs,
path: &Path,
populate: Populate,
) -> Result<()> {
let path = Self::tracker_file_name(path);
let populate = populate.or_partial(0..size_of::<TrackerHeader>() as u64);
fs.schedule_prefetch(&path, Some(tracker_open_options(populate, false)), None)?;
Ok(())
}
pub fn open<Fs: UniversalReadFs<File = S>>(
fs: &Fs,
path: &Path,
populate: Populate,
writeable: bool,
) -> Result<Self> {
let path = Self::tracker_file_name(path);
let storage = Self::open_storage(fs, &path, populate, writeable)?;
let header: TrackerHeader = Self::read_header(&storage)?;
let pending_updates = AHashMap::new();
Ok(Self {
next_pointer_offset: header.next_pointer_offset,
path,
header,
storage,
pending_updates,
})
}
fn read_header(storage: &S) -> Result<TrackerHeader> {
let header = storage.read(ReadRange::one(0), Random)?[0];
Ok(header)
}
fn open_storage<Fs: UniversalReadFs<File = S>>(
fs: &Fs,
path: &Path,
populate: Populate,
writeable: bool,
) -> Result<S> {
let storage = match fs.open(
path,
tracker_open_options(populate, writeable),
Default::default(),
) {
Err(UniversalIoError::NotFound { .. }) => {
return Err(BlobstoreError::service_error(format!(
"Tracker file does not exist: {}",
path.display()
)));
}
other => other?,
};
Ok(storage)
}
fn get_raw(&self, point_offset: PointOffset) -> Result<Option<ValuePointer>> {
read_slot(&self.storage, point_offset)
}
pub fn get(&self, point_offset: PointOffset) -> Result<Option<ValuePointer>> {
match self.pending_updates.get(&point_offset) {
Some(pending) if pending.is_empty() => {
debug_assert!(false, "pending updates must not be empty");
self.get_raw(point_offset)
}
Some(pending) => Ok(pending.current),
None => self.get_raw(point_offset),
}
}
pub fn iter<U, I>(&self, point_offsets: I) -> Result<Iter<'_, U, I, S>>
where
U: UserData,
I: Iterator<Item = (U, PointOffset)>,
{
Iter::new(point_offsets, &self.storage, &self.pending_updates)
}
pub fn has_pointer(&self, point_offset: PointOffset) -> Result<bool> {
Ok(self.get(point_offset)?.is_some())
}
pub fn populate(&self) -> Result<()> {
self.storage.populate().map_err(Into::into)
}
}
impl<S: UniversalRead> TrackerRead<S> for Tracker<S> {
fn max_point_offset(&self) -> Result<PointOffset> {
Ok(self.pointer_count())
}
fn get(&self, point_offset: PointOffset) -> Result<Option<ValuePointer>> {
Tracker::get(self, point_offset)
}
fn iter<U, I>(&self, point_offsets: I) -> Result<Iter<'_, U, I, S>>
where
U: UserData,
I: Iterator<Item = (U, PointOffset)>,
{
Tracker::iter(self, point_offsets)
}
}
impl<S> Tracker<S>
where
S: UniversalWrite,
{
const DEFAULT_SIZE: usize = 1024 * 1024;
pub fn new(fs: &S::Fs, path: &Path, size_hint: Option<usize>) -> Result<Self> {
let path = Self::tracker_file_name(path);
let size = size_hint.unwrap_or(Self::DEFAULT_SIZE).next_power_of_two();
assert!(
size > std::mem::size_of::<TrackerHeader>(),
"Size hint is too small"
);
create_and_ensure_length(&path, size)?;
let storage = fs.open(
&path,
tracker_open_options(Populate::No, true),
Default::default(),
)?;
let header = TrackerHeader::default();
let pending_updates = AHashMap::new();
let mut page_tracker = Self {
path,
header,
storage,
pending_updates,
next_pointer_offset: 0,
};
page_tracker.write_header()?;
Ok(page_tracker)
}
#[must_use = "The old pointers need to be freed in the bitmask"]
pub fn write_pending(
&mut self,
pending_updates: AHashMap<PointOffset, PointerUpdates>,
) -> Result<Vec<ValuePointer>> {
let mut old_pointers = Vec::new();
for (point_offset, updates) in pending_updates {
match updates.current {
Some(new_pointer) => {
if let Some(old_pointer) = self.get_raw(point_offset)? {
old_pointers.push(old_pointer);
}
self.persist_pointer(point_offset, Some(new_pointer))?;
}
None => self.persist_pointer(point_offset, None)?,
}
old_pointers.extend(&updates.to_free);
if let Some(latest_updates) = self.pending_updates.get_mut(&point_offset) {
let is_empty = latest_updates.drain_persisted(&updates);
if is_empty {
let prev = self.pending_updates.remove(&point_offset);
if let Some(prev) = prev {
debug_assert!(
prev.is_empty(),
"remove pending element should be empty but got {prev:?}"
);
}
}
}
}
self.write_pointer_count()?;
Ok(old_pointers)
}
pub fn flusher(&self) -> crate::blobstore::blobstore::Flusher {
let inner = self.storage.flusher();
Box::new(move || inner().map_err(Into::into))
}
fn write_header(&mut self) -> Result<()> {
self.storage.write(0, &[self.header])?;
Ok(())
}
fn persist_pointer(
&mut self,
point_offset: PointOffset,
pointer: Option<ValuePointer>,
) -> Result<()> {
let storage_len = self.storage.len::<u8>()? as usize;
if pointer.is_none() && point_offset as usize >= storage_len {
return Ok(());
}
let point_offset = point_offset as usize;
let start_offset = size_of::<TrackerHeader>() + point_offset * size_of::<OptionalPointer>();
let end_offset = start_offset + size_of::<OptionalPointer>();
if storage_len < end_offset {
self.storage.flusher()()?;
let new_size = end_offset.next_power_of_two();
create_and_ensure_length(&self.path, new_size)?;
self.storage.reopen()?;
}
let pointer = OptionalPointer::from(pointer);
self.storage.write(start_offset as u64, &[pointer])?;
Ok(())
}
fn write_pointer_count(&mut self) -> Result<()> {
self.header.next_pointer_offset = self.next_pointer_offset;
self.write_header()
}
pub fn set(&mut self, point_offset: PointOffset, value_pointer: ValuePointer) {
self.pending_updates
.entry(point_offset)
.or_default()
.set(value_pointer);
self.next_pointer_offset = self.next_pointer_offset.max(point_offset + 1);
}
pub fn unset(&mut self, point_offset: PointOffset) -> Result<Option<ValuePointer>> {
let pointer_opt = self.get(point_offset)?;
if let Some(pointer) = pointer_opt {
self.pending_updates
.entry(point_offset)
.or_default()
.unset(pointer);
}
Ok(pointer_opt)
}
}