use kevy_index::{IndexValue, ValType};
use kevy_store::Store;
use super::args::{ComposeQuery, HybridArgs, KnnArgs, MatchArgs, Shape};
use super::wire::{
decode_gstats_arg, encode_agg_chunk, encode_hydration_row, encode_stats_chunk, peek_hydration,
};
use super::{ST_BADARGS, ST_BUILDING, ST_NOINDEX, ST_OK, ST_OVERBUDGET};
use crate::index_runtime;
use crate::state::Ctx;
#[path = "ops_encode.rs"]
mod encode;
use encode::encode_hits;
use super::ops_clauses::{boxed_preds, distinct_field, facet_fields, scope_positions, sort_field};
pub(super) fn op_match(ctx: &Ctx<'_>, store: &mut Store, argv: &[Vec<u8>]) -> Vec<u8> {
let q = match MatchArgs::parse_terminal(argv) {
crate::cmd_index_query::args::MatchParse::Ok(q) => q,
crate::cmd_index_query::args::MatchParse::BadArgs => return vec![ST_BADARGS],
crate::cmd_index_query::args::MatchParse::NotYet(clause) => {
let mut chunk = vec![crate::cmd_index_query::ST_NOTYET];
chunk.extend_from_slice(clause);
return chunk;
}
};
let res = index_runtime::with_ready_text_segment(ctx, store, &q.name, |ts, spec| {
let want = scope_positions(spec, &q.scope)?;
let opts = kevy_text::QueryOpts {
stats: None,
typo: q.typo,
fields: &want,
filter: &[],
sort: None,
distinct: None,
};
Ok((ts.docs(), ts.total_len_in(&want), ts.query_df_in(&q.text, opts)))
});
match res {
Ok(Ok((n_docs, total_len, tokdf))) => {
let mut chunk = Vec::new();
encode_stats_chunk(&mut chunk, n_docs, total_len, &tokdf);
chunk
}
Ok(Err(chunk)) => 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 op_match_score(ctx: &Ctx<'_>, store: &mut Store, argv: &[Vec<u8>]) -> Vec<u8> {
let Some(q) = super::args::parse_match_score(argv) else {
return vec![ST_BADARGS];
};
let Some(stats) = argv.get(4).and_then(|b| decode_gstats_arg(b)) else {
return vec![ST_BADARGS];
};
let res = index_runtime::with_ready_text_segment(ctx, store, &q.name, |ts, spec| {
let (hits, sort_field, distinct_field, facets) = scored_hits(ts, spec, &q, &stats)?;
let spans = q.highlight.as_ref().map(|want| {
hits.iter().map(|h| hit_highlight(ts, spec, &h.key, &q.text, want)).collect::<Vec<_>>()
});
let hit_keys = |f: usize| {
hits.iter()
.map(|h| {
ts.stored_value(&h.key, f)
.and_then(|raw| kevy_index::order_key(spec.values[f].ty, raw))
})
.collect::<Vec<_>>()
};
let okeys = sort_field.map(hit_keys);
let dkeys = distinct_field.map(hit_keys);
Ok((hits, spans, okeys, dkeys, facets))
});
match res {
Ok(Err(chunk)) => chunk,
Ok(Ok((hits, spans, okeys, dkeys, facets))) => {
encode_hits(store, &hits, &spans, &okeys, &dkeys, &facets, &q.fields)
}
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],
}
}
type ShardPage =
(Vec<kevy_text::TextMatch>, Option<usize>, Option<usize>, Vec<Vec<kevy_text::Bucket>>);
fn scored_hits(
ts: &kevy_text::TextSegment,
spec: &kevy_index::IndexSpec,
q: &super::args::MatchArgs,
stats: &kevy_text::CorpusStats,
) -> Result<ShardPage, Vec<u8>> {
let scope = scope_positions(spec, &q.scope)?;
let tests = boxed_preds(spec, &q.filters)?;
let filter: Vec<kevy_text::Filter> = tests
.iter()
.map(|(field, test)| kevy_text::Filter { field: *field, test: test.as_ref() })
.collect();
let sorted = sort_field(spec, &q.sort)?;
let grouped = distinct_field(spec, &q.distinct)?;
let dkey = grouped.map(|(_, ty)| move |raw: &[u8]| kevy_index::order_key(ty, raw));
let distinct = grouped
.zip(dkey.as_ref())
.map(|((field, _), k)| kevy_text::Distinct { field, key: k });
let key = sorted.map(|(_, _, ty)| move |raw: &[u8]| kevy_index::order_key(ty, raw));
let sort = sorted.zip(key.as_ref()).map(|((field, desc, _), k)| kevy_text::Sort {
field,
desc,
key: k,
});
let opts = kevy_text::QueryOpts {
stats: Some(stats),
typo: q.typo,
fields: &scope,
filter: &filter,
sort,
distinct,
};
let counted = facet_fields(spec, &q.facets)?;
let fkeys: Vec<_> =
counted.iter().map(|(_, ty)| { let ty = *ty; move |raw: &[u8]| kevy_index::order_key(ty, raw) }).collect();
let facets: Vec<kevy_text::Facet> = counted
.iter()
.zip(&fkeys)
.map(|((field, _), k)| kevy_text::Facet { field: *field, key: k })
.collect();
let r = ts.matches_query_faceted(&q.text, q.limit + q.offset, opts, &facets);
Ok((r.hits, sorted.map(|(field, _, _)| field), grouped.map(|(field, _)| field), r.facets))
}
fn hit_highlight(
ts: &kevy_text::TextSegment,
spec: &kevy_index::IndexSpec,
key: &[u8],
text: &[u8],
want: &[Vec<u8>],
) -> super::HitSpans {
ts.highlight_spans(key, text)
.into_iter()
.filter_map(|(fi, spans)| {
let name = spec.fields.get(fi)?.name.clone();
if !want.is_empty() && !want.contains(&name) {
return None;
}
let ranges = spans.into_iter().map(|(s, e)| (s as u32, e as u32)).collect();
Some((name, ranges))
})
.collect()
}
pub(super) fn op_agg(ctx: &Ctx<'_>, store: &mut Store, argv: &[Vec<u8>]) -> Vec<u8> {
let single = argv[2].eq_ignore_ascii_case(b"GROUP");
if single {
let res = index_runtime::with_ready_agg(ctx, store, &argv[1], |a| {
argv.get(3).map(|g| vec![(g.clone(), a.group(g))])
});
return match res {
Ok(None) => vec![ST_BADARGS],
Ok(Some(rows)) => encode_agg_chunk(&rows),
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],
};
}
let Some((by, limit)) = super::args::parse_groups_args(argv) else {
return vec![ST_BADARGS];
};
let depth: usize = argv
.iter()
.find_map(|a| std::str::from_utf8(a).ok()?.strip_prefix("DEPTH=")?.parse().ok())
.unwrap_or(1);
let res = index_runtime::with_ready_agg(ctx, store, &argv[1], |a| {
if depth == 0 {
return (a.all_groups(), true);
}
let fetch = (limit * 4 * depth).max(64 * depth);
let rows = a.top_groups(by, fetch);
let exhausted = rows.len() < fetch;
(rows, exhausted)
});
match res {
Ok((rows, exhausted)) => {
let mut chunk = encode_agg_chunk(&rows);
chunk.push(u8::from(exhausted));
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 op_agg_fetch(ctx: &Ctx<'_>, store: &mut Store, argv: &[Vec<u8>]) -> Vec<u8> {
let res = index_runtime::with_ready_agg(ctx, store, &argv[1], |a| {
argv[2..].iter().map(|g| (g.clone(), a.group(g))).collect::<Vec<_>>()
});
match res {
Ok(rows) => encode_agg_chunk(&rows),
Err(e) if e.as_wire().starts_with("INDEXBUILDING") => vec![ST_BUILDING],
Err(_) => vec![ST_NOINDEX],
}
}
pub(super) fn op_rebuild(ctx: &Ctx<'_>, store: &mut Store, argv: &[Vec<u8>]) -> Vec<u8> {
let Some(name) = argv.get(1) else {
return vec![ST_BADARGS];
};
match index_runtime::with_ready_ann(ctx, store, name, |g| g.rebuild()) {
Ok(()) => vec![ST_OK],
Err(e) if e.as_wire().starts_with("INDEXBUILDING") => vec![ST_BUILDING],
Err(_) => vec![ST_NOINDEX],
}
}
pub(super) fn op_knn(ctx: &Ctx<'_>, store: &mut Store, argv: &[Vec<u8>]) -> Vec<u8> {
let Some(q) = KnnArgs::parse(argv) else {
return vec![ST_BADARGS];
};
let res = index_runtime::with_ready_ann(ctx, store, &q.name, |g| {
kevy_vector::parse_vector(&q.vec, g.dim()).map(|v| g.knn(&v, q.limit, q.ef))
});
match res {
Ok(None) => vec![ST_BADARGS], Ok(Some(hits)) => {
let mut chunk = vec![ST_OK];
chunk.extend_from_slice(&(hits.len() as u32).to_le_bytes());
let keys: Vec<&[u8]> = hits.iter().map(|(k, _)| k.as_slice()).collect();
let rows = peek_hydration(store, &keys, &q.fields);
for (i, (key, dist)) in hits.iter().enumerate() {
chunk.extend_from_slice(&(key.len() as u32).to_le_bytes());
chunk.extend_from_slice(key);
chunk.extend_from_slice(&f64::from(*dist).to_le_bytes());
encode_hydration_row(&mut chunk, q.fields.len(), &rows[i]);
}
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 op_hybrid(ctx: &Ctx<'_>, store: &mut Store, argv: &[Vec<u8>]) -> Vec<u8> {
let Some(q) = HybridArgs::parse(argv) else {
return vec![ST_BADARGS];
};
let depth = q.limit * 4;
let m = index_runtime::with_ready_text_segment(ctx, store, &q.text_idx, |ts, _| {
ts.matches(&q.text, depth)
});
let k = index_runtime::with_ready_ann(ctx, store, &q.ann_idx, |g| {
kevy_vector::parse_vector(&q.vec, g.dim()).map(|v| g.knn(&v, depth, q.ef))
});
let (m, k) = match (m, k) {
(Ok(m), Ok(Some(k))) => (m, k),
(_, Ok(None)) => return vec![ST_BADARGS],
(Err(e), _) | (_, Err(e)) if e.as_wire().starts_with("INDEXBUILDING") => {
return vec![ST_BUILDING];
}
(Err(e), _) | (_, Err(e)) if e.as_wire().starts_with("INDEXOVERBUDGET") => {
return vec![ST_OVERBUDGET];
}
_ => return vec![ST_NOINDEX],
};
let mut chunk = vec![ST_OK];
chunk.extend_from_slice(&(m.len() as u32).to_le_bytes());
let keys: Vec<&[u8]> = m
.iter()
.map(|h| h.key.as_slice())
.chain(k.iter().map(|(key, _)| key.as_slice()))
.collect();
let rows = peek_hydration(store, &keys, &q.fields);
for (i, h) in m.iter().enumerate() {
chunk.extend_from_slice(&(h.key.len() as u32).to_le_bytes());
chunk.extend_from_slice(&h.key);
chunk.extend_from_slice(&h.score.to_le_bytes());
encode_hydration_row(&mut chunk, q.fields.len(), &rows[i]);
}
chunk.extend_from_slice(&(k.len() as u32).to_le_bytes());
for (i, (key, dist)) in k.iter().enumerate() {
chunk.extend_from_slice(&(key.len() as u32).to_le_bytes());
chunk.extend_from_slice(key);
chunk.extend_from_slice(&f64::from(*dist).to_le_bytes());
encode_hydration_row(&mut chunk, q.fields.len(), &rows[m.len() + i]);
}
chunk
}
pub(super) fn op_compose(ctx: &Ctx<'_>, store: &mut Store, argv: &[Vec<u8>]) -> Vec<u8> {
let Some(cq) = ComposeQuery::parse(argv) else {
return vec![ST_BADARGS];
};
let res = index_runtime::with_two_ready_segments(
ctx,
store,
&cq.a.name,
&cq.b.name,
|spec_a, seg_a, spec_b, seg_b| compose_keys(&cq, spec_a.ty, seg_a, spec_b.ty, seg_b),
);
match res {
Ok(Some(keys)) => {
let mut chunk = vec![ST_OK];
chunk.extend_from_slice(&(keys.len() as u32).to_le_bytes());
let krefs: Vec<&[u8]> = keys.iter().map(Vec::as_slice).collect();
let rows = peek_hydration(store, &krefs, &cq.fields);
for (i, k) in keys.iter().enumerate() {
chunk.extend_from_slice(&(k.len() as u32).to_le_bytes());
chunk.extend_from_slice(k);
encode_hydration_row(&mut chunk, cq.fields.len(), &rows[i]);
}
chunk
}
Ok(None) => vec![ST_BADARGS],
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 compose_keys(
cq: &ComposeQuery,
ty_a: ValType,
seg_a: &kevy_index::Segment,
ty_b: ValType,
seg_b: &kevy_index::Segment,
) -> Option<Vec<Vec<u8>>> {
let (min_a, max_a) = sub_bounds(&cq.a.shape, ty_a)?;
let (min_b, max_b) = sub_bounds(&cq.b.shape, ty_b)?;
let (a_hits, _) = seg_a.range(&min_a, &max_a, None, usize::MAX);
let mut keys: Vec<Vec<u8>> = if cq.and {
a_hits
.into_iter()
.filter(|(k, _)| {
seg_b
.verify_entry(k)
.is_some_and(|v| *v >= min_b && *v <= max_b)
})
.map(|(k, _)| k)
.collect()
} else {
let (b_hits, _) = seg_b.range(&min_b, &max_b, None, usize::MAX);
let mut all: Vec<Vec<u8>> =
a_hits.into_iter().chain(b_hits).map(|(k, _)| k).collect();
all.sort();
all.dedup();
all
};
keys.sort();
if let Some(cur) = &cq.cursor_key {
keys.retain(|k| k.as_slice() > cur.as_slice());
}
keys.truncate(cq.limit);
Some(keys)
}
fn sub_bounds(shape: &Shape, ty: ValType) -> Option<(IndexValue, IndexValue)> {
match shape {
Shape::Range { min, max } => Some((
IndexValue::parse_literal(ty, min)?,
IndexValue::parse_literal(ty, max)?,
)),
Shape::Eq { value } => {
let v = IndexValue::parse_literal(ty, value)?;
Some((v.clone(), v))
}
Shape::Where(_) | Shape::Verify => None,
}
}