use super::segment::Segment;
use crate::storage::lsm::columnar::{ColumnTypeTag, ColumnarSSTable, FixedSegment, TextSegment};
use crate::types::{ColumnType, Value};
use std::cmp::Reverse;
use std::collections::BinaryHeap;
use std::sync::Arc;
enum ColData {
Fixed(FixedSegment),
Text(TextSegment),
Vector(Vec<Option<Vec<f32>>>),
Spatial(Vec<Option<crate::types::Geometry>>),
Opaque,
}
struct SegmentCursor {
row_map_keys: Vec<u64>,
row_map_ts: Vec<u64>,
row_map_deleted: Vec<bool>,
order: Vec<usize>,
pos: usize,
col_data: Vec<ColData>,
col_types: Vec<ColumnType>,
}
impl SegmentCursor {
fn new(seg: &Segment, col_types: Vec<ColumnType>) -> Self {
let n = seg.sst.num_rows;
let _ = seg.sst.load_all_timestamps();
let _ = seg.sst.load_full_keys();
let mut row_map_keys = Vec::with_capacity(n);
let mut row_map_ts = Vec::with_capacity(n);
let mut row_map_deleted = Vec::with_capacity(n);
for i in 0..n {
row_map_keys.push(seg.sst.row_map.key(i));
row_map_ts.push(seg.sst.row_map.timestamp_loaded(i));
row_map_deleted.push(seg.sst.row_map.is_deleted(i));
}
let mut order: Vec<usize> = (0..n).collect();
order.sort_by_key(|&i| row_map_keys[i]);
let mut col_data = Vec::with_capacity(col_types.len());
for ci in 0..col_types.len() {
let cd = if ci < seg.sst.column_tags.len() && seg.sst.column_tags[ci].is_fixed() {
match seg.sst.read_fixed_i64(ci) {
Ok(seg_data) => ColData::Fixed(seg_data),
Err(_) => ColData::Opaque,
}
} else if ci < seg.sst.column_tags.len()
&& matches!(seg.sst.column_tags[ci], ColumnTypeTag::Text)
{
match seg.sst.read_text(ci) {
Ok(seg_data) => ColData::Text(seg_data),
Err(_) => ColData::Opaque,
}
} else if ci < seg.sst.column_tags.len()
&& matches!(seg.sst.column_tags[ci], ColumnTypeTag::Vector)
{
let decoded = seg.sst.read_vectors(ci).unwrap_or_default();
let mut per = vec![None; n];
let mut di = 0usize;
for i in 0..n {
if seg.sst.row_map.is_deleted(i) {
continue;
}
let ek = seg.sst.row_map.key(i) & 0xFFFFFFFF;
while di < decoded.len() && decoded[di].0 != ek {
di += 1;
}
if di < decoded.len() {
per[i] = Some(decoded[di].1.clone());
di += 1;
}
}
ColData::Vector(per)
} else if ci < seg.sst.column_tags.len()
&& matches!(seg.sst.column_tags[ci], ColumnTypeTag::Spatial)
{
let decoded = seg.sst.read_spatial(ci).unwrap_or_default();
let mut per = vec![None; n];
let mut di = 0usize;
for i in 0..n {
if seg.sst.row_map.is_deleted(i) {
continue;
}
let ek = seg.sst.row_map.key(i) & 0xFFFFFFFF;
while di < decoded.len() && decoded[di].0 != ek {
di += 1;
}
if di < decoded.len() {
per[i] = Some(decoded[di].1.clone());
di += 1;
}
}
ColData::Spatial(per)
} else {
ColData::Opaque
};
col_data.push(cd);
}
Self {
row_map_keys,
row_map_ts,
row_map_deleted,
order,
pos: 0,
col_data,
col_types,
}
}
#[inline]
fn peek_key(&self) -> Option<u64> {
self.order.get(self.pos).map(|&i| self.row_map_keys[i])
}
fn advance(&mut self) -> Option<(u64, u64, bool, Vec<Value>)> {
let &i = self.order.get(self.pos)?;
self.pos += 1;
let key = self.row_map_keys[i];
let ts = self.row_map_ts[i];
let deleted = self.row_map_deleted[i];
let row = if deleted {
Vec::new()
} else {
self.decode_row(i)
};
Some((key, ts, deleted, row))
}
fn decode_row(&self, i: usize) -> Vec<Value> {
let mut row = Vec::with_capacity(self.col_types.len());
for (ci, ct) in self.col_types.iter().enumerate() {
let v = match self.col_data.get(ci) {
Some(ColData::Fixed(f)) => match ct {
ColumnType::Integer => f.get_i64(i).map(Value::Integer),
ColumnType::Float => f.get_f64(i).map(Value::Float),
ColumnType::Boolean => f.get_bool(i).map(Value::Bool),
ColumnType::Timestamp => f
.get_i64(i)
.map(|v| Value::Timestamp(crate::types::Timestamp::from_micros(v))),
_ => None,
},
Some(ColData::Text(t)) => t.get_str(i).map(|s| Value::Text(s.into())),
Some(ColData::Vector(cols)) => cols
.get(i)
.cloned()
.flatten()
.map(|v| Value::Vector(crate::types::ArcVec(Arc::new(v)))),
Some(ColData::Spatial(cols)) => cols
.get(i)
.cloned()
.flatten()
.map(|g| Value::Spatial(std::boxed::Box::new(g))),
_ => None,
}
.unwrap_or(Value::Null);
row.push(v);
}
row
}
}
pub struct MergeCursor {
cursors: Vec<SegmentCursor>,
heap: BinaryHeap<Reverse<(u64, usize)>>,
}
impl MergeCursor {
pub fn new(segments: &[Arc<Segment>], col_types: &[ColumnType]) -> Self {
let cursors: Vec<SegmentCursor> = segments
.iter()
.map(|s| SegmentCursor::new(s, col_types.to_vec()))
.collect();
let mut heap = BinaryHeap::with_capacity(cursors.len());
for (idx, c) in cursors.iter().enumerate() {
if let Some(k) = c.peek_key() {
heap.push(Reverse((k, idx)));
}
}
Self { cursors, heap }
}
}
impl Iterator for MergeCursor {
type Item = (u64 , u64 , Vec<Value>);
fn next(&mut self) -> Option<Self::Item> {
loop {
let min_key = self.heap.peek().map(|Reverse((k, _))| *k)?;
let mut at_key: Vec<usize> = Vec::new();
while let Some(&Reverse((k, _))) = self.heap.peek() {
if k != min_key {
break;
}
let Reverse((_, idx)) = self.heap.pop().unwrap();
at_key.push(idx);
}
let mut best_idx: Option<usize> = None;
let mut best_ts: u64 = 0;
for &idx in &at_key {
let c = &self.cursors[idx];
if let Some(&row_i) = c.order.get(c.pos) {
let ts = c.row_map_ts[row_i];
if best_idx.is_none() || ts >= best_ts {
best_idx = Some(idx);
best_ts = ts;
}
}
}
let mut emitted: Option<(u64, u64, Vec<Value>)> = None;
for &idx in &at_key {
let c = &mut self.cursors[idx];
if let Some((key, ts, deleted, row)) = c.advance() {
if Some(idx) == best_idx && !deleted {
emitted = Some((key, ts, row));
}
}
if let Some(k) = c.peek_key() {
self.heap.push(Reverse((k, idx)));
}
}
if let Some(e) = emitted {
return Some(e);
}
}
}
}
#[allow(dead_code)]
fn _type_anchor(_s: &ColumnarSSTable) {}