use gnitz_expr::SchemaFacts;
use gnitz_wire::{PkKeys, ReadBound, ReadSpec, SinkKind, WireFault, WireStatus};
use std::rc::Rc;
use super::SkeletonHydrator;
use crate::relation::RelationRegistry;
use crate::relation::{delta_round, delta_round_prefix};
use gnitz_expr::RowFilter;
use gnitz_zset::algebra::SinkPlan;
use gnitz_zset::repr::{Batch, SkeletonKeys, SourceCursor};
use gnitz_zset::schema::key::{key_range_between_cuts, KeyCut};
use gnitz_zset::schema::SchemaDescriptor;
impl RelationRegistry {
pub fn scan_spec(
&self,
target_id: u64,
spec: ReadSpec,
reply_layout: u64,
hydrator: Option<&mut dyn SkeletonHydrator>,
) -> Result<Rc<Batch>, String> {
let ReadSpec { bound, predicate, sink } = spec;
let entry = self.relation_or_err(target_id)?;
let src_schema = entry.schema();
if let (ReadBound::None, true, None, SinkKind::Rows { cut: None }) =
(&bound, predicate.is_empty(), &sink.map, &sink.kind)
{
if !entry.table().has_skeleton_rows() {
check_layout(reply_layout, &src_schema)?;
return Ok(entry.full_scan());
}
}
let (source, unapplied) = self.open_bound(target_id, bound)?;
let filter = RowFilter::for_read(&predicate, &unapplied, &src_schema).map_err(|e| format!("scan_spec: {e}"))?;
let mut sink = SinkPlan::from_wire(&src_schema, &sink, self.config.adhoc_group_cap)?;
check_layout(reply_layout, sink.output_schema())?;
let mut rows = Survivors {
registry: self,
id: target_id,
source,
hydrator,
filter,
ranges: Vec::new(),
};
let chunk_rows = self.config.scan_chunk_rows;
for drain_rows in drain_ramp(sink.first_drain(chunk_rows), chunk_rows) {
let Some((chunk, ranges)) = rows.next(drain_rows)? else {
break;
};
if sink.push(&chunk, ranges)? {
break;
}
}
Ok(Rc::new(sink.finish()))
}
pub fn delta_read(
&self,
id: u64,
after_tick: u64,
cut_tick: u64,
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());
}
let view = entry.schema();
check_layout(reply_layout, &view)?;
if after_tick == 0 {
return Ok(entry.full_scan());
}
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(),
);
let rows = feed.range_cursor(band).materialize();
Ok(Rc::new(rows.without_key_prefix(&view)))
}
}
fn check_layout(reply_layout: u64, produced: &SchemaDescriptor) -> Result<(), String> {
if reply_layout != produced.layout_digest() {
return Err("reply schema does not match the output layout".to_string());
}
Ok(())
}
type SurvivorChunk<'r> = (Batch, &'r mut Vec<(usize, usize)>);
struct Survivors<'a, 'h> {
registry: &'a RelationRegistry,
id: u64,
source: SourceCursor,
hydrator: Option<&'h mut dyn SkeletonHydrator>,
filter: RowFilter,
ranges: Vec<(usize, usize)>,
}
impl Survivors<'_, '_> {
fn next(&mut self, max_rows: usize) -> Result<Option<SurvivorChunk<'_>>, String> {
let mut skeletons = SkeletonKeys::default();
let Some(mut chunk) = self.source.drain_live_chunk(max_rows, &mut skeletons) else {
return Ok(None);
};
if !skeletons.keys.is_empty() {
chunk = self.hydrate(chunk, skeletons)?;
}
self.filter.ranges(&chunk.as_mem_batch(), &mut self.ranges);
Ok(Some((chunk, &mut self.ranges)))
}
fn hydrate(&mut self, live: Batch, mut skeletons: SkeletonKeys) -> Result<Batch, String> {
let Some(hydrator) = self.hydrator.as_deref_mut() else {
return Err(format!(
"relation {} holds skeleton rows but this process maintains no circuit",
self.id
));
};
let keys = PkKeys::from_sorted(live.schema().pk_stride(), std::mem::take(&mut skeletons.keys));
#[cfg(debug_assertions)]
let asked = keys.clone();
let hydrated = hydrator
.hydrate_keys(self.registry, self.id, keys)
.map_err(|e| format!("hydrate: view {}: {e}", self.id))?;
#[cfg(debug_assertions)]
assert_hydration_matches(&hydrated, &asked, &skeletons.coarse);
let schema = *live.schema();
Ok(match live.is_empty() {
true => hydrated,
false => hydrated.merged_consolidated(&live, &schema),
})
}
}
#[cfg(debug_assertions)]
fn assert_hydration_matches(out: &Batch, keys: &PkKeys, coarse: &[i64]) {
use gnitz_zset::repr::pk_group_end;
use gnitz_zset::schema::key::compare_pk_bytes;
assert_eq!(keys.len(), coarse.len());
let mut expected = keys.iter().zip(coarse).peekable();
let mut i = 0;
while i < out.len() {
let pk = out.get_pk_bytes(i);
if let Some((key, _)) = expected.next_if(|(key, _)| compare_pk_bytes(key, pk).is_lt()) {
panic!("hydration produced no rows for skeleton key {key:?}");
}
let Some((_, &weight)) = expected.next().filter(|(key, _)| *key == pk) else {
panic!("hydration produced rows for a PK no skeleton row named");
};
let j = pk_group_end(out, i);
let sum = out.as_mem_batch().sum_weights(i, j);
assert_eq!(sum, weight, "hydration weight mismatch for key {pk:?}");
i = j;
}
assert!(
expected.next().is_none(),
"hydration produced no rows for a trailing skeleton key"
);
}
fn drain_ramp(first: usize, chunk_rows: usize) -> impl Iterator<Item = usize> {
let (mut step, mut drained) = (first, 0usize);
std::iter::from_fn(move || {
let rows = step.min(chunk_rows - drained % chunk_rows);
drained += rows;
step = step.saturating_mul(2);
Some(rows)
})
}
#[cfg(test)]
#[path = "tests/scan_spec.rs"]
mod tests;
#[cfg(test)]
#[path = "benches/scan_spec.rs"]
mod bench;