use std::ops::Range;
use std::rc::Rc;
use crate::repr::batch::Batch;
use crate::repr::merge::ColumnarSource;
use crate::repr::merge::{ColPtr, UnifiedSource};
use crate::repr::scatter::DecodedColumns;
use crate::repr::shard_reader::MappedShard;
use crate::schema::payload_order::{with_payload_cmp, PayloadOrder};
use crate::schema::SchemaDescriptor;
use gnitz_wire::RowSource;
#[derive(Clone)]
pub enum Run {
Mem(Rc<Batch>),
Shard(Rc<MappedShard>),
}
impl Run {
#[inline(always)]
pub fn get_weight(&self, row: usize) -> i64 {
ColumnarSource::get_weight(self, row)
}
pub(crate) fn find_lower_bound_bytes(&self, key: &[u8]) -> usize {
match self {
Run::Mem(b) => b.find_lower_bound_bytes(key),
Run::Shard(s) => s.find_lower_bound_bytes(key),
}
}
pub(crate) fn advance_to(&self, key: &[u8], hint: usize) -> usize {
match self {
Run::Mem(b) => b.advance_to(key, hint),
Run::Shard(s) => s.advance_to(key, hint),
}
}
pub(crate) fn slice_to_owned_batch(&self, start: usize, row_count: usize, schema: &SchemaDescriptor) -> Batch {
match self {
Run::Mem(b) => {
let mut slice = Batch::from_ranges(b, &[(start, start + row_count)], 0);
slice.set_schema(schema);
slice
}
Run::Shard(s) => {
debug_assert!(*s.schema() == *schema, "a shard slices under the schema it is bound to");
s.slice_to_owned_batch(start, row_count)
}
}
}
}
pub struct StoredRow {
pub(crate) run: Run,
pub(crate) row: usize,
}
impl StoredRow {
#[inline]
pub fn new(run: Run, row: usize) -> Self {
StoredRow { run, row }
}
#[inline]
pub(crate) fn weight(&self) -> i64 {
self.run.get_weight(self.row)
}
#[inline]
pub fn source(&self) -> (&impl gnitz_wire::RowSource, usize) {
(&self.run, self.row)
}
}
pub fn first_live_payload_group(schema: &SchemaDescriptor, pool: &[StoredRow]) -> Option<usize> {
with_payload_cmp!(schema, first_live_payload_group_with, schema, pool)
}
fn first_live_payload_group_with<P: PayloadOrder>(
schema: &SchemaDescriptor,
pool: &[StoredRow],
payload: P,
) -> Option<usize> {
(0..pool.len()).find(|&i| {
let (run, row) = (&pool[i].run, pool[i].row);
let others: i64 = pool
.iter()
.enumerate()
.filter(|&(j, c)| j != i && payload.compare(schema, run, row, &c.run, c.row).is_eq())
.map(|(_, c)| c.weight())
.sum();
pool[i].weight() + others > 0
})
}
impl RowSource for Run {
#[inline(always)]
fn get_pk_bytes(&self, row: usize) -> &[u8] {
match self {
Run::Mem(b) => b.get_pk_bytes(row),
Run::Shard(s) => s.get_pk_bytes(row),
}
}
#[inline(always)]
fn get_null_word(&self, row: usize) -> u64 {
match self {
Run::Mem(b) => b.get_null_word(row),
Run::Shard(s) => s.get_null_word(row),
}
}
#[inline(always)]
fn get_col_ptr(&self, row: usize, payload_col: usize, col_size: usize) -> &[u8] {
match self {
Run::Mem(b) => b.get_col_ptr(row, payload_col, col_size),
Run::Shard(s) => s.get_col_ptr(row, payload_col, col_size),
}
}
#[inline(always)]
fn blob(&self) -> &[u8] {
match self {
Run::Mem(b) => b.blob(),
Run::Shard(s) => s.blob(),
}
}
#[inline(always)]
fn row_count(&self) -> usize {
match self {
Run::Mem(b) => b.count,
Run::Shard(s) => s.row_count(),
}
}
}
impl ColumnarSource for Run {
fn to_unified(
&self,
schema: &SchemaDescriptor,
cols: &mut Vec<ColPtr>,
window: Range<usize>,
decoded: &mut DecodedColumns,
) -> UnifiedSource<'_> {
match self {
Run::Mem(b) => crate::repr::merge::mem_batch_to_unified(&b.as_mem_batch(), schema, cols),
Run::Shard(s) => s.to_unified(schema, cols, window, decoded),
}
}
#[inline(always)]
fn get_weight(&self, row: usize) -> i64 {
match self {
Run::Mem(b) => b.get_weight(row),
Run::Shard(s) => s.get_weight(row),
}
}
#[inline(always)]
fn is_skeleton(&self) -> bool {
match self {
Run::Mem(_) => false,
Run::Shard(s) => s.is_skeleton(),
}
}
}