gnitz-store 0.1.3

The Z-set store of the gnitz database: LSM storage, relation registry and read executor
//! `RunSet` — a set of in-heap sorted runs with a PK bloom and a fold trigger.

use std::cell::OnceCell;
use std::rc::Rc;

use super::bloom::BloomFilter;
use gnitz_wire::PkBuf;
use gnitz_zset::repr::{merge_consolidated, Batch, MemBatch};
use gnitz_zset::schema::key::{pk_bytes_eq, pk_in_range, pk_ranges_overlap, probe_key};
use gnitz_zset::schema::SchemaDescriptor;

/// Runs to accumulate before folding them into one: bounds how many runs a
/// cursor merges and a PK probe walks.
const FOLD_THRESHOLD: usize = 16;

/// The share of its budget a set must have free not to be
/// [crowded](RunSet::is_crowded): one part in this many.
const ROOM_SHARE: usize = 16;

fn pk_min(run: &Batch) -> &[u8] {
    run.get_pk_bytes(0)
}

fn pk_max(run: &Batch) -> &[u8] {
    run.get_pk_bytes(run.len() - 1)
}

pub(super) struct RunSet {
    /// Each [`Batch::trimmed`], since the set charges only its rows' bytes.
    runs: Vec<Rc<Batch>>,
    /// PK bloom over every key pushed since it was built — a superset of the
    /// live runs' keys — built on the first probe.
    bloom: OnceCell<BloomFilter>,
    /// Heap budget: [`is_full`](Self::is_full) reports crossing it, and the
    /// bloom's key capacity is derived from it. What crossing it *means* — fold
    /// into the next tier, or spill to a shard — is `Table`'s policy.
    budget: usize,
    bytes: usize,
}

impl RunSet {
    pub(super) fn new(budget: usize) -> Self {
        RunSet {
            runs: Vec::with_capacity(FOLD_THRESHOLD),
            bloom: OnceCell::new(),
            budget,
            bytes: 0,
        }
    }

    /// Append a consolidated run, folding the set when it gets crowded. Empty
    /// runs are never stored.
    pub(super) fn push(&mut self, run: Batch, schema: &SchemaDescriptor) {
        debug_assert!(run.is_consolidated(), "RunSet::push requires a consolidated run",);
        self.push_run(Rc::new(run.trimmed()), schema);
    }

    /// [`Self::push`] for a run another set held.
    fn push_run(&mut self, run: Rc<Batch>, schema: &SchemaDescriptor) {
        if run.is_empty() {
            return;
        }
        if let Some(bloom) = self.bloom.get_mut() {
            bloom_add_batch(bloom, &run);
        }
        self.bytes += run.total_bytes();
        self.runs.push(run);
        if self.runs.len() >= FOLD_THRESHOLD {
            self.fold(schema);
        }
    }

    /// The runs whose PK extent meets the inclusive `bound`; `None` takes them all.
    pub(super) fn runs_overlapping(&self, bound: Option<(PkBuf, PkBuf)>) -> impl Iterator<Item = &Rc<Batch>> {
        self.runs.iter().filter(move |run| {
            bound.is_none_or(|(lo, hi)| pk_ranges_overlap(pk_min(run), pk_max(run), lo.pk_bytes(), hi.pk_bytes()))
        })
    }

    /// Visit each run holding `key`, newest first, with the row its matches start
    /// at; `fingerprint` is `key`'s [`probe_key`].
    pub(super) fn find_pk_bytes(&self, key: &[u8], fingerprint: u64, mut visitor: impl FnMut(&Rc<Batch>, usize)) {
        if !self.may_contain(fingerprint) {
            return;
        }
        for run in self.runs.iter().rev() {
            if !pk_in_range(pk_min(run), pk_max(run), key) {
                continue;
            }
            let start = run.find_lower_bound_bytes(key);
            if start < run.len() && pk_bytes_eq(run.get_pk_bytes(start), key) {
                visitor(run, start);
            }
        }
    }

    /// How many runs this set holds — one cursor source each.
    pub(super) fn len(&self) -> usize {
        self.runs.len()
    }

    /// The set has outgrown its heap budget and must be drained by its owner.
    pub(super) fn is_full(&self) -> bool {
        self.bytes > self.budget
    }

    /// The set has less than a [`ROOM_SHARE`]th of its budget free.
    pub(super) fn is_crowded(&self) -> bool {
        self.bytes > self.budget - self.budget / ROOM_SHARE
    }

    pub(super) fn row_count(&self) -> usize {
        self.runs.iter().map(|r| r.len()).sum()
    }

    pub(super) fn clear(&mut self) {
        self.runs.clear();
        self.bytes = 0;
        self.bloom.take();
    }

    /// Widen every run narrower than `schema` to it, filling the appended
    /// trailing columns with NULL (`ALTER TABLE … ADD COLUMN`).
    pub(super) fn widen_runs(&mut self, schema: &SchemaDescriptor) {
        let npc = schema.num_payload_cols();
        let mut bytes = 0;
        for run in &mut self.runs {
            if run.num_payload_cols() < npc {
                let widened = run.widened_with_nulls(schema, false);
                debug_assert!(
                    widened.is_consolidated(),
                    "widen_runs: the widened run must still be consolidated",
                );
                *run = Rc::new(widened.trimmed());
            }
            bytes += run.total_bytes();
        }
        self.bytes = bytes;
        // The widen leaves PK bytes untouched, so the bloom stays valid.
    }

    /// Fold every run into one consolidated run, dropping net-zero
    /// (PK, payload) rows.
    pub(super) fn fold(&mut self, schema: &SchemaDescriptor) {
        if self.runs.len() <= 1 {
            return;
        }
        let input_rows = self.row_count();
        // A dominant run — usually the previous fold's output — is galloped
        // against the fold of the rest, so its stretches are bulk-copied.
        let big = (0..self.runs.len()).max_by_key(|&i| self.runs[i].len()).unwrap();
        let merged = if self.runs[big].len() * 2 >= input_rows {
            // `remove`, not `swap_remove`: the rest stay in push order, which
            // is what lets an ascending load's runs fold by concatenation.
            let dominant = self.runs.remove(big);
            match &self.runs[..] {
                [one] => dominant.merged_consolidated(one, schema),
                _ => dominant.merged_consolidated(&self.consolidate_all(schema), schema),
            }
        } else {
            self.consolidate_all(schema)
        };
        self.runs.clear();
        // A cancelled row's key stays in the filter as a false positive.
        if self.bloom.get().is_some_and(|b| b.stale(merged.len())) {
            self.bloom.take();
        }
        self.bytes = 0;
        if !merged.is_empty() {
            let run = Rc::new(merged.trimmed());
            self.bytes = run.total_bytes();
            self.runs.push(run);
        }
    }

    /// Every run folded N-way into one consolidated batch.
    fn consolidate_all(&self, schema: &SchemaDescriptor) -> Batch {
        let views: Vec<MemBatch> = self.runs.iter().map(|r| r.as_mem_batch()).collect();
        merge_consolidated(&views, schema)
    }

    /// Fold to a single run and hand it to `write`, emptying the set once that
    /// succeeds. `Ok(false)` when the set held no row to write.
    pub(super) fn spill<E>(
        &mut self,
        schema: &SchemaDescriptor,
        write: impl FnOnce(&Batch) -> Result<(), E>,
    ) -> Result<bool, E> {
        self.fold(schema);
        let Some(run) = self.runs.first() else {
            return Ok(false);
        };
        write(run)?;
        self.clear();
        Ok(true)
    }

    /// Fold to a single run, move it into `dst` and answer it; this set is left
    /// empty. `None` when the set is empty or fully cancelled.
    pub(super) fn drain_into(&mut self, dst: &mut RunSet, schema: &SchemaDescriptor) -> Option<Rc<Batch>> {
        self.fold(schema);
        let run = self.runs.pop();
        self.clear();
        if let Some(run) = &run {
            dst.push_run(Rc::clone(run), schema);
        }
        run
    }

    /// Bloom probe for a PK by its [`probe_key`]. The first probe builds the
    /// filter from all live runs.
    fn may_contain(&self, probe_key: u64) -> bool {
        // Answered without building a budget-sized filter.
        if self.runs.is_empty() {
            return false;
        }
        let bloom = self.bloom.get_or_init(|| {
            // Sized to the rows the budget holds, not the current ones: later
            // pushes add to it.
            let row_width = self.runs[0].schema().row_width();
            let mut bloom = BloomFilter::new((self.budget / row_width).max(16));
            for run in &self.runs {
                bloom_add_batch(&mut bloom, run);
            }
            bloom
        });
        bloom.may_contain(probe_key)
    }
}

/// Insert every row's PK into `bloom`, keyed by [`probe_key`].
fn bloom_add_batch(bloom: &mut BloomFilter, batch: &Batch) {
    for i in 0..batch.len() {
        bloom.add(probe_key(batch.get_pk_bytes(i)));
    }
}

#[cfg(test)]
#[path = "tests/run_set.rs"]
mod tests;

#[cfg(test)]
#[path = "benches/run_set.rs"]
mod bench;