use kevy_index::{IndexSpec, Segment};
use kevy_resp::CmdError;
use kevy_store::Store;
use crate::state::{CatalogState, Ctx};
enum BuildState {
Backfilling { keys: Vec<Vec<u8>>, pos: usize },
Ready,
FailedOverBudget,
}
struct ShardIndex {
spec: IndexSpec,
seg: Segment,
window: Option<kevy_window::WindowRt>,
cold_text: Option<TextColdDir>,
text: Option<kevy_text::TextSegment>,
ann: Option<kevy_vector::Hnsw>,
agg: Option<kevy_index::AggSegment>,
build: BuildState,
}
#[derive(Default)]
pub(crate) struct ShardIndexes {
generation: u64,
idx: Vec<ShardIndex>,
stats_dirty: bool,
reserved_cache: u64,
}
#[inline]
pub(crate) fn on_write(ctx: &Ctx<'_>, store: &mut Store, key: &[u8]) {
let mut st = ctx.shard.indexes.borrow_mut();
refresh(&ctx.state.catalogs, &mut st, store);
let st = &mut *st;
for si in &mut st.idx {
if key.starts_with(&si.spec.prefix) {
apply_row(store, si, key);
st.stats_dirty = true;
}
}
}
pub(crate) fn on_tick(ctx: &Ctx<'_>, store: &mut Store) {
let mut st = ctx.shard.indexes.borrow_mut();
refresh(&ctx.state.catalogs, &mut st, store);
let st = &mut *st;
let segs_dir = shard_segs_dir(ctx.state, ctx.shard.shard_id());
let mut batches: Vec<(Vec<u8>, Vec<Vec<u8>>)> = Vec::new();
for si in &mut st.idx {
if matches!(si.build, BuildState::Backfilling { .. }) {
st.stats_dirty = true;
}
advance_backfill(store, si, 2048);
if let (Some(win), Some(dir), BuildState::Ready) = (&mut si.window, &segs_dir, &si.build) {
let drives = window_driver(&ctx.state.catalogs, &si.spec.name);
if drives && let Some(rows) = win.pending_rows(&si.seg) {
batches.push((table_of(&si.spec.name).to_vec(), rows));
}
st.stats_dirty |= evict_and_slide(win, &si.spec.name, &mut si.seg, store, dir, drives);
}
}
if let Some(dir) = &segs_dir {
for (table, keys) in &batches {
freeze_text_batches(st, table, keys, dir);
}
}
}
pub(crate) fn on_flush(ctx: &Ctx<'_>, store: &mut Store) {
let mut st = ctx.shard.indexes.borrow_mut();
refresh(&ctx.state.catalogs, &mut st, store);
let st = &mut *st;
for si in &mut st.idx {
si.seg = new_scalar_seg(&si.spec);
si.text = new_text_seg(&si.spec);
si.ann = new_ann_seg(&si.spec);
si.agg = (si.spec.kind == kevy_index::IndexKind::Agg).then(kevy_index::AggSegment::new);
si.build = BuildState::Ready;
st.stats_dirty = true;
}
}
pub(crate) fn reserved_bytes(ctx: &Ctx<'_>, store: &mut Store) -> u64 {
let mut st = ctx.shard.indexes.borrow_mut();
refresh(&ctx.state.catalogs, &mut st, store);
if !st.stats_dirty {
return st.reserved_cache;
}
let sum = st
.idx
.iter()
.map(|si| {
si.seg.stats().approx_bytes
+ si.text.as_ref().map_or(0, |t| t.stats().approx_bytes)
+ si.ann.as_ref().map_or(0, |g| g.stats().approx_bytes)
+ si.agg.as_ref().map_or(0, |a| a.stats().approx_bytes)
})
.sum();
st.reserved_cache = sum;
st.stats_dirty = false;
sum
}
pub(crate) fn with_ready_segment<R>(
ctx: &Ctx<'_>,
store: &mut Store,
name: &[u8],
f: impl FnOnce(&IndexSpec, &Segment, Option<&kevy_window::WindowRt>) -> R,
) -> Result<R, CmdError> {
let mut st = ctx.shard.indexes.borrow_mut();
refresh(&ctx.state.catalogs, &mut st, store);
let si = st.idx.iter().find(|si| si.spec.name == name).ok_or("ERR no such index")?;
match si.build {
BuildState::Ready => Ok(f(&si.spec, &si.seg, si.window.as_ref())),
BuildState::Backfilling { .. } => {
Err(CmdError::Wire("INDEXBUILDING index is still building"))
}
BuildState::FailedOverBudget => {
Err(CmdError::Wire("INDEXOVERBUDGET index build exceeded MAXMEM"))
}
}
}
pub(crate) fn with_ready_agg<R>(
ctx: &Ctx<'_>,
store: &mut Store,
name: &[u8],
f: impl FnOnce(&kevy_index::AggSegment) -> R,
) -> Result<R, CmdError> {
let mut st = ctx.shard.indexes.borrow_mut();
refresh(&ctx.state.catalogs, &mut st, store);
let si = st.idx.iter().find(|si| si.spec.name == name).ok_or("ERR no such index")?;
match (&si.build, &si.agg) {
(BuildState::Ready, Some(a)) => Ok(f(a)),
(BuildState::Backfilling { .. }, _) => {
Err(CmdError::Wire("INDEXBUILDING index is still building"))
}
(BuildState::FailedOverBudget, _) => {
Err(CmdError::Wire("INDEXOVERBUDGET index build exceeded MAXMEM"))
}
(_, None) => Err(CmdError::Wire("ERR not an aggregate index")),
}
}
pub(crate) fn with_ready_ann<R>(
ctx: &Ctx<'_>,
store: &mut Store,
name: &[u8],
f: impl FnOnce(&mut kevy_vector::Hnsw) -> R,
) -> Result<R, CmdError> {
let mut st = ctx.shard.indexes.borrow_mut();
refresh(&ctx.state.catalogs, &mut st, store);
let si = st.idx.iter_mut().find(|si| si.spec.name == name).ok_or("ERR no such index")?;
match (&si.build, &mut si.ann) {
(BuildState::Ready, Some(g)) => Ok(f(g)),
(BuildState::Backfilling { .. }, _) => {
Err(CmdError::Wire("INDEXBUILDING index is still building"))
}
(BuildState::FailedOverBudget, _) => {
Err(CmdError::Wire("INDEXOVERBUDGET index build exceeded MAXMEM"))
}
(_, None) => Err(CmdError::Wire("ERR not a vector index")),
}
}
pub(crate) fn with_ready_text_segment<R>(
ctx: &Ctx<'_>,
store: &mut Store,
name: &[u8],
f: impl FnOnce(
&mut Store,
&kevy_text::TextSegment,
&kevy_index::IndexSpec,
Option<&TextColdDir>,
) -> R,
) -> Result<R, CmdError> {
let mut st = ctx.shard.indexes.borrow_mut();
refresh(&ctx.state.catalogs, &mut st, store);
let si = st.idx.iter().find(|si| si.spec.name == name).ok_or("ERR no such index")?;
match (&si.build, &si.text) {
(BuildState::Ready, Some(ts)) => Ok(f(store, ts, &si.spec, si.cold_text.as_ref())),
(BuildState::Backfilling { .. }, _) => {
Err(CmdError::Wire("INDEXBUILDING index is still building"))
}
(BuildState::FailedOverBudget, _) => {
Err(CmdError::Wire("INDEXOVERBUDGET index build exceeded MAXMEM"))
}
(_, None) => Err(CmdError::Wire("ERR not a text index")),
}
}
pub(crate) fn with_segment_resolver<R>(
ctx: &Ctx<'_>,
store: &mut Store,
f: impl for<'s> FnOnce(&'s dyn Fn(&[u8]) -> Option<&'s Segment>) -> R,
) -> R {
let mut st = ctx.shard.indexes.borrow_mut();
refresh(&ctx.state.catalogs, &mut st, store);
let idx = &st.idx;
let resolver = |name: &[u8]| -> Option<&Segment> {
idx.iter()
.find(|si| si.spec.name == name && matches!(si.build, BuildState::Ready))
.map(|si| &si.seg)
};
f(&resolver)
}
pub(crate) fn with_two_ready_segments<R>(
ctx: &Ctx<'_>,
store: &mut Store,
a: &[u8],
b: &[u8],
f: impl FnOnce(&IndexSpec, &Segment, &IndexSpec, &Segment) -> R,
) -> Result<R, CmdError> {
let mut st = ctx.shard.indexes.borrow_mut();
refresh(&ctx.state.catalogs, &mut st, store);
let ia = st.idx.iter().position(|si| si.spec.name == a).ok_or("ERR no such index")?;
let ib = st.idx.iter().position(|si| si.spec.name == b).ok_or("ERR no such index")?;
for i in [ia, ib] {
if matches!(st.idx[i].build, BuildState::Backfilling { .. }) {
return Err(CmdError::Wire("INDEXBUILDING index is still building"));
}
}
let (sa, sb) = (&st.idx[ia], &st.idx[ib]);
Ok(f(&sa.spec, &sa.seg, &sb.spec, &sb.seg))
}
pub(crate) fn segment_building(ctx: &Ctx<'_>, store: &mut Store, name: &[u8]) -> bool {
let mut st = ctx.shard.indexes.borrow_mut();
refresh(&ctx.state.catalogs, &mut st, store);
st.idx
.iter()
.find(|si| si.spec.name == name)
.is_some_and(|si| matches!(si.build, BuildState::Backfilling { .. }))
}
fn new_scalar_seg(spec: &IndexSpec) -> Segment {
let scalar = matches!(spec.kind, kevy_index::IndexKind::Range | kevy_index::IndexKind::Unique);
if scalar && !spec.values.is_empty() {
Segment::with_values(spec.values.len())
} else {
Segment::new()
}
}
fn new_ann_seg(spec: &kevy_index::IndexSpec) -> Option<kevy_vector::Hnsw> {
spec.ann.as_ref().map(|a| {
kevy_vector::Hnsw::new(
a.dim as usize,
kevy_vector::HnswParams {
m: a.m as usize,
ef_construction: a.ef as usize,
distance: match a.distance {
1 => kevy_vector::Distance::L2,
2 => kevy_vector::Distance::Ip,
_ => kevy_vector::Distance::Cosine,
},
},
)
})
}
fn new_text_seg(spec: &kevy_index::IndexSpec) -> Option<kevy_text::TextSegment> {
(spec.kind == kevy_index::IndexKind::Text).then(|| {
kevy_text::TextSegment::with_shape(kevy_text::SegmentShape {
fields: spec.fields.len(),
positions: spec.with_positions,
values: spec.values.len(),
})
})
}
fn refresh(catalogs: &CatalogState, st: &mut ShardIndexes, store: &mut Store) {
let generation = catalogs.index_gen();
if st.generation == generation {
return;
}
st.stats_dirty = true;
let cat = catalogs.index();
let mut next: Vec<ShardIndex> = Vec::new();
if let Some(cat) = cat {
for (spec, _state) in cat.iter() {
match st.idx.iter().position(|si| si.spec == *spec) {
Some(i) => {
let mut si = st.idx.swap_remove(i);
let want = window_for(catalogs, &si.spec);
let have = si.window.as_ref().map(|w| (w.spec.clone(), w.shape));
if have != want {
si.window = want.map(|(w, sh)| kevy_window::WindowRt::new(w, sh));
}
let want_text = text_window_for(catalogs, &si.spec);
if si.cold_text.is_some() != want_text {
si.cold_text = want_text.then(TextColdDir::new);
}
next.push(si);
}
None => next.push(fresh_shard_index(catalogs, spec, store)),
}
}
}
st.idx = next;
st.generation = generation;
}
fn fresh_shard_index(catalogs: &CatalogState, spec: &IndexSpec, store: &mut Store) -> ShardIndex {
let mut pat = spec.prefix.clone();
pat.push(b'*');
let keys = store.collect_keys(Some(&pat), None);
ShardIndex {
agg: (spec.kind == kevy_index::IndexKind::Agg).then(kevy_index::AggSegment::new),
text: new_text_seg(spec),
ann: new_ann_seg(spec),
seg: new_scalar_seg(spec),
window: window_for(catalogs, spec).map(|(w, sh)| kevy_window::WindowRt::new(w, sh)),
cold_text: text_window_for(catalogs, spec).then(TextColdDir::new),
spec: spec.clone(),
build: BuildState::Backfilling { keys, pos: 0 },
}
}
pub(crate) use kevy_window::{ColdHit, ColdPageQuery, TextColdDir, WindowRt};
mod window_slide;
use window_slide::{
evict_and_slide, freeze_text_batches, shard_segs_dir, table_of, text_window_for, window_driver,
window_for,
};
mod row_apply;
pub(crate) use row_apply::{RowValue, row_value};
use row_apply::{advance_backfill, apply_row};
#[cfg(test)]
mod tests;