use core::ops::{Bound, RangeBounds};
use alloc::{vec, vec::Vec};
use crate::comparator::UserComparator;
use crate::table::SeqnoVisibility;
use crate::table::columnar::{
COL_SEQNO, COL_USER_KEY, COL_VALUE_TYPE, ColumnBatch, TypeTag, bytes_column_row, fixed_u64_row,
};
use crate::table::columnar_predicate::{ColumnRangePredicate, filter_batch, take_rows};
use crate::{Error, SeqNo, Table, Tree, UserKey};
struct Segment {
table: Table,
min: UserKey,
max: UserKey,
global: SeqNo,
visibility: SeqnoVisibility,
may_dup: bool,
recency_rank: usize,
}
struct Group {
segments: Vec<Segment>,
max: UserKey,
}
impl Tree {
pub fn columnar_scan<R: RangeBounds<UserKey>>(
&self,
projection: &[u16],
predicate: Option<&ColumnRangePredicate>,
seqno: SeqNo,
range: R,
) -> crate::Result<ColumnarScan> {
if self.config.merge_operator.is_some() {
return Err(Error::FeatureUnsupported(
"columnar scan of a tree with a merge operator: merge chains \
would be returned unresolved",
));
}
let comparator = self.config.comparator.clone();
let lo = clone_bound(range.start_bound());
let hi = clone_bound(range.end_bound());
let bounds_ref = (bound_as_ref(&lo), bound_as_ref(&hi));
let super_version = self.version_history.read().get_version_for_snapshot(seqno);
let mut segments: Vec<Segment> = Vec::new();
for (recency_rank, table) in super_version.version.iter_tables().enumerate() {
if !table.check_key_range_overlap_cmp(&bounds_ref, comparator.as_ref()) {
continue;
}
let visibility = table.seqno_visibility(seqno);
if visibility == SeqnoVisibility::None {
continue;
}
if !table.metadata.columnar {
return Err(Error::FeatureUnsupported(
"columnar_scan: a non-columnar segment overlaps the range (mixed-mode tree)",
));
}
let key_range = &table.metadata.key_range;
let may_dup = table
.metadata
.key_count
.is_none_or(|k| k != table.metadata.item_count);
segments.push(Segment {
min: key_range.min().clone(),
max: key_range.max().clone(),
global: table.global_seqno(),
visibility,
may_dup,
recency_rank,
table: table.clone(),
});
}
let groups = group_by_overlap(segments, comparator.as_ref());
Ok(ColumnarScan {
groups: groups.into_iter().collect(),
buffered: Vec::new().into(),
projection: projection.to_vec(),
predicate: predicate.cloned(),
comparator,
seqno,
lo,
hi,
})
}
}
fn group_by_overlap(mut segments: Vec<Segment>, cmp: &dyn UserComparator) -> Vec<Group> {
use core::cmp::Ordering;
segments.sort_by(|a, b| cmp.compare(a.min.as_ref(), b.min.as_ref()));
let mut groups: Vec<Group> = Vec::new();
for seg in segments {
match groups.last_mut() {
Some(g) if cmp.compare(seg.min.as_ref(), g.max.as_ref()) != Ordering::Greater => {
if cmp.compare(seg.max.as_ref(), g.max.as_ref()) == Ordering::Greater {
g.max = seg.max.clone();
}
g.segments.push(seg);
}
_ => groups.push(Group {
max: seg.max.clone(),
segments: vec![seg],
}),
}
}
groups
}
pub struct ColumnarScan {
groups: alloc::collections::VecDeque<Group>,
buffered: alloc::collections::VecDeque<ColumnBatch>,
projection: Vec<u16>,
predicate: Option<ColumnRangePredicate>,
comparator: alloc::sync::Arc<dyn UserComparator>,
seqno: SeqNo,
lo: Bound<UserKey>,
hi: Bound<UserKey>,
}
impl ColumnarScan {
fn process_group(&self, group: &Group) -> crate::Result<Vec<ColumnBatch>> {
let rts = self.visible_group_range_tombstones(&group.segments)?;
if let [seg] = group.segments.as_slice() {
return self.process_singleton(seg, &rts);
}
self.merge_group(group, &rts)
}
fn visible_group_range_tombstones(
&self,
segments: &[Segment],
) -> crate::Result<Vec<(UserKey, UserKey, SeqNo)>> {
let mut rts = Vec::new();
for seg in segments {
for rt in seg.table.visible_range_tombstones() {
let eff = rt
.seqno
.checked_add(seg.global)
.ok_or(Error::InvalidHeader(
"columnar_scan: effective range-tombstone seqno overflows",
))?;
if eff < self.seqno {
rts.push((rt.start.clone(), rt.end.clone(), eff));
}
}
}
Ok(rts)
}
fn rt_covered(&self, rts: &[(UserKey, UserKey, SeqNo)], key: &[u8], eff: SeqNo) -> bool {
let cmp = self.comparator.as_ref();
rts.iter().any(|(start, end, rt_eff)| {
eff < *rt_eff
&& cmp.compare(key, start.as_ref()) != core::cmp::Ordering::Less
&& cmp.compare(key, end.as_ref()) == core::cmp::Ordering::Less
})
}
fn range_is_full(&self) -> bool {
matches!(self.lo, Bound::Unbounded) && matches!(self.hi, Bound::Unbounded)
}
fn globalize_seqnos(batch: &mut ColumnBatch, global: SeqNo) -> crate::Result<()> {
if global == 0 {
return Ok(());
}
let Some(col) = batch.columns.iter_mut().find(|c| c.column_id == COL_SEQNO) else {
return Ok(());
};
for row in 0..batch.row_count as usize {
let at = row * 8;
let bytes = col
.data
.get_mut(at..at + 8)
.ok_or(Error::InvalidHeader("columnar_scan: short seqno column"))?;
let local = u64::from_le_bytes(
(&*bytes)
.try_into()
.map_err(|_| Error::InvalidHeader("columnar_scan: short seqno column"))?,
);
let effective = local.checked_add(global).ok_or(Error::InvalidHeader(
"columnar_scan: effective seqno overflows",
))?;
bytes.copy_from_slice(&effective.to_le_bytes());
}
Ok(())
}
fn process_singleton(
&self,
seg: &Segment,
rts: &[(UserKey, UserKey, SeqNo)],
) -> crate::Result<Vec<ColumnBatch>> {
if seg.may_dup
|| seg.table.tombstone_count() > 0
|| seg.table.weak_tombstone_count() > 0
|| !rts.is_empty()
{
return self.process_singleton_dedup(seg, rts);
}
let range_filter = !self.range_is_full();
if seg.visibility == SeqnoVisibility::All && !range_filter {
let mut out = seg
.table
.columnar_scan(&self.projection, self.predicate.as_ref())?;
out.retain(|b| b.row_count > 0);
for batch in &mut out {
Self::globalize_seqnos(batch, seg.global)?;
}
return Ok(out);
}
let partial = seg.visibility == SeqnoVisibility::Partial;
let seqno_projected = self.projection.contains(&COL_SEQNO);
let key_projected = self.projection.contains(&COL_USER_KEY);
let mut augmented = self.projection.clone();
if partial && !seqno_projected {
augmented.push(COL_SEQNO);
}
if range_filter && !key_projected {
augmented.push(COL_USER_KEY);
}
let threshold = self.seqno.saturating_sub(seg.global);
let cmp = self.comparator.as_ref();
let mut out = Vec::new();
for batch in seg
.table
.columnar_scan(&augmented, self.predicate.as_ref())?
{
if batch.row_count == 0 {
continue;
}
let seqno_col = if partial {
Some(
batch
.columns
.iter()
.find(|c| c.column_id == COL_SEQNO)
.ok_or(Error::InvalidHeader(
"columnar_scan: partial-visibility batch missing the seqno column",
))?,
)
} else {
None
};
let key_col = if range_filter {
Some(
batch
.columns
.iter()
.find(|c| c.column_id == COL_USER_KEY)
.ok_or(Error::InvalidHeader(
"columnar_scan: range-filtered batch missing the key column",
))?,
)
} else {
None
};
let mut mask = Vec::with_capacity(batch.row_count as usize);
for row in 0..batch.row_count {
let seqno_ok = match seqno_col {
Some(seqno_col) => fixed_u64_row(&seqno_col.data, row)? < threshold,
None => true,
};
let keep = if !seqno_ok {
false
} else if let Some(key_col) = key_col {
let key = bytes_column_row(&key_col.data, batch.row_count, row)?;
key_in_bounds(key, &self.lo, &self.hi, cmp)
} else {
true
};
mask.push(keep);
}
let mut visible = filter_batch(&batch, &mask);
if partial && !seqno_projected {
visible.columns.retain(|c| c.column_id != COL_SEQNO);
}
if range_filter && !key_projected {
visible.columns.retain(|c| c.column_id != COL_USER_KEY);
}
if visible.row_count > 0 {
Self::globalize_seqnos(&mut visible, seg.global)?;
out.push(visible);
}
}
Ok(out)
}
fn process_singleton_dedup(
&self,
seg: &Segment,
rts: &[(UserKey, UserKey, SeqNo)],
) -> crate::Result<Vec<ColumnBatch>> {
let key_projected = self.projection.contains(&COL_USER_KEY);
let seqno_projected = self.projection.contains(&COL_SEQNO);
let partial = seg.visibility == SeqnoVisibility::Partial;
let seqno_needed = partial || !rts.is_empty();
let mut augmented = self.projection.clone();
if !key_projected {
augmented.push(COL_USER_KEY);
}
if seqno_needed && !seqno_projected {
augmented.push(COL_SEQNO);
}
let predicate_col = self.predicate.as_ref().map(|p| p.column_id);
let predicate_col_projected = predicate_col.is_some_and(|c| self.projection.contains(&c));
if let Some(pc) = predicate_col
&& !augmented.contains(&pc)
{
augmented.push(pc);
}
let deletes = seg.table.tombstone_count() > 0 || seg.table.weak_tombstone_count() > 0;
let vt_projected = self.projection.contains(&COL_VALUE_TYPE);
if deletes && !vt_projected {
augmented.push(COL_VALUE_TYPE);
}
let threshold = self.seqno.saturating_sub(seg.global);
let range_filter = !self.range_is_full();
let cmp = self.comparator.as_ref();
let mut out = Vec::new();
let mut last_key: Option<alloc::vec::Vec<u8>> = None;
for batch in seg.table.columnar_scan(&augmented, None)? {
if batch.row_count == 0 {
continue;
}
let key_col = batch
.columns
.iter()
.find(|c| c.column_id == COL_USER_KEY)
.ok_or(Error::InvalidHeader(
"columnar_scan: dedup batch missing the key column",
))?;
let vt_col = if deletes {
Some(
batch
.columns
.iter()
.find(|c| c.column_id == COL_VALUE_TYPE)
.ok_or(Error::InvalidHeader(
"columnar_scan: dedup batch missing the value-type column",
))?,
)
} else {
None
};
let seqno_col = if seqno_needed {
Some(
batch
.columns
.iter()
.find(|c| c.column_id == COL_SEQNO)
.ok_or(Error::InvalidHeader(
"columnar_scan: dedup batch missing the seqno column",
))?,
)
} else {
None
};
let mut mask = Vec::with_capacity(batch.row_count as usize);
for row in 0..batch.row_count {
let local = match seqno_col {
Some(seqno_col) => Some(fixed_u64_row(&seqno_col.data, row)?),
None => None,
};
let visible = !partial || local.is_some_and(|l| l < threshold);
if !visible {
mask.push(false);
continue;
}
let key = bytes_column_row(&key_col.data, batch.row_count, row)?;
if last_key
.as_deref()
.is_some_and(|p| cmp.compare(p, key) == core::cmp::Ordering::Equal)
{
mask.push(false);
continue;
}
last_key = Some(key.to_vec());
if !rts.is_empty() {
let eff =
local
.unwrap_or(0)
.checked_add(seg.global)
.ok_or(Error::InvalidHeader(
"columnar_scan: effective seqno overflows",
))?;
if self.rt_covered(rts, key, eff) {
mask.push(false);
continue;
}
}
if let Some(vt_col) = vt_col {
let byte = *vt_col.data.get(row as usize).ok_or(Error::InvalidHeader(
"columnar_scan: value-type column shorter than the row count",
))?;
let value_type = crate::ValueType::try_from(byte)
.map_err(|()| Error::InvalidTag(("ValueType", byte)))?;
if value_type.is_tombstone() {
mask.push(false);
continue;
}
}
mask.push(!range_filter || key_in_bounds(key, &self.lo, &self.hi, cmp));
}
let mut visible = filter_batch(&batch, &mask);
if let Some(pred) = self.predicate.as_ref() {
let pred_mask = pred.matching_rows(&visible);
visible = filter_batch(&visible, &pred_mask);
}
if !key_projected {
visible.columns.retain(|c| c.column_id != COL_USER_KEY);
}
if !seqno_projected {
visible.columns.retain(|c| c.column_id != COL_SEQNO);
}
if deletes && !vt_projected {
visible.columns.retain(|c| c.column_id != COL_VALUE_TYPE);
}
if let Some(pc) = predicate_col
&& !predicate_col_projected
{
visible.columns.retain(|c| c.column_id != pc);
}
if visible.row_count > 0 {
Self::globalize_seqnos(&mut visible, seg.global)?;
out.push(visible);
}
}
Ok(out)
}
fn merge_group(
&self,
group: &Group,
rts: &[(UserKey, UserKey, SeqNo)],
) -> crate::Result<Vec<ColumnBatch>> {
let key_projected = self.projection.contains(&COL_USER_KEY);
let seqno_projected = self.projection.contains(&COL_SEQNO);
let mut augmented = self.projection.clone();
if !key_projected {
augmented.push(COL_USER_KEY);
}
if !seqno_projected {
augmented.push(COL_SEQNO);
}
let predicate_col = self.predicate.as_ref().map(|p| p.column_id);
let predicate_col_projected = predicate_col.is_some_and(|c| self.projection.contains(&c));
if let Some(pc) = predicate_col
&& !augmented.contains(&pc)
{
augmented.push(pc);
}
let deletes = group
.segments
.iter()
.any(|s| s.table.tombstone_count() > 0 || s.table.weak_tombstone_count() > 0);
let vt_projected = self.projection.contains(&COL_VALUE_TYPE);
if deletes && !vt_projected {
augmented.push(COL_VALUE_TYPE);
}
let mut combined: Option<ColumnBatch> = None;
let mut effective: Vec<SeqNo> = Vec::new();
let mut source_rank: Vec<usize> = Vec::new();
for seg in &group.segments {
let threshold = self.seqno.saturating_sub(seg.global);
for batch in seg.table.columnar_scan(&augmented, None)? {
if batch.row_count == 0 {
continue;
}
let seqno_col = batch
.columns
.iter()
.find(|c| c.column_id == COL_SEQNO)
.ok_or(Error::InvalidHeader(
"columnar_scan: merged group missing the seqno column",
))?;
let mut mask = Vec::with_capacity(batch.row_count as usize);
for row in 0..batch.row_count {
let local = fixed_u64_row(&seqno_col.data, row)?;
let visible = seg.visibility == SeqnoVisibility::All || local < threshold;
mask.push(visible);
if visible {
let eff = local.checked_add(seg.global).ok_or(Error::InvalidHeader(
"columnar_scan: effective seqno overflow",
))?;
effective.push(eff);
source_rank.push(seg.recency_rank);
}
}
let visible = filter_batch(&batch, &mask);
if visible.row_count == 0 {
continue;
}
match &mut combined {
Some(acc) => acc.append(&visible)?,
None => combined = Some(visible),
}
}
}
let Some(combined) = combined else {
return Ok(Vec::new());
};
let key_col = combined
.columns
.iter()
.find(|c| c.column_id == COL_USER_KEY)
.ok_or(Error::InvalidHeader(
"columnar_scan: merged group missing the key column",
))?;
if key_col.type_tag != TypeTag::Bytes {
return Err(Error::InvalidHeader(
"columnar_scan: key column is not a bytes column",
));
}
let rows = combined.row_count;
debug_assert_eq!(rows as usize, effective.len(), "seqno tracked per row");
debug_assert_eq!(rows as usize, source_rank.len(), "rank tracked per row");
let mut keys: Vec<&[u8]> = Vec::with_capacity(rows as usize);
for i in 0..rows {
keys.push(bytes_column_row(&key_col.data, rows, i)?);
}
let key_at = |i: u32| keys.get(i as usize).copied().unwrap_or(&[]);
let eff_at = |i: u32| effective.get(i as usize).copied().unwrap_or(0);
let rank_at = |i: u32| source_rank.get(i as usize).copied().unwrap_or(usize::MAX);
let cmp = self.comparator.as_ref();
let mut order: Vec<u32> = (0..rows).collect();
order.sort_by(|&a, &b| {
cmp.compare(key_at(a), key_at(b))
.then_with(|| eff_at(b).cmp(&eff_at(a)))
.then_with(|| rank_at(a).cmp(&rank_at(b)))
});
let range_filter = !self.range_is_full();
let vt_col = if deletes {
Some(
combined
.columns
.iter()
.find(|c| c.column_id == COL_VALUE_TYPE)
.ok_or(Error::InvalidHeader(
"columnar_scan: merged group missing the value-type column",
))?,
)
} else {
None
};
let mut kept: Vec<u32> = Vec::with_capacity(order.len());
let mut prev: Option<&[u8]> = None;
for &i in &order {
let key = key_at(i);
if let Some(p) = prev
&& cmp.compare(p, key) == core::cmp::Ordering::Equal
{
continue;
}
prev = Some(key);
if range_filter && !key_in_bounds(key, &self.lo, &self.hi, cmp) {
continue;
}
if let Some(vt_col) = vt_col {
let byte = *vt_col.data.get(i as usize).ok_or(Error::InvalidHeader(
"columnar_scan: value-type column shorter than the row count",
))?;
let value_type = crate::ValueType::try_from(byte)
.map_err(|()| Error::InvalidTag(("ValueType", byte)))?;
if value_type.is_tombstone() {
continue;
}
}
if self.rt_covered(rts, key, eff_at(i)) {
continue;
}
kept.push(i);
}
let mut merged = take_rows(&combined, &kept)?;
if let Some(col) = merged.columns.iter_mut().find(|c| c.column_id == COL_SEQNO) {
for (row, &i) in kept.iter().enumerate() {
let at = row * 8;
let bytes = col
.data
.get_mut(at..at + 8)
.ok_or(Error::InvalidHeader("columnar_scan: short seqno column"))?;
bytes.copy_from_slice(&eff_at(i).to_le_bytes());
}
}
if let Some(pred) = self.predicate.as_ref() {
let mask = pred.matching_rows(&merged);
merged = filter_batch(&merged, &mask);
}
if !key_projected {
merged.columns.retain(|c| c.column_id != COL_USER_KEY);
}
if !seqno_projected {
merged.columns.retain(|c| c.column_id != COL_SEQNO);
}
if deletes && !vt_projected {
merged.columns.retain(|c| c.column_id != COL_VALUE_TYPE);
}
if let Some(pc) = predicate_col
&& !predicate_col_projected
{
merged.columns.retain(|c| c.column_id != pc);
}
if merged.row_count == 0 {
return Ok(Vec::new());
}
Ok(vec![merged])
}
}
fn key_in_bounds(
key: &[u8],
lo: &Bound<UserKey>,
hi: &Bound<UserKey>,
cmp: &dyn UserComparator,
) -> bool {
use core::cmp::Ordering;
let above_lo = match lo {
Bound::Unbounded => true,
Bound::Included(k) => cmp.compare(key, k.as_ref()) != Ordering::Less,
Bound::Excluded(k) => cmp.compare(key, k.as_ref()) == Ordering::Greater,
};
let below_hi = match hi {
Bound::Unbounded => true,
Bound::Included(k) => cmp.compare(key, k.as_ref()) != Ordering::Greater,
Bound::Excluded(k) => cmp.compare(key, k.as_ref()) == Ordering::Less,
};
above_lo && below_hi
}
impl Iterator for ColumnarScan {
type Item = crate::Result<ColumnBatch>;
fn next(&mut self) -> Option<Self::Item> {
loop {
if let Some(batch) = self.buffered.pop_front() {
return Some(Ok(batch));
}
let group = self.groups.pop_front()?;
match self.process_group(&group) {
Ok(batches) => self.buffered.extend(batches),
Err(e) => return Some(Err(e)),
}
}
}
}
fn clone_bound(bound: Bound<&UserKey>) -> Bound<UserKey> {
match bound {
Bound::Included(k) => Bound::Included(k.clone()),
Bound::Excluded(k) => Bound::Excluded(k.clone()),
Bound::Unbounded => Bound::Unbounded,
}
}
fn bound_as_ref(bound: &Bound<UserKey>) -> Bound<&[u8]> {
match bound {
Bound::Included(k) => Bound::Included(k.as_ref()),
Bound::Excluded(k) => Bound::Excluded(k.as_ref()),
Bound::Unbounded => Bound::Unbounded,
}
}