use gnitz_wire::PkKeys;
use super::{PkSetGather, ReadCursor, SkeletonKeys};
use crate::repr::batch::Batch;
use crate::schema::key::{sort_indices, KeySpec};
pub enum SourceCursor {
Full(Box<ReadCursor>),
Bounded(Box<BoundedIndexCursor>),
PkSet(Box<PkSetGather>),
}
impl SourceCursor {
pub fn drain_chunk(&mut self, max_rows: usize) -> Option<Batch> {
let mut skeletons = SkeletonKeys::default();
let chunk = self.drain_live_chunk(max_rows, &mut skeletons);
skeletons.assert_none();
chunk
}
pub fn drain_live_chunk(&mut self, max_rows: usize, skeletons: &mut SkeletonKeys) -> Option<Batch> {
match self {
SourceCursor::Full(c) => c.drain_live_chunk(max_rows, skeletons),
SourceCursor::Bounded(c) => c.drain_live_chunk(max_rows, skeletons),
SourceCursor::PkSet(g) => g.drain_live_chunk(max_rows, skeletons),
}
}
}
pub struct BoundedIndexCursor {
idx: ReadCursor,
src: PkSetGather,
spec: KeySpec,
pks: Vec<u8>,
order: Vec<u32>,
}
impl BoundedIndexCursor {
pub fn new(idx: ReadCursor, src: ReadCursor, spec: KeySpec) -> Self {
let none = PkKeys::from_sorted(src.schema.pk_stride(), Vec::new());
BoundedIndexCursor {
idx,
src: PkSetGather::new(src, none),
spec,
pks: Vec::new(),
order: Vec::new(),
}
}
fn drain_live_chunk(&mut self, max_rows: usize, skeletons: &mut SkeletonKeys) -> Option<Batch> {
assert!(
max_rows > 0,
"BoundedIndexCursor::drain_live_chunk: max_rows must be positive"
);
loop {
if let Some(chunk) = self.src.drain_live_chunk(max_rows, skeletons) {
debug_assert!(
(1..chunk.len()).all(|i| chunk.get_pk_bytes(i - 1) != chunk.get_pk_bytes(i)),
"index owner holds two rows under one PK"
);
return Some(chunk);
}
if !self.refill(max_rows) {
return None;
}
}
}
fn refill(&mut self, max_rows: usize) -> bool {
let stride = self.src.schema().pk_stride();
self.pks.clear();
let mut taken = 0;
while taken < max_rows && self.idx.valid {
debug_assert!(self.idx.current_weight > 0, "an index entry at a non-positive weight");
self.pks
.extend_from_slice(self.spec.split_entry(self.idx.current_pk_bytes()).1);
self.idx.advance();
taken += 1;
}
if taken == 0 {
return false;
}
sort_indices(&self.pks, stride, &mut self.order);
let mut keys = Vec::with_capacity(self.pks.len());
for &i in &self.order {
keys.extend_from_slice(&self.pks[i as usize * stride..][..stride]);
}
self.src.reload(PkKeys::from_sorted(stride, keys));
true
}
}