use std::collections::HashSet;
use std::rc::Rc;
use super::manifest::{self, ManifestEntry, ShardSet};
use gnitz_wire::PkBuf;
use gnitz_wire::RowSource;
use gnitz_zset::repr::StorageError;
use gnitz_zset::repr::{guard_slot, pk_group_end, MappedShard};
use gnitz_zset::schema::key::{pk_bytes_eq, pk_in_range};
use gnitz_zset::schema::SchemaDescriptor;
mod index;
#[derive(Clone, Copy, PartialEq, Eq, PartialOrd, Ord, Debug)]
pub(crate) enum CompactionKind {
L0Fold,
GuardSplit,
TierFold,
BandCut,
Vertical,
GuardMerge,
Dehydrate,
}
#[cfg(test)]
pub(crate) mod cstats {
use super::CompactionKind;
use std::cell::RefCell;
use std::collections::BTreeMap;
#[derive(Default, Clone)]
pub(crate) struct Phase {
pub(crate) n: usize,
pub(crate) in_bytes: u64,
pub(crate) max_in: u64,
pub(crate) out_bytes: u64,
pub(crate) in_files: usize,
}
thread_local! {
static STATS: RefCell<BTreeMap<CompactionKind, Phase>> = const { RefCell::new(BTreeMap::new()) };
}
pub(super) fn record(kind: CompactionKind, inb: u64, outb: u64, inf: usize) {
STATS.with(|stats| {
let mut stats = stats.borrow_mut();
let p = stats.entry(kind).or_default();
p.n += 1;
p.in_bytes += inb;
p.max_in = p.max_in.max(inb);
p.out_bytes += outb;
p.in_files += inf;
});
}
pub(crate) fn reset() {
STATS.with(|stats| stats.borrow_mut().clear());
}
pub(crate) fn dump() -> BTreeMap<CompactionKind, Phase> {
STATS.with(|stats| stats.borrow().clone())
}
}
#[cfg(test)]
impl ShardIndex {
pub(crate) fn level_shape(&self) -> (usize, [usize; 2]) {
(
self.levels[L0].entries().count(),
[L1, TERMINAL].map(|l| self.levels[l].guards.len()),
)
}
}
const LEVELS: usize = 3;
pub(super) const L0: usize = 0;
pub(super) const L1: usize = 1;
pub(super) const TERMINAL: usize = 2;
pub(super) const L0_COMPACT_THRESHOLD: usize = 4;
const GUARD_FILE_THRESHOLD: usize = 4;
const CANCEL_PERCENT: usize = 25;
const MIN_GUARD_BYTES: u64 = 64 * 1024;
const SWEEP_STEPS: u64 = 8;
const SPLIT_SAMPLES: usize = 1024;
const MAX_PARTS: u64 = 64;
pub(super) struct ShardEntry {
shard: Rc<MappedShard>,
seq: u64,
newest: u64,
pk_min: PkBuf,
pk_max: PkBuf,
}
impl ShardEntry {
pub(crate) fn open(dir: &str, seq: u64, schema: &SchemaDescriptor, newest: u64) -> Result<Self, StorageError> {
let shard = Rc::new(MappedShard::open(&manifest::shard_path(dir, seq), schema)?);
let pk_min = PkBuf::from_bytes(shard.get_pk_bytes(0));
let pk_max = PkBuf::from_bytes(shard.get_pk_bytes(shard.row_count() - 1));
Ok(ShardEntry { shard, seq, newest, pk_min, pk_max })
}
fn first_row_above(&self, key: &PkBuf) -> usize {
let row = self.shard.find_lower_bound_bytes(key.pk_bytes());
match row < self.shard.row_count() && pk_bytes_eq(self.shard.get_pk_bytes(row), key.pk_bytes()) {
true => pk_group_end(&*self.shard, row),
false => row,
}
}
fn probe_pk_bytes(&self, key: &[u8], filter_key: u64) -> Option<usize> {
if !pk_in_range(self.pk_min.pk_bytes(), self.pk_max.pk_bytes(), key) {
return None;
}
if !self.shard.shard_filter_may_contain(filter_key) {
return None;
}
let idx = self.shard.find_lower_bound_bytes(key);
(idx < self.shard.row_count() && pk_bytes_eq(self.shard.get_pk_bytes(idx), key)).then_some(idx)
}
}
struct LevelGuard {
guard_key: PkBuf,
entries: Vec<ShardEntry>,
}
impl LevelGuard {
fn dehydrated(&self) -> bool {
self.entries.iter().any(|e| e.shard.is_skeleton())
}
fn newest(&self) -> u64 {
self.entries
.iter()
.map(|e| e.newest)
.max()
.expect("a guard holds at least one shard")
}
fn bytes(&self) -> u64 {
self.entries.iter().map(|e| e.shard.file_len()).sum()
}
fn rows(&self) -> usize {
self.entries.iter().map(|e| e.shard.row_count()).sum()
}
fn retractions(&self) -> usize {
self.entries.iter().skip(1).map(|e| e.shard.retraction_rows()).sum()
}
fn key_extent(&self) -> (PkBuf, PkBuf) {
let lo = self.entries.iter().map(|e| e.pk_min).min();
let hi = self.entries.iter().map(|e| e.pk_max).max();
lo.zip(hi).expect("a guard holds at least one shard")
}
}
fn fold_destinations<'a>(key: PkBuf, entries: impl Iterator<Item = &'a ShardEntry> + Clone, target: u64) -> Vec<PkBuf> {
let mut keys = vec![key];
let bytes: u64 = entries.clone().map(|e| e.shard.file_len()).sum();
let parts = bytes.div_ceil(target).min(MAX_PARTS) as usize;
let lo = entries.clone().map(|e| e.pk_min).min();
let hi = entries.clone().map(|e| e.pk_max).max();
if parts >= 2 && lo < hi {
let rows: usize = entries.clone().map(|e| e.shard.row_count()).sum();
let step = (rows / SPLIT_SAMPLES).max(1);
let mut sample: Vec<PkBuf> = entries
.flat_map(|e| {
(0..e.shard.row_count())
.step_by(step)
.map(|r| PkBuf::from_bytes(e.shard.get_pk_bytes(r)))
})
.collect();
sample.sort_unstable();
sample.dedup();
let parts = parts.min(sample.len());
keys.extend((1..parts).map(|i| sample[i * sample.len() / parts]));
keys.sort_unstable();
keys.dedup();
}
keys
}
#[derive(Default)]
struct FLSMLevel {
guards: Vec<LevelGuard>,
}
impl FLSMLevel {
fn slot(&self, key: &[u8]) -> usize {
if self.guards.len() <= 1 {
return 0;
}
guard_slot(&self.guards, key, |g| g.guard_key.pk_bytes())
}
fn find_guards_for_range(&self, range_min: &[u8], range_max: &[u8]) -> std::ops::Range<usize> {
self.slot(range_min)..(self.slot(range_max) + 1).min(self.guards.len())
}
fn entries(&self) -> impl Iterator<Item = &ShardEntry> {
self.guards.iter().flat_map(|g| g.entries.iter())
}
fn bytes(&self) -> u64 {
self.guards.iter().map(LevelGuard::bytes).sum()
}
fn get_or_create_guard(&mut self, gk: PkBuf) -> &mut LevelGuard {
let pos = match self.guards.binary_search_by(|g| g.guard_key.cmp(&gk)) {
Ok(pos) => pos,
Err(pos) => {
self.guards
.insert(pos, LevelGuard { guard_key: gk, entries: Vec::new() });
pos
}
};
&mut self.guards[pos]
}
}
#[derive(Clone, Copy)]
pub(crate) enum ShardBudget {
Unbounded,
Dehydrate(u64),
Drop(u64),
}
impl ShardBudget {
fn cap(self) -> Option<u64> {
match self {
ShardBudget::Unbounded => None,
ShardBudget::Dehydrate(cap) | ShardBudget::Drop(cap) => Some(cap),
}
}
}
pub(super) struct ShardIndex {
pub(super) output_dir: String,
pub schema: SchemaDescriptor,
levels: [FLSMLevel; LEVELS],
pending: Vec<ShardEntry>,
shard_seq: u64,
published_through: u64,
retired: Vec<u64>,
l0_run_bytes: u64,
budget: ShardBudget,
dropped_max: PkBuf,
skip_pk_filter: bool,
cancel_yield: CancelYield,
owed: bool,
running: Option<index::Running>,
bands: Vec<PkBuf>,
}
#[derive(Clone, Copy)]
struct CancelYield {
cancelled: usize,
retractions: usize,
}
impl CancelYield {
const FRESH: Self = CancelYield { cancelled: 2, retractions: 1 };
fn expect(self, retractions: usize) -> usize {
(retractions as u128 * self.cancelled as u128 / self.retractions as u128) as usize
}
fn observe(&mut self, retractions: usize, cancelled: usize) {
if retractions > 0 {
self.cancelled = self.cancelled / 2 + cancelled.min(2 * retractions);
self.retractions = self.retractions / 2 + retractions;
}
}
}
impl ShardIndex {
pub(super) fn open(
output_dir: &str,
schema: SchemaDescriptor,
budget: ShardBudget,
skip_pk_filter: bool,
shards: &ShardSet,
) -> Result<Self, StorageError> {
let mut idx = ShardIndex {
output_dir: output_dir.to_string(),
schema,
levels: Default::default(),
pending: Vec::new(),
shard_seq: 0,
published_through: 0,
retired: Vec::new(),
l0_run_bytes: shards.run_bytes.max(MIN_GUARD_BYTES),
budget,
dropped_max: PkBuf::zeroed(schema.pk_stride()),
skip_pk_filter,
cancel_yield: CancelYield::FRESH,
owed: !shards.entries.is_empty(),
running: None,
bands: Vec::new(),
};
for e in &shards.entries {
let entry = ShardEntry::open(output_dir, e.seq, &idx.schema, e.newest)?;
let level = idx
.levels
.get_mut(e.level as usize)
.ok_or(StorageError::Corrupt("manifest level"))?;
level.get_or_create_guard(e.guard_key).entries.push(entry);
}
idx.shard_seq = idx.all_entries().map(|e| e.seq).max().unwrap_or(0);
idx.published_through = idx.shard_seq;
let live: HashSet<u64> = idx.all_entries().map(|e| e.seq).collect();
manifest::remove_stale_files(output_dir, |seq| live.contains(&seq))?;
Ok(idx)
}
pub(super) fn shard_set(&self) -> ShardSet {
let whole = PkBuf::zeroed(self.schema.pk_stride());
let entries = (0u64..)
.zip(&self.levels)
.flat_map(|(level, l)| {
l.guards.iter().flat_map(move |g| {
g.entries.iter().map(move |e| ManifestEntry {
seq: e.seq,
newest: e.newest,
level,
guard_key: g.guard_key,
})
})
})
.chain(self.pending.iter().map(|e| ManifestEntry {
seq: e.seq,
newest: e.newest,
level: L0 as u64,
guard_key: whole,
}))
.collect();
ShardSet { run_bytes: self.l0_run_bytes, entries }
}
pub(super) fn dropped_max(&self) -> PkBuf {
self.dropped_max
}
pub(super) fn swap_schema(&mut self, schema: SchemaDescriptor) -> Result<(), StorageError> {
if schema != self.schema {
self.abandon_fold();
let rebound: Vec<Rc<MappedShard>> = self
.all_entries()
.map(|e| e.shard.rebind(&schema).map(Rc::new))
.collect::<Result<_, _>>()?;
for (e, shard) in self.all_entries_mut().zip(rebound) {
e.shard = shard;
}
}
self.schema = schema;
Ok(())
}
pub(crate) fn verify_shards(&self) -> Result<(), StorageError> {
self.all_entries().try_for_each(|e| e.shard.verify_body())
}
}