use kevy_index::{ScalarClauses, ValueTest};
use kevy_store::Store;
use super::args::Query;
use super::ops_clauses::{distinct_field, facet_fields, filter_tests, sort_field};
use super::wire::{encode_hydration_row, encode_value, peek_hydration};
use super::{ST_BUILDING, ST_CLAUSE, ST_NOINDEX, ST_OK, ST_OVERBUDGET};
use crate::index_runtime;
use crate::state::Ctx;
pub(super) const CURSOR_CLAUSE_CONFLICT: &str =
"CURSOR cannot combine with SORT|DISTINCT|FACET|OFFSET";
pub(super) fn clause_chunk(msg: &str) -> Vec<u8> {
let mut chunk = vec![ST_CLAUSE];
chunk.extend_from_slice(msg.as_bytes());
chunk
}
pub(super) fn run_claused_count(ctx: &Ctx<'_>, store: &mut Store, q: &Query) -> Vec<u8> {
let res = index_runtime::with_ready_segment(ctx, store, &q.name, |spec, seg, win| {
let now = (kevy_store::now_unix_ms() / 1000) as i64;
let (min, max) = q.bounds_for(spec, now)?;
super::probe_window(ctx, &q.name, win, &min);
let filters: Vec<(usize, ValueTest)> = filter_tests(spec, &q.filters, now)?;
let cold = match win
.filter(|w| w.has_cold())
.map(|w| w.cold_claused_count(spec.ty, &min, &max, &filters))
.transpose()
{
Ok(n) => n.unwrap_or(0),
Err(_) => return Err(vec![super::ST_NOINDEX]),
};
Ok(seg.count_claused(&min, &max, &filters) + cold)
});
match res {
Ok(Err(chunk)) => chunk,
Ok(Ok(n)) => {
let mut chunk = vec![ST_OK];
chunk.extend_from_slice(&n.to_le_bytes());
chunk
}
Err(e) if e.as_wire().starts_with("INDEXBUILDING") => vec![ST_BUILDING],
Err(e) if e.as_wire().starts_with("INDEXOVERBUDGET") => vec![ST_OVERBUDGET],
Err(_) => vec![ST_NOINDEX],
}
}
pub(super) fn run_claused_query(ctx: &Ctx<'_>, store: &mut Store, q: &Query) -> Vec<u8> {
let res = index_runtime::with_ready_segment(ctx, store, &q.name, |spec, seg, win| {
let now = (kevy_store::now_unix_ms() / 1000) as i64;
let (min, max) = q.bounds_for(spec, now)?;
super::probe_window(ctx, &q.name, win, &min);
let filters: Vec<(usize, ValueTest)> = filter_tests(spec, &q.filters, now)?;
let sort = sort_field(spec, &q.sort)?;
let distinct = distinct_field(spec, &q.distinct)?;
let facets = facet_fields(spec, &q.facets)?;
let clauses = ScalarClauses {
filters: &filters,
sort,
distinct,
facets: &facets,
fetch: q.limit + q.offset,
};
let cursor = q.cursor(spec.ty);
let mut page = seg.query_claused(&min, &max, cursor.as_ref(), &clauses);
if let Some(w) = win.filter(|w| w.has_cold()) {
match w.cold_claused(spec.ty, &min, &max, cursor.as_ref(), &clauses) {
Ok((chits, cfacets)) => merge_cold_claused(&mut page, chits, cfacets, &clauses),
Err(_) => return Err(vec![super::ST_NOINDEX]),
}
}
Ok(page)
});
match res {
Ok(Err(chunk)) => chunk,
Ok(Ok(page)) => encode_claused_chunk(store, q, &page),
Err(e) if e.as_wire().starts_with("INDEXBUILDING") => vec![ST_BUILDING],
Err(e) if e.as_wire().starts_with("INDEXOVERBUDGET") => vec![ST_OVERBUDGET],
Err(_) => vec![ST_NOINDEX],
}
}
fn merge_cold_claused(
page: &mut kevy_index::ClausedPage,
chits: Vec<kevy_index::ScalarHit>,
cfacets: Vec<Vec<kevy_index::FacetBucket>>,
c: &ScalarClauses<'_>,
) {
let all: Vec<(kevy_index::ScalarHit, ())> =
page.hits.drain(..).chain(chits).map(|h| (h, ())).collect();
let merged =
kevy_index::merge_claused(all, c.sort.map(|(_, d, _)| d), c.distinct.is_some(), 0, c.fetch);
page.hits = merged.into_iter().map(|(h, ())| h).collect();
kevy_index::fold_facets(&mut page.facets, cfacets);
kevy_index::sort_facets(&mut page.facets);
page.cursor = match c.selects() || page.hits.len() < c.fetch {
true => None,
false => page
.hits
.last()
.map(|h| kevy_index::Cursor { value: h.value.clone(), key: h.key.clone() }),
};
}
fn encode_claused_chunk(store: &mut Store, q: &Query, page: &kevy_index::ClausedPage) -> Vec<u8> {
let mut chunk = vec![ST_OK];
chunk.extend_from_slice(&(page.hits.len() as u32).to_le_bytes());
let keys: Vec<&[u8]> = page.hits.iter().map(|h| h.key.as_slice()).collect();
let rows = peek_hydration(store, &keys, &q.fields);
for (i, h) in page.hits.iter().enumerate() {
chunk.extend_from_slice(&(h.key.len() as u32).to_le_bytes());
chunk.extend_from_slice(&h.key);
encode_value(&mut chunk, &h.value);
encode_hydration_row(&mut chunk, q.fields.len(), &rows[i]);
if q.sort.is_some() {
encode_okey(&mut chunk, h.okey.as_deref());
}
if q.distinct.is_some() {
encode_okey(&mut chunk, h.dkey.as_deref());
}
}
for field in &page.facets {
chunk.extend_from_slice(&(field.len() as u32).to_le_bytes());
for (id, label, n) in field {
chunk.extend_from_slice(&(id.len() as u32).to_le_bytes());
chunk.extend_from_slice(id);
chunk.extend_from_slice(&(label.len() as u32).to_le_bytes());
chunk.extend_from_slice(label);
chunk.extend_from_slice(&n.to_le_bytes());
}
}
chunk
}
fn encode_okey(chunk: &mut Vec<u8>, k: Option<&[u8]>) {
match k {
None => chunk.push(0),
Some(b) => {
chunk.push(1);
chunk.extend_from_slice(&(b.len() as u32).to_le_bytes());
chunk.extend_from_slice(b);
}
}
}