use std::ops::{Bound, Range, RangeBounds};
use std::sync::Arc;
use pigeonhole_engine::{
FamilyId, QualifierFilter, ReadSpec, ScanCursor, ScanSpec, TableInfo, ValuePredicate,
};
use smallvec::SmallVec;
use crate::cell::RowBuf;
use crate::table::{TableCore, family_not_found};
use crate::{Result, Row, RowRef, Snapshot};
#[derive(Debug, Clone, PartialEq)]
pub enum ValueFilter {
Equals(Vec<u8>),
Prefix(Vec<u8>),
I64(std::cmp::Ordering, i64),
}
impl ValueFilter {
pub(crate) fn to_engine(&self) -> ValuePredicate {
match self {
ValueFilter::Equals(v) => ValuePredicate::Equals(v.clone()),
ValueFilter::Prefix(v) => ValuePredicate::Prefix(v.clone()),
ValueFilter::I64(o, v) => ValuePredicate::I64(*o, *v),
}
}
}
#[derive(Debug, Clone, PartialEq)]
pub enum Condition {
Exists {
family: String,
qualifier: Vec<u8>,
},
Absent {
family: String,
qualifier: Vec<u8>,
},
Value {
family: String,
qualifier: Vec<u8>,
filter: ValueFilter,
},
}
type FamilyIds = SmallVec<[FamilyId; 4]>;
#[derive(Debug, Clone)]
struct Selection {
names: SmallVec<[u8; 32]>,
ends: SmallVec<[u32; 4]>,
spec: ReadSpec,
snapshot: Option<Snapshot>,
}
impl Selection {
fn new() -> Self {
let mut spec = ReadSpec::default();
spec.versions = 1;
Self {
names: SmallVec::new(),
ends: SmallVec::new(),
spec,
snapshot: None,
}
}
fn push_family(&mut self, name: &str) {
self.names.extend_from_slice(name.as_bytes());
let end = u32::try_from(self.names.len()).expect("family names under 4 GiB");
self.ends.push(end);
}
fn family_names(&self) -> impl Iterator<Item = &str> {
let mut start = 0;
self.ends.iter().map(move |&end| {
let name = &self.names[start..end as usize];
start = end as usize;
std::str::from_utf8(name).expect("pushed from a str")
})
}
fn qualifier_bounds(&mut self, start: Bound<&[u8]>, end: Bound<&[u8]>) {
self.spec.qualifiers =
QualifierFilter::Range(start.map(<[u8]>::to_vec), end.map(<[u8]>::to_vec));
}
fn time_range(&mut self, range: Range<u64>) {
self.spec.time_range = Some((range.start, range.end));
}
fn start_latest(&self, core: &TableCore) -> Result<Option<(Arc<TableInfo>, FamilyIds)>> {
core.db.check_open()?;
let info = core.current_info();
let mut ids = SmallVec::new();
for name in self.family_names() {
let Some(family) = info.family(name) else {
return Ok(None);
};
if !ids.contains(&family.id) {
ids.push(family.id);
}
}
if self.ends.is_empty() {
ids.extend(info.families.iter().map(|f| f.id));
}
if ids.is_empty() {
return Ok(None);
}
Ok(Some((info, ids)))
}
fn start(
&mut self,
core: &TableCore,
) -> Result<(Arc<TableInfo>, pigeonhole_engine::Snapshot, FamilyIds)> {
core.db.check_open()?;
let snapshot = match self.snapshot.take() {
Some(s) => {
core.db.check_snapshot(&s)?;
s.inner
}
None => core.db.engine.snapshot()?,
};
let info = core.current_info();
let mut ids = SmallVec::new();
for name in self.family_names() {
let id = info
.family(name)
.ok_or_else(|| family_not_found(&info.name, name))?
.id;
if !ids.contains(&id) {
ids.push(id);
}
}
Ok((info, snapshot, ids))
}
}
macro_rules! selection_methods {
() => {
pub fn families<'f>(mut self, families: impl IntoIterator<Item = &'f str>) -> Self {
self.sel.names.clear();
self.sel.ends.clear();
for name in families {
self.sel.push_family(name);
}
self
}
pub fn family(mut self, family: &str) -> Self {
self.sel.push_family(family);
self
}
pub fn qualifier_prefix(mut self, prefix: &[u8]) -> Self {
self.sel.spec.qualifiers = QualifierFilter::Prefix(prefix.to_vec());
self
}
pub fn qualifier_range<'k, K: AsRef<[u8]> + ?Sized + 'k>(
mut self,
range: impl RangeBounds<&'k K>,
) -> Self {
self.sel.qualifier_bounds(
range.start_bound().map(|k| k.as_ref()),
range.end_bound().map(|k| k.as_ref()),
);
self
}
pub fn qualifier_bounds(mut self, start: Bound<&[u8]>, end: Bound<&[u8]>) -> Self {
self.sel.qualifier_bounds(start, end);
self
}
pub fn latest(mut self) -> Self {
self.sel.spec.versions = 1;
self
}
pub fn time_range(mut self, range: Range<u64>) -> Self {
self.sel.time_range(range);
self
}
pub fn value_filter(mut self, filter: ValueFilter) -> Self {
self.sel.spec.value = Some(filter.to_engine());
self
}
pub fn snapshot(mut self, snapshot: &Snapshot) -> Self {
self.sel.snapshot = Some(snapshot.clone());
self
}
};
}
#[derive(Debug)]
#[must_use = "a row read does nothing until .read()"]
pub struct RowRead<'t> {
core: &'t TableCore,
row: Vec<u8>,
sel: Selection,
}
impl<'t> RowRead<'t> {
pub(crate) fn new(core: &'t TableCore, row: &[u8]) -> Self {
Self {
core,
row: row.to_vec(),
sel: Selection::new(),
}
}
selection_methods!();
pub fn versions(mut self, n: u32) -> Self {
self.sel.spec.versions = n;
self
}
pub fn column_limit(mut self, n: u32) -> Self {
self.sel.spec.columns_per_row = n;
self
}
pub fn read(mut self) -> Result<Option<RowRef<'t>>> {
let latest = match self.sel.snapshot {
None => self.sel.start_latest(self.core)?,
Some(_) => None,
};
let (info, snapshot, families) = match latest {
Some((info, families)) => (info, None, families),
None => {
let (info, snapshot, families) = self.sel.start(self.core)?;
(info, Some(snapshot), families)
}
};
let mut buf = RowBuf::new(info, 0);
let (engine, table) = (&self.core.db.engine, self.core.info.id);
let any = match &snapshot {
None => engine.read_row_latest_into(
table,
&self.row,
&families,
&self.sel.spec,
&mut buf,
)?,
Some(snapshot) => engine.read_row_into_families(
snapshot,
table,
&self.row,
&families,
&self.sel.spec,
&mut buf,
)?,
};
if !any {
return Ok(None);
}
buf.key = self.row;
Ok(Some(RowRef::owned(buf)))
}
}
#[derive(Debug)]
#[must_use = "a scan does nothing until .iter()"]
pub struct Scan<'t> {
core: &'t TableCore,
start: Bound<Vec<u8>>,
end: Bound<Vec<u8>>,
sel: Selection,
limit: Option<u64>,
}
impl<'t> Scan<'t> {
pub(crate) fn new(core: &'t TableCore, start: Bound<Vec<u8>>, end: Bound<Vec<u8>>) -> Self {
Self {
core,
start,
end,
sel: Selection::new(),
limit: None,
}
}
selection_methods!();
pub fn versions(mut self, n: u32) -> Self {
self.sel.spec.versions = n;
self
}
pub fn columns_per_row(mut self, n: u32) -> Self {
self.sel.spec.columns_per_row = n;
self
}
pub fn limit(mut self, n: u64) -> Self {
self.limit = Some(n);
self
}
pub fn iter(mut self) -> Result<RowIter<'t>> {
let (info, snapshot, families) = self.sel.start(self.core)?;
let mut spec = ScanSpec::new(self.start, self.end);
spec.read = self.sel.spec;
spec.read.families = families.into_vec();
spec.limit = self.limit.unwrap_or(0);
let cursor = self
.core
.db
.engine
.scan(&snapshot, self.core.info.id, spec)?;
Ok(RowIter {
cursor,
buf: RowBuf::new(info, 0),
done: self.limit == Some(0),
_table: std::marker::PhantomData,
})
}
}
#[derive(Debug)]
pub struct RowIter<'t> {
cursor: ScanCursor,
buf: RowBuf,
done: bool,
_table: std::marker::PhantomData<&'t TableCore>,
}
impl RowIter<'_> {
fn fill(&mut self) -> Result<bool> {
if self.done {
return Ok(false);
}
let r = self.fill_inner();
if !matches!(r, Ok(true)) {
self.done = true;
}
r
}
fn fill_inner(&mut self) -> Result<bool> {
self.buf.clear();
if !self.cursor.next_row()? {
return Ok(false);
}
self.buf.key.extend_from_slice(self.cursor.row());
loop {
let start = self.buf.qualifiers.len();
let Some(family) = self.cursor.next_cell_into(&mut self.buf.qualifiers)? else {
break;
};
let end = self.buf.qualifiers.len();
self.cursor.push_current(family, start..end, &mut self.buf);
}
Ok(true)
}
pub fn next_ref(&mut self) -> Result<Option<RowRef<'_>>> {
Ok(if self.fill()? {
Some(RowRef::borrowed(&self.buf))
} else {
None
})
}
}
impl Iterator for RowIter<'_> {
type Item = Result<Row>;
fn next(&mut self) -> Option<Self::Item> {
match self.fill() {
Ok(true) => {
let next = self.buf.empty_like();
Some(Ok(Row::new(std::mem::replace(&mut self.buf, next))))
}
Ok(false) => None,
Err(e) => Some(Err(e)),
}
}
}
#[cfg(test)]
mod tests {
use std::sync::Arc;
use pigeonhole_engine::{FamilyOptions, ValueRef, WriteBatch};
use super::Selection;
use crate::{Family, Options, Pigeonhole};
#[test]
fn a_family_added_before_the_view_loads_is_not_read_unnamed() {
let dir = crate::doc_support::temp_dir();
let db = Pigeonhole::open(dir.join("t.phdb"), Options::default().shards(1)).unwrap();
let t = db
.table("t")
.unwrap()
.family("f", Family::default())
.create_if_missing()
.unwrap();
t.mutate(b"r").put("f", b"q", b"v").commit().unwrap();
let engine = Arc::clone(&t.core.db.engine);
let table = t.core.info.id;
engine.before_latest_view_load(Box::new({
let engine = Arc::clone(&engine);
move || {
let info = engine
.add_family(table, "x", FamilyOptions::default())
.unwrap();
let x = info.family("x").unwrap().id;
let mut wb = WriteBatch::new();
wb.put(table, x, b"r", b"q", None, ValueRef::Bytes(b"new"))
.unwrap();
engine.commit(wb, None).unwrap();
}
}));
let row = t.row(b"r").read().unwrap().unwrap();
let families: Vec<&str> = row.iter().map(|e| e.family).collect();
assert_eq!(
families,
["f"],
"a cell of a family the read could not name"
);
let row = t.row(b"r").read().unwrap().unwrap();
let families: Vec<&str> = row.iter().map(|e| e.family).collect();
assert_eq!(families, ["f", "x"]);
db.close().unwrap();
}
#[test]
fn family_names_come_back_in_order_inline_or_spilled() {
let mut sel = Selection::new();
assert_eq!(sel.family_names().count(), 0);
sel.push_family("f");
sel.push_family("");
sel.push_family("été");
assert_eq!(sel.family_names().collect::<Vec<_>>(), ["f", "", "été"]);
assert!(!sel.names.spilled());
let long = "a-family-name-longer-than-the-inline-buffer";
for _ in 0..5 {
sel.push_family(long);
}
assert!(sel.names.spilled());
let names: Vec<&str> = sel.family_names().collect();
assert_eq!(names.len(), 8);
assert_eq!(&names[..3], ["f", "", "été"]);
assert!(names[3..].iter().all(|n| *n == long));
sel.names.clear();
sel.ends.clear();
sel.push_family("g");
assert_eq!(sel.family_names().collect::<Vec<_>>(), ["g"]);
}
}