use gnitz_expr::RowFilter;
use gnitz_wire::{PkKeys, ReadBound, ReadSpec, SinkKind, WireFault, WireStatus};
use std::rc::Rc;
use super::scan_spec::check_layout;
use crate::relation::{delta_round, delta_round_prefix, Cut, RelationRegistry};
use gnitz_zset::algebra::SinkPlan;
use gnitz_zset::repr::{Batch, SourceCursor};
use gnitz_zset::schema::key::{key_range_between_cuts, KeyCut};
impl RelationRegistry {
pub fn delta_read(
&self,
id: u64,
after_tick: u64,
cut_tick: u64,
spec: ReadSpec,
reply_layout: u64,
) -> Result<Rc<Batch>, WireFault> {
let entry = self.relation_or_err(id)?;
if !entry.kind().has_delta_feed() {
return Err(format!(
"delta_read: relation {id} carries no delta feed; \
create the view WITH (delta = '<size>') to subscribe to it"
)
.into());
}
if !matches!(spec.sink.kind, SinkKind::Rows { cut: None }) {
return Err("delta_read: a delta read forwards rows, uncut and unfolded"
.to_string()
.into());
}
if after_tick == 0 {
return Ok(self.scan_spec(id, spec, reply_layout, None)?);
}
let view = entry.schema();
let whole = spec.is_whole();
let reply = match whole {
true => view,
false => *SinkPlan::from_wire(&view, &spec.sink, self.config.adhoc_group_cap)?.output_schema(),
};
check_layout(reply_layout, &reply)?;
let feed = entry
.delta()
.ok_or_else(|| format!("delta_read: this process holds no delta store for relation {id}"))?;
let dropped_through = delta_round(feed.dropped_max().pk_bytes());
if after_tick < dropped_through {
return Err(WireFault {
status: WireStatus::DeltaExpired,
text: format!(
"delta cursor {after_tick} of relation {id} is below the \
retained floor {dropped_through}; re-read at 0"
),
});
}
let band = key_range_between_cuts(
KeyCut::above(&delta_round_prefix(after_tick)),
KeyCut::above(&delta_round_prefix(cut_tick)),
feed.schema().pk_stride(),
);
if whole {
let rows = feed.range_cursor(band).materialize();
return Ok(Rc::new(rows.without_key_prefix(&view)));
}
let ReadSpec { bound, predicate, sink } = spec;
let probes = match &bound {
ReadBound::PkSet(keys) if keys.stride() != view.pk_stride() => {
return Err(format!(
"delta_read: PkSet key stride {} != pk_stride {} (relation {id})",
keys.stride(),
view.pk_stride()
)
.into());
}
ReadBound::PkSet(keys) => {
let rounds = cut_tick.saturating_sub(after_tick);
(rounds.saturating_mul(keys.len() as u64) <= DELTA_GATHER_MAX_PROBES)
.then(|| stamped_keys(after_tick + 1..=cut_tick, keys))
}
_ => None,
};
let (mut source, unapplied) = match probes {
Some(probes) => (
SourceCursor::PkSet(Box::new(feed.gather(probes, Cut::Now))),
ReadBound::None,
),
None => (SourceCursor::Full(Box::new(feed.range_cursor(band))), bound),
};
let stamped = *feed.schema();
let mut filter =
RowFilter::for_read(&predicate, &unapplied, &stamped).map_err(|e| format!("delta_read: {e}"))?;
let mut plan = SinkPlan::from_wire(&stamped, &sink, self.config.adhoc_group_cap)?;
let mut ranges = Vec::new();
while let Some(chunk) = source.drain_chunk(self.config.scan_chunk_rows) {
filter.ranges(&chunk.as_mem_batch(), &mut ranges);
if let ReadBound::PkSet(keys) = &unapplied {
keep_keyed(&chunk, keys, &mut ranges);
}
plan.push(&chunk, &mut ranges)?;
}
let rows = plan.finish().without_key_prefix(&reply);
Ok(Rc::new(match sink.map {
Some(_) => rows.into_consolidated(),
None => rows,
}))
}
}
const DELTA_GATHER_MAX_PROBES: u64 = 4096;
fn stamped_keys(rounds: std::ops::RangeInclusive<u64>, keys: &PkKeys) -> PkKeys {
let stride = keys.stride() + delta_round_prefix(0).len();
let mut bytes = Vec::with_capacity(keys.len() * stride * rounds.clone().count());
for round in rounds {
for key in keys.iter() {
bytes.extend_from_slice(&delta_round_prefix(round));
bytes.extend_from_slice(key);
}
}
PkKeys::from_sorted(stride, bytes)
}
fn keep_keyed(chunk: &Batch, keys: &PkKeys, ranges: &mut Vec<(usize, usize)>) {
let stamp = chunk.schema().pk_stride() - keys.stride();
let mut kept = Vec::with_capacity(ranges.len());
for &(start, end) in ranges.iter() {
let mut run = None;
for row in start..end {
match (keys.contains(&chunk.get_pk_bytes(row)[stamp..]), run) {
(true, None) => run = Some(row),
(false, Some(first)) => {
kept.push((first, row));
run = None;
}
_ => {}
}
}
if let Some(first) = run {
kept.push((first, end));
}
}
*ranges = kept;
}
#[cfg(test)]
#[path = "tests/delta_read.rs"]
mod tests;
#[cfg(test)]
#[path = "benches/delta_read.rs"]
mod bench;