use std::collections::BTreeSet;
use std::path::{Path, PathBuf};
use io_uring::types::FsyncFlags;
use super::super::batch_fsync::{new_ring, sync_paths};
use super::super::manifest::{self, Manifest};
use super::Table;
use gnitz_zset::repr::StorageError;
pub(super) struct FlushWork {
bytes: Vec<u8>,
dirs: Vec<PathBuf>,
}
pub(crate) fn flush_barrier<'a>(
tables: impl IntoIterator<Item = &'a mut Table>,
checkpoint_mark: u64,
) -> Result<(), StorageError> {
let mut work: Vec<(&'a mut Table, FlushWork)> = Vec::new();
for t in tables {
if let Some(w) = t.flush_prepare(checkpoint_mark)? {
work.push((t, w));
}
}
if work.is_empty() {
return Ok(());
}
let mut ring = new_ring()?;
sync_paths(
&mut ring,
work.iter().flat_map(|(t, _)| {
t.shard_index
.unsynced_paths()
.chain([manifest::staging_path(&t.shard_index.output_dir)])
}),
FsyncFlags::DATASYNC,
)?;
let mut dirs = BTreeSet::new();
for (t, w) in &mut work {
manifest::commit(&t.shard_index.output_dir)?;
t.shard_index.mark_published();
dirs.extend(std::mem::take(&mut w.dirs));
}
sync_paths(&mut ring, &dirs, FsyncFlags::empty())?;
for (t, w) in work {
t.durable_manifest = Some(w.bytes);
t.shard_index.unlink_retired();
}
Ok(())
}
impl Table {
#[cfg(test)]
pub(crate) fn flush(&mut self) -> Result<(), StorageError> {
assert!(!self.is_rederived(), "the base round never visits a rederived table");
flush_barrier([&mut *self], 0)
}
fn fold_memtable_into_ram_tier(&mut self) {
self.memtable.drain_into(&mut self.ram_tier, &self.shard_index.schema);
}
pub(crate) fn fold_to_ram(&mut self) -> Result<(), StorageError> {
self.fold_memtable_into_ram_tier();
if self.held_in_ram || !self.ram_tier.is_full() {
return Ok(());
}
self.ram_tier.fold(&self.shard_index.schema);
if !self.ram_tier.is_crowded() {
return Ok(());
}
self.spill_ram_tier()
}
pub(super) fn flush_prepare(&mut self, checkpoint_mark: u64) -> Result<Option<FlushWork>, StorageError> {
self.fold_memtable_into_ram_tier();
self.spill_ram_tier()?;
let schema = self.shard_index.schema;
self.pending
.spill(&schema, |run| self.shard_index.append_pending_run(run))?;
let done = self.shard_index.finish_fold()?;
self.evicted(done.evicted);
let bytes = manifest::encode(&Manifest {
checkpoint_mark,
caller_record: self.caller_record.clone(),
shards: self.shard_index.shard_set(),
});
if self.durable_manifest.as_deref() == Some(&bytes[..]) {
debug_assert!(
self.shard_index.unsynced_paths().next().is_none(),
"a durable manifest names an unsynced shard"
);
return Ok(None);
}
let first_publish = self.durable_manifest.is_none();
self.durable_manifest = None;
manifest::prepare(&self.shard_index.output_dir, &bytes)?;
let entry_dirs = if first_publish { 2 } else { 0 };
let dirs = Path::new(&self.shard_index.output_dir)
.ancestors()
.take(1 + entry_dirs)
.filter(|d| !d.as_os_str().is_empty())
.map(Path::to_path_buf)
.collect();
Ok(Some(FlushWork { bytes, dirs }))
}
fn spill_ram_tier(&mut self) -> Result<(), StorageError> {
let schema = self.shard_index.schema;
self.ram_tier
.spill(&schema, |run| self.shard_index.append_l0_run(run))?;
Ok(())
}
}