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::scope_positions;
#[path = "ops_cold.rs"]
mod cold_seam;
use cold_seam::{cold_refusal, merge_cold_stats};
#[path = "ops_match.rs"]
mod match_page;
use match_page::{hit_highlight, order_keys, scored_hits};
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, cold| {
if let Some(chunk) = cold_refusal(&q, cold) {
return Err(chunk);
}
let want = scope_positions(spec, &q.scope)?;
let opts = kevy_text::QueryOpts {
stats: None,
typo: q.typo,
fields: &want,
filter: &[],
sort: None,
distinct: None,
};
let (mut n_docs, mut total_len, mut tokdf) =
(ts.docs(), ts.total_len_in(&want), ts.query_df_in(&q.text, opts));
merge_cold_stats(cold, &mut n_docs, &mut total_len, &mut tokdf);
Ok((n_docs, total_len, tokdf))
});
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, |store, ts, spec, cold| {
if let Some(chunk) = cold_refusal(&q, cold) {
return Err(chunk);
}
let ((hits, sort_field, distinct_field, facets), cold_vals) =
scored_hits(ts, spec, &q, &stats, cold)?;
let spans = q.highlight.as_ref().map(|want| {
hits.iter()
.map(|h| {
let chilled = cold.is_some_and(index_runtime::TextColdDir::has_cold);
hit_highlight(store, ts, spec, &h.key, &q.text, want, chilled)
})
.collect::<Vec<_>>()
});
let hit_keys = |f: usize| order_keys(ts, spec, &cold_vals, &hits, f);
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],
}
}
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,
}
}