use super::batch::Batch;
use super::error::StorageError;
use super::merge::run_merge_in;
use super::scatter::UnifiedSet;
use super::shard_reader::MappedShard;
use crate::schema::key::compare_pk_ordering;
use crate::schema::SchemaDescriptor;
use gnitz_wire::PkBuf;
use gnitz_wire::RowSource;
use std::ops::Range;
pub fn guard_slot<T>(guards: &[T], key: &[u8], gk: impl Fn(&T) -> &[u8]) -> usize {
guards
.partition_point(|g| compare_pk_ordering(gk(g), key).is_le())
.saturating_sub(1)
}
#[inline(never)]
pub fn merge_guard(
shards: &[&MappedShard],
guards: &[PkBuf],
g: usize,
starts: &mut [usize],
dehydrate: bool,
schema: &SchemaDescriptor,
) -> Option<(bool, Batch)> {
debug_assert!(guards.is_sorted_by(|a, b| a < b), "guards must be sorted and distinct");
let skeleton = dehydrate || shards.iter().any(|s| s.is_skeleton());
let out_schema = if skeleton { schema.pk_only() } else { *schema };
let windows: Vec<Range<usize>> = shards
.iter()
.zip(starts)
.map(|(shard, start)| {
let end = match guards.get(g + 1) {
Some(next) => shard.advance_to(next.pk_bytes(), *start),
None => shard.row_count(),
};
std::mem::replace(start, end)..end
})
.collect();
let mut survivors: Vec<(u32, u32, i64)> = Vec::with_capacity(windows.iter().map(Range::len).sum());
run_merge_in(shards, &out_schema, windows.iter().cloned(), |src, row, w| {
survivors.push((src as u32, row as u32, w));
});
if survivors.is_empty() {
return None;
}
let set = UnifiedSet::of(shards, &out_schema, windows);
Some((skeleton, set.materialize(&survivors, set.src_rows())))
}
pub fn merge_and_route(
shards: &[&MappedShard],
guards: &[PkBuf],
dehydrate: bool,
schema: &SchemaDescriptor,
emit: &mut dyn FnMut(PkBuf, bool, Batch) -> Result<(), StorageError>,
) -> Result<(), StorageError> {
assert!(!guards.is_empty(), "merge_and_route requires at least one guard");
for s in shards {
s.verify_body()?;
}
let mut starts = vec![0usize; shards.len()];
for (g, &key) in guards.iter().enumerate() {
if let Some((skeleton, batch)) = merge_guard(shards, guards, g, &mut starts, dehydrate, schema) {
emit(key, skeleton, batch)?;
}
}
Ok(())
}
#[cfg(test)]
#[path = "tests/compact.rs"]
mod tests;
#[cfg(test)]
#[path = "benches/compact.rs"]
mod bench;