use std::sync::atomic::{AtomicU8, Ordering};
use std::sync::{Arc, OnceLock};
use rudb_common::{Error, Result};
use rudb_vector::Vector;
#[derive(Debug, Default)]
pub(crate) struct Peel {
answers: OnceLock<Answers>,
}
#[derive(Debug)]
struct Answers {
dictionary: Arc<Vector>,
decided: Vec<AtomicU8>,
}
impl Peel {
pub(crate) fn answer<M, D>(
&self,
column: &Vector,
len: usize,
map: M,
decide: D,
) -> Option<Result<Vec<bool>>>
where
M: Fn(usize) -> usize,
D: Fn(&Vector, usize) -> Result<bool>,
{
let (codes, dictionary) = column.shared_dictionary_parts()?;
let answers = self.answers.get_or_init(|| Answers {
dictionary: Arc::clone(dictionary),
decided: (0..dictionary.len()).map(|_| AtomicU8::new(0)).collect(),
});
if !Arc::ptr_eq(&answers.dictionary, dictionary) {
return None;
}
Some(answers.run(dictionary, codes, len, map, decide))
}
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub(crate) enum Found {
At(u32),
Absent,
}
#[derive(Debug, Default)]
pub(crate) struct Lookup {
memo: OnceLock<Searched>,
}
#[derive(Debug)]
struct Searched {
dictionary: Arc<Vector>,
found: Found,
}
impl Lookup {
pub(crate) fn find(&self, column: &Vector, wanted: &[u8]) -> Option<Result<Found>> {
let (_, dictionary) = column.shared_dictionary_parts()?;
if let Some(memo) = self.memo.get() {
return Arc::ptr_eq(&memo.dictionary, dictionary).then_some(Ok(memo.found));
}
let ranks = dictionary.ranks()?;
let found = match search(dictionary, ranks, wanted) {
Ok(found) => found,
Err(error) => return Some(Err(error)),
};
let _ = self.memo.set(Searched { dictionary: Arc::clone(dictionary), found });
Some(Ok(found))
}
}
fn search(dictionary: &Vector, ranks: usize, wanted: &[u8]) -> Result<Found> {
let mut low = 0;
let mut high = ranks;
while low < high {
let middle = low + (high - low) / 2;
match dictionary.compare_rank(middle, wanted)? {
std::cmp::Ordering::Less => low = middle + 1,
std::cmp::Ordering::Greater => high = middle,
std::cmp::Ordering::Equal => {
return Ok(Found::At(dictionary.code_at_rank(middle)?));
}
}
}
Ok(Found::Absent)
}
impl Answers {
fn run<M, D>(
&self,
dictionary: &Vector,
codes: &[u32],
len: usize,
map: M,
decide: D,
) -> Result<Vec<bool>>
where
M: Fn(usize) -> usize,
D: Fn(&Vector, usize) -> Result<bool>,
{
let mut out = Vec::with_capacity(len);
for slot in 0..len {
let code = *codes
.get(map(slot))
.ok_or_else(|| Error::internal("a peeled row is past the end of its codes"))?;
let state = self.decided.get(code as usize).ok_or_else(|| {
Error::internal("a peeled code is past the end of its dictionary")
})?;
let mut held = state.load(Ordering::Relaxed);
if held == 0 {
held = u8::from(decide(dictionary, code as usize)?) + 1;
state.store(held, Ordering::Relaxed);
}
out.push(held == 2);
}
Ok(out)
}
}
#[cfg(test)]
mod tests {
use std::sync::atomic::AtomicUsize;
use rudb_common::{LogicalType, Value};
use super::*;
fn letters(values: &[&str]) -> Vector {
let values: Vec<Value> = values.iter().map(|text| Value::Varchar((*text).into())).collect();
Vector::from_values(LogicalType::Varchar, &values).expect("a vector of text")
}
fn holds(dictionary: &Vector, code: usize) -> Result<bool> {
Ok(dictionary.try_bytes_at(code)?.is_some_and(|bytes| bytes.starts_with(b"a")))
}
#[test]
fn a_peel_decides_once_per_distinct_code_and_reads_the_memo_after_that() {
let values = Arc::new(letters(&["apple", "pear", "avocado"]));
let column = Vector::stable_dictionary(vec![0, 1, 2, 1, 0, 0], values).expect("in range");
let peel = Peel::default();
let calls = AtomicUsize::new(0);
let answers = peel
.answer(
&column,
6,
|slot| slot,
|dictionary, code| {
calls.fetch_add(1, Ordering::Relaxed);
holds(dictionary, code)
},
)
.expect("a shared dictionary is peelable")
.expect("the predicate answers");
assert_eq!(answers, [true, false, true, false, true, true]);
assert_eq!(calls.load(Ordering::Relaxed), 3, "six rows over three distinct values");
}
#[test]
fn a_second_chunk_over_the_same_dictionary_asks_the_predicate_nothing() {
let values = Arc::new(letters(&["apple", "pear"]));
let first = Vector::stable_dictionary(vec![0, 1], Arc::clone(&values)).expect("in range");
let second = Vector::stable_dictionary(vec![1, 1, 0], values).expect("in range");
let peel = Peel::default();
let calls = AtomicUsize::new(0);
let count = |column: &Vector, len: usize| {
peel.answer(
column,
len,
|slot| slot,
|dictionary, code| {
calls.fetch_add(1, Ordering::Relaxed);
holds(dictionary, code)
},
)
.expect("peelable")
.expect("answers")
};
assert_eq!(count(&first, 2), [true, false]);
assert_eq!(calls.load(Ordering::Relaxed), 2);
assert_eq!(count(&second, 3), [false, false, true]);
assert_eq!(calls.load(Ordering::Relaxed), 2, "the second chunk decided nothing new");
}
#[test]
fn a_different_dictionary_is_declined_rather_than_answered_from_the_first_ones_memo() {
let peel = Peel::default();
let first = Vector::stable_dictionary(vec![0], Arc::new(letters(&["apple"]))).expect("one");
let second = Vector::stable_dictionary(vec![0], Arc::new(letters(&["pear"]))).expect("one");
assert!(peel.answer(&first, 1, |slot| slot, holds).is_some());
assert!(peel.answer(&second, 1, |slot| slot, holds).is_none());
}
#[test]
fn a_column_that_is_not_a_shared_dictionary_has_nothing_to_peel() {
let flat = letters(&["apple", "pear"]);
assert!(Peel::default().answer(&flat, 2, |slot| slot, holds).is_none());
}
#[derive(Debug)]
struct Filed {
values: Vec<Vec<u8>>,
order: Vec<(u64, u32)>,
reads: AtomicUsize,
}
fn head(bytes: &[u8]) -> u64 {
let mut word = [0; 8];
let take = bytes.len().min(8);
word[..take].copy_from_slice(&bytes[..take]);
u64::from_be_bytes(word)
}
impl Filed {
fn new(values: &[&str]) -> Self {
let values: Vec<Vec<u8>> = values.iter().map(|text| text.as_bytes().to_vec()).collect();
let mut order = (0..values.len() as u32)
.map(|code| (head(&values[code as usize]), code))
.collect::<Vec<_>>();
order.sort_by(|&(_, left), &(_, right)| {
values[left as usize].cmp(&values[right as usize])
});
Self { values, order, reads: AtomicUsize::new(0) }
}
fn at(&self, rank: usize) -> Result<(u64, u32)> {
self.order
.get(rank)
.copied()
.ok_or_else(|| Error::internal("a rank past the end of the order"))
}
}
impl rudb_vector::TextSource for Filed {
fn len(&self) -> usize {
self.values.len()
}
fn bytes_at(&self, index: usize) -> Result<Option<&[u8]>> {
self.reads.fetch_add(1, Ordering::Relaxed);
Ok(self.values.get(index).map(Vec::as_slice))
}
fn footprint(&self) -> usize {
self.values.iter().map(Vec::len).sum()
}
fn ranks(&self) -> Option<usize> {
Some(self.order.len())
}
fn compare_rank(&self, rank: usize, wanted: &[u8]) -> Result<std::cmp::Ordering> {
let (found, code) = self.at(rank)?;
let settled = found.cmp(&head(wanted));
if settled != std::cmp::Ordering::Equal {
return Ok(settled);
}
Ok(self.bytes_at(code as usize)?.unwrap_or_default().cmp(wanted))
}
fn code_at_rank(&self, rank: usize) -> Result<u32> {
Ok(self.at(rank)?.1)
}
}
fn filed(count: usize) -> (Vector, Arc<Filed>) {
let spellings =
(0..count).map(|at| format!("value-{:04}", (at * 7919) % count)).collect::<Vec<_>>();
let source =
Arc::new(Filed::new(&spellings.iter().map(String::as_str).collect::<Vec<_>>()));
let values = Vector::external_text(LogicalType::Varchar, Arc::clone(&source) as Arc<_>)
.expect("a filed vector");
(values, source)
}
#[test]
fn a_literal_is_found_in_a_sorted_dictionary_without_reading_every_value() {
let (values, source) = filed(1024);
let column = Vector::stable_dictionary(vec![3, 900, 3], Arc::new(values)).expect("codes");
let found = Lookup::default()
.find(&column, b"value-0700")
.expect("a sorted dictionary can be searched")
.expect("the search reads");
let Found::At(code) = found else { panic!("the dictionary holds it") };
assert_eq!(source.values[code as usize], b"value-0700");
let reads = source.reads.load(Ordering::Relaxed);
assert!(reads <= 11, "a search of 1024 values read {reads} of them");
}
#[test]
fn a_search_over_values_that_differ_early_reads_only_the_one_it_finds() {
let source = Arc::new(Filed::new(&["cherry", "apple", "date", "banana"]));
let values = Vector::external_text(LogicalType::Varchar, Arc::clone(&source) as Arc<_>)
.expect("a filed vector");
let column = Vector::stable_dictionary(vec![0, 1, 2, 3], Arc::new(values)).expect("codes");
assert_eq!(
Lookup::default().find(&column, b"cherry").expect("searchable").expect("read"),
Found::At(0)
);
assert_eq!(source.reads.load(Ordering::Relaxed), 1, "only the value it found");
assert_eq!(
Lookup::default().find(&column, b"fig").expect("searchable").expect("read"),
Found::Absent
);
assert_eq!(source.reads.load(Ordering::Relaxed), 1, "and nothing for the one it did not");
}
#[test]
fn values_that_share_their_first_eight_bytes_are_still_told_apart() {
let source = Arc::new(Filed::new(&["prefixed-two", "prefixed-one", "prefixed-three"]));
let values = Vector::external_text(LogicalType::Varchar, Arc::clone(&source) as Arc<_>)
.expect("a filed vector");
let column = Vector::stable_dictionary(vec![0, 1, 2], Arc::new(values)).expect("codes");
for (wanted, expected) in [
(&b"prefixed-one"[..], Found::At(1)),
(b"prefixed-two", Found::At(0)),
(b"prefixed-three", Found::At(2)),
(b"prefixed-four", Found::Absent),
] {
assert_eq!(
Lookup::default().find(&column, wanted).expect("searchable").expect("read"),
expected,
"searching for {}",
String::from_utf8_lossy(wanted)
);
}
}
#[test]
fn a_literal_the_dictionary_does_not_hold_is_answered_absent() {
let (values, _) = filed(64);
let column = Vector::stable_dictionary(vec![0], Arc::new(values)).expect("codes");
let found = Lookup::default().find(&column, b"nothing").expect("searchable").expect("read");
assert_eq!(found, Found::Absent);
}
#[test]
fn the_empty_string_is_found_like_any_other_value() {
let source = Arc::new(Filed::new(&["pear", "", "apple"]));
let values = Vector::external_text(LogicalType::Varchar, source).expect("a filed vector");
let column = Vector::stable_dictionary(vec![0, 1, 2], Arc::new(values)).expect("codes");
assert_eq!(
Lookup::default().find(&column, b"").expect("searchable").expect("read"),
Found::At(1)
);
}
#[test]
fn a_second_chunk_over_the_same_dictionary_searches_nothing() {
let (values, source) = filed(256);
let values = Arc::new(values);
let first = Vector::stable_dictionary(vec![0, 1], Arc::clone(&values)).expect("codes");
let second = Vector::stable_dictionary(vec![2], Arc::clone(&values)).expect("codes");
let lookup = Lookup::default();
let found = lookup.find(&first, b"value-0100").expect("searchable").expect("read");
let reads = source.reads.load(Ordering::Relaxed);
assert!(reads > 0, "the first chunk did the search");
assert_eq!(lookup.find(&second, b"value-0100").expect("searchable").expect("read"), found);
assert_eq!(source.reads.load(Ordering::Relaxed), reads, "the second chunk read nothing");
}
#[test]
fn a_search_is_not_reused_across_two_dictionaries() {
let lookup = Lookup::default();
let (first, _) = filed(8);
let (second, _) = filed(8);
let first = Vector::stable_dictionary(vec![0], Arc::new(first)).expect("codes");
let second = Vector::stable_dictionary(vec![0], Arc::new(second)).expect("codes");
assert!(lookup.find(&first, b"value-0000").is_some());
assert!(lookup.find(&second, b"value-0000").is_none());
}
#[test]
fn a_dictionary_with_no_order_is_declined() {
let values = Arc::new(letters(&["apple", "pear"]));
let column = Vector::stable_dictionary(vec![0, 1], values).expect("codes");
assert!(Lookup::default().find(&column, b"apple").is_none());
}
}