use std::collections::HashSet;
use std::path::PathBuf;
use std::sync::atomic::{AtomicU64, Ordering};
use std::time::Instant;
use anyhow::{Context, Result};
use clap::Parser;
use serde::Serialize;
use uuid::Uuid;
use khive_mcp::serve::{resolve_runtime_config, RuntimeConfigInputs};
use khive_runtime::{entity_fts_document, note_fts_document, KhiveRuntime, Namespace};
use khive_storage::entity::Entity;
use khive_storage::error::StorageError;
use khive_storage::note::Note;
use khive_storage::VectorStore;
use khive_types::SubstrateKind;
const MAX_EMBED_BYTES: usize = 32_768;
struct ProgressBar {
label: &'static str,
start: Instant,
current: AtomicU64,
total: AtomicU64,
window_current: AtomicU64,
window_nanos: AtomicU64,
rate: std::sync::Mutex<f64>,
}
const RATE_WINDOW_SECS: f64 = 10.0;
impl ProgressBar {
fn new(label: &'static str) -> Self {
Self {
label,
start: Instant::now(),
current: AtomicU64::new(0),
total: AtomicU64::new(0),
window_current: AtomicU64::new(0),
window_nanos: AtomicU64::new(0),
rate: std::sync::Mutex::new(0.0),
}
}
fn update(&self, current: u64, total: u64) {
self.current.store(current, Ordering::Relaxed);
self.total.store(total, Ordering::Relaxed);
let now_ns = self.start.elapsed().as_nanos() as u64;
let prev_ns = self.window_nanos.load(Ordering::Relaxed);
let delta_secs = (now_ns - prev_ns) as f64 / 1e9;
if delta_secs >= RATE_WINDOW_SECS {
let prev_current = self.window_current.load(Ordering::Relaxed);
let delta_items = current.saturating_sub(prev_current);
if delta_secs > 0.1 {
let window_rate = delta_items as f64 / delta_secs;
if let Ok(mut r) = self.rate.lock() {
if *r < 0.1 {
*r = window_rate;
} else {
*r = 0.3 * *r + 0.7 * window_rate;
}
}
}
self.window_current.store(current, Ordering::Relaxed);
self.window_nanos.store(now_ns, Ordering::Relaxed);
}
self.render();
}
fn render(&self) {
use std::io::Write;
let current = self.current.load(Ordering::Relaxed);
let total = self.total.load(Ordering::Relaxed);
let pct = if total > 0 {
(current as f64 / total as f64 * 100.0).min(100.0)
} else {
0.0
};
const BAR_WIDTH: usize = 30;
let filled = (pct / 100.0 * BAR_WIDTH as f64) as usize;
let empty = BAR_WIDTH.saturating_sub(filled);
let bar: String = format!("{}{}", "\u{2588}".repeat(filled), "\u{2591}".repeat(empty),);
let rate = self.rate.lock().map(|r| *r).unwrap_or(0.0);
let eta = if rate > 0.1 && current < total {
let remaining = (total - current) as f64 / rate;
if remaining >= 60.0 {
format!(
"ETA {}m {:02}s",
remaining as u64 / 60,
remaining as u64 % 60
)
} else {
format!("ETA {:.0}s", remaining)
}
} else if current >= total && total > 0 {
"done".into()
} else {
"warming up…".into()
};
eprint!(
"\r {:<10} [{bar}] {pct:>5.1}% ({current}/{total}) {rate:>6.0}/s {eta} ",
self.label,
);
let _ = std::io::stderr().flush();
}
fn finish(&self) {
self.render();
eprintln!();
}
}
#[derive(Parser, Debug)]
pub struct ReindexArgs {
#[arg(long, env = "KHIVE_DB")]
pub db: Option<String>,
#[arg(long = "config", env = "KHIVE_CONFIG")]
pub config: Option<PathBuf>,
#[arg(long)]
pub model: Option<String>,
#[arg(long, default_value = "100")]
pub batch_size: u32,
#[arg(long)]
pub keep_existing: bool,
#[arg(long, env = "KHIVE_NAMESPACE")]
pub namespace: Option<String>,
#[arg(long, conflicts_with = "no_knowledge")]
pub knowledge_only: bool,
#[arg(long)]
pub no_knowledge: bool,
#[arg(long)]
pub best_effort: bool,
#[arg(long, conflicts_with = "sections_only")]
pub no_sections: bool,
#[arg(long, conflicts_with = "no_knowledge")]
pub sections_only: bool,
#[arg(long)]
pub human: bool,
}
#[derive(Serialize)]
struct ReindexReport {
entities_processed: u64,
notes_processed: u64,
#[serde(skip_serializing_if = "Option::is_none")]
knowledge_atoms_indexed: Option<u64>,
#[serde(skip_serializing_if = "Option::is_none")]
knowledge_sections_indexed: Option<u64>,
knowledge_atoms_failed: u64,
knowledge_pass_errored: bool,
knowledge_ann_failed: bool,
knowledge_sections_failed: u64,
models_used: Vec<String>,
elapsed_ms: u64,
errors_skipped: u64,
entities_fts_failed: u64,
notes_fts_failed: u64,
}
impl ReindexReport {
fn has_failures(&self) -> bool {
self.errors_skipped > 0
|| self.entities_fts_failed > 0
|| self.notes_fts_failed > 0
|| self.knowledge_atoms_failed > 0
|| self.knowledge_pass_errored
|| self.knowledge_ann_failed
|| self.knowledge_sections_failed > 0
}
}
#[allow(clippy::too_many_arguments)]
async fn embed_and_store_batch(
rt: &KhiveRuntime,
token: &khive_runtime::NamespaceToken,
model_names: &[String],
namespace: &str,
staged: &[(Uuid, String)],
kind: SubstrateKind,
field: &str,
drop_existing: bool,
) -> u64 {
let mut errors: u64 = 0;
for model_name in model_names {
let vectors = match rt.vectors_for_model(token, model_name) {
Ok(v) => v,
Err(e) => {
tracing::warn!(model = %model_name, error = %e, "vector store unavailable");
errors += staged.len() as u64;
continue;
}
};
let subset: Vec<&(Uuid, String)> = if drop_existing {
staged.iter().collect()
} else {
let ids: Vec<Uuid> = staged.iter().map(|(id, _)| *id).collect();
match filter_unembedded(vectors.as_ref(), &ids, namespace).await {
Ok(unembedded) => {
let keep: HashSet<Uuid> = unembedded.into_iter().collect();
staged.iter().filter(|(id, _)| keep.contains(id)).collect()
}
Err(e) => {
tracing::error!(model = %model_name, error = %e, "filter_unembedded failed; skipping batch for this model");
errors += staged.len() as u64;
continue;
}
}
};
if subset.is_empty() {
continue;
}
let texts: Vec<String> = subset.iter().map(|(_, t)| truncate_text(t)).collect();
match rt.embed_document_batch_with_model(model_name, &texts).await {
Ok(embeddings) if embeddings.len() == subset.len() => {
for ((id, _), emb) in subset.iter().zip(embeddings.iter()) {
if drop_existing {
let _ = vectors.delete(*id).await;
}
if let Err(e) = vectors
.insert(*id, kind, namespace, field, vec![emb.clone()])
.await
{
tracing::warn!(id = %id, model = %model_name, error = %e, "vector insert failed");
errors += 1;
}
}
}
Ok(_) => {
tracing::warn!(model = %model_name, "embedding count mismatch for batch");
errors += subset.len() as u64;
}
Err(e) => {
tracing::warn!(model = %model_name, error = %e, "embed_batch failed");
errors += subset.len() as u64;
}
}
}
errors
}
async fn fts_backfill_notes_batch(
rt: &KhiveRuntime,
token: &khive_runtime::NamespaceToken,
batch: &[Note],
) -> u64 {
let fts = match rt.text_for_notes(token) {
Ok(f) => f,
Err(e) => {
tracing::error!(error = %e, "FTS store unavailable; counting whole batch as failed");
return batch.len() as u64;
}
};
let mut errors: u64 = 0;
for note in batch {
let doc = note_fts_document(note);
if let Err(e) = fts.upsert_document(doc).await {
tracing::warn!(id = %note.id, error = %e, "FTS upsert failed for note");
errors += 1;
}
}
errors
}
async fn fts_backfill_entities_batch(
rt: &KhiveRuntime,
token: &khive_runtime::NamespaceToken,
batch: &[Entity],
) -> u64 {
let fts = match rt.text(token) {
Ok(f) => f,
Err(e) => {
tracing::error!(error = %e, "FTS store unavailable; counting whole batch as failed");
return batch.len() as u64;
}
};
let mut errors: u64 = 0;
for entity in batch {
let doc = entity_fts_document(entity);
if let Err(e) = fts.upsert_document(doc).await {
tracing::warn!(id = %entity.id, error = %e, "FTS upsert failed for entity");
errors += 1;
}
}
errors
}
async fn filter_unembedded(
vectors: &dyn VectorStore,
ids: &[Uuid],
namespace: &str,
) -> Result<Vec<Uuid>> {
match vectors.batch_exists(ids, namespace).await {
Ok(existing) => Ok(ids
.iter()
.filter(|id| !existing.contains(id))
.copied()
.collect()),
Err(StorageError::Unsupported { .. }) => Ok(ids.to_vec()),
Err(e) => Err(anyhow::anyhow!("{e}")),
}
}
pub async fn run_reindex(args: ReindexArgs) -> Result<()> {
let explicit = args.namespace.is_some();
let raw = args.namespace.as_deref().unwrap_or("local");
let ns = Namespace::parse(raw).map_err(|e| anyhow::anyhow!("{e}"))?;
let cfg = resolve_runtime_config(RuntimeConfigInputs {
db: args.db.as_deref(),
config: args.config.as_deref(),
namespace: ns,
namespace_explicit: explicit,
no_embed: false,
packs: None,
brain_profile: None,
})?;
let resolved_ns = cfg.default_namespace.clone();
let rt = KhiveRuntime::new(cfg).map_err(|e| anyhow::anyhow!("{e}"))?;
let token = rt
.authorize(resolved_ns)
.map_err(|e| anyhow::anyhow!("{e}"))
.context("failed to authorize namespace")?;
let do_graph = !args.knowledge_only && !args.sections_only; let do_knowledge = !args.no_knowledge; let do_atoms = do_knowledge && !args.sections_only;
let do_sections = do_knowledge && !args.no_sections;
let model_names: Vec<String> = if !do_graph {
vec![]
} else {
match args.model.as_deref().filter(|s| !s.is_empty()) {
Some(name) => vec![name.to_string()],
None => {
let names = rt.registered_embedding_model_names();
if names.is_empty() {
eprintln!("warning: no embedding model configured — skipping vector embedding; FTS backfill will still run");
}
names
}
}
};
let batch_size = args.batch_size.clamp(1, 500);
let drop_existing = !args.keep_existing;
let ns_str = token.namespace().as_str().to_owned();
let start = std::time::Instant::now();
let mut entities_processed: u64 = 0;
let mut notes_processed: u64 = 0;
let mut errors_skipped: u64 = 0;
let mut entities_fts_failed: u64 = 0;
let mut notes_fts_failed: u64 = 0;
if do_graph {
let entity_total = rt.count_entities(&token, None).await.unwrap_or(0);
let entity_bar = ProgressBar::new("entities");
entity_bar.update(0, entity_total);
let mut entity_offset: u32 = 0;
loop {
let batch = rt
.list_entities(&token, None, None, batch_size, entity_offset)
.await
.map_err(|e| anyhow::anyhow!("{e}"))?;
let n = batch.len();
if n == 0 {
break;
}
let mut staged: Vec<(Uuid, String)> = Vec::with_capacity(n);
for entity in &batch {
let text = match &entity.description {
Some(d) if !d.is_empty() => format!("{} {}", entity.name, d),
_ => entity.name.clone(),
};
if !text.trim().is_empty() {
staged.push((entity.id, text));
}
}
if !staged.is_empty() {
errors_skipped += embed_and_store_batch(
&rt,
&token,
&model_names,
&ns_str,
&staged,
SubstrateKind::Entity,
"entity.body",
drop_existing,
)
.await;
entities_processed += staged.len() as u64;
}
entities_fts_failed += fts_backfill_entities_batch(&rt, &token, &batch).await;
entity_bar.update(entities_processed, entity_total);
if n < batch_size as usize {
break;
}
entity_offset += n as u32;
}
entity_bar.finish();
let note_total = count_notes(&rt, &ns_str).await;
let note_bar = ProgressBar::new("notes");
note_bar.update(0, note_total);
let mut note_offset: u32 = 0;
loop {
let batch = rt
.list_notes(&token, None, batch_size, note_offset)
.await
.map_err(|e| anyhow::anyhow!("{e}"))?;
let n = batch.len();
if n == 0 {
break;
}
let mut staged: Vec<(Uuid, String)> = Vec::with_capacity(n);
for note in &batch {
let text = note.content.clone();
if !text.trim().is_empty() {
staged.push((note.id, text));
}
}
if !staged.is_empty() {
errors_skipped += embed_and_store_batch(
&rt,
&token,
&model_names,
&ns_str,
&staged,
SubstrateKind::Note,
"note.content",
drop_existing,
)
.await;
notes_processed += staged.len() as u64;
}
notes_fts_failed += fts_backfill_notes_batch(&rt, &token, &batch).await;
note_bar.update(notes_processed, note_total);
if n < batch_size as usize {
break;
}
note_offset += n as u32;
}
note_bar.finish();
if let Err(e) = invalidate_vamana_snapshots(&rt, &ns_str).await {
tracing::warn!(error = %e, "failed to invalidate Vamana snapshots after reindex");
}
}
let mut knowledge_atoms_indexed: Option<u64> = None;
let mut knowledge_sections_indexed: Option<u64> = None;
let mut knowledge_atoms_failed: u64 = 0;
let mut knowledge_pass_errored = false;
let mut knowledge_ann_failed = false;
let mut knowledge_sections_failed: u64 = 0;
if do_atoms || do_sections {
let atom_bar = ProgressBar::new("atoms");
let section_bar = ProgressBar::new("sections");
let on_atom = |c: u64, t: u64| atom_bar.update(c, t);
let on_section = |c: u64, t: u64| section_bar.update(c, t);
let opts = khive_pack_knowledge::KnowledgeReindexOptions {
atoms: do_atoms,
sections: do_sections,
drop_existing,
rebuild_ann: true,
batch_size: Some(batch_size),
};
match khive_pack_knowledge::reindex_knowledge(
&rt,
&token,
opts,
if do_atoms { Some(&on_atom) } else { None },
if do_sections { Some(&on_section) } else { None },
)
.await
{
Ok(v) => {
if do_atoms {
knowledge_atoms_indexed =
Some(v.get("atoms_indexed").and_then(|n| n.as_u64()).unwrap_or(0));
knowledge_atoms_failed = v.get("failed").and_then(|n| n.as_u64()).unwrap_or(0);
knowledge_ann_failed = v
.get("ann_failed")
.and_then(|b| b.as_bool())
.unwrap_or(false);
}
if do_sections {
knowledge_sections_indexed = Some(
v.get("sections_indexed")
.and_then(|n| n.as_u64())
.unwrap_or(0),
);
knowledge_sections_failed = v
.get("sections_failed")
.and_then(|n| n.as_u64())
.unwrap_or(0);
}
}
Err(e) => {
tracing::error!(error = %e, "knowledge reindex failed");
eprintln!("\nerror: knowledge reindex failed: {e}");
knowledge_pass_errored = true;
}
}
if do_atoms {
atom_bar.finish();
}
if do_sections {
section_bar.finish();
}
}
let elapsed_ms = start.elapsed().as_millis() as u64;
let report = ReindexReport {
entities_processed,
notes_processed,
knowledge_atoms_indexed,
knowledge_sections_indexed,
knowledge_atoms_failed,
knowledge_pass_errored,
knowledge_ann_failed,
knowledge_sections_failed,
models_used: model_names,
elapsed_ms,
errors_skipped,
entities_fts_failed,
notes_fts_failed,
};
print_report(&report, args.human);
finish(&report, args.best_effort)
}
fn decide_result(has_failures: bool, best_effort: bool) -> Result<()> {
if has_failures && !best_effort {
anyhow::bail!(
"reindex completed with failures; recall/search state may be stale. \
Re-run, or pass --best-effort to accept a partial rebuild."
);
}
Ok(())
}
fn finish(report: &ReindexReport, best_effort: bool) -> Result<()> {
let result = decide_result(report.has_failures(), best_effort);
if report.has_failures() && best_effort {
eprintln!("warning: reindex completed with failures (best-effort mode; exiting 0)");
}
result
}
async fn invalidate_vamana_snapshots(rt: &KhiveRuntime, namespace: &str) -> anyhow::Result<()> {
use khive_storage::types::{SqlStatement, SqlValue};
let pattern = format!("{namespace}::vamana::%");
let sql = rt.sql();
let mut writer = sql
.writer()
.await
.context("open SQL writer for Vamana snapshot invalidation")?;
match writer
.execute(SqlStatement {
sql: "DELETE FROM retrieval_snapshots WHERE namespace LIKE ?1".into(),
params: vec![SqlValue::Text(pattern)],
label: Some("invalidate_vamana_snapshots".into()),
})
.await
{
Ok(deleted) => {
tracing::info!(
deleted,
namespace,
"invalidated Vamana snapshots after reindex"
);
Ok(())
}
Err(e) => {
let msg = e.to_string();
if msg.contains("no such table") {
tracing::debug!("retrieval_snapshots absent; no Vamana snapshots to invalidate");
Ok(())
} else {
Err(anyhow::anyhow!("{e}"))
}
}
}
}
async fn count_notes(rt: &KhiveRuntime, ns: &str) -> u64 {
use khive_storage::types::{SqlStatement, SqlValue};
let sql = rt.sql();
let Ok(mut reader) = sql.reader().await else {
return 0;
};
let row = reader
.query_row(SqlStatement {
sql: "SELECT count(*) AS cnt FROM notes WHERE namespace = ?1 AND deleted_at IS NULL"
.into(),
params: vec![SqlValue::Text(ns.to_owned())],
label: None,
})
.await;
match row {
Ok(Some(r)) => match r.get("cnt") {
Some(SqlValue::Integer(n)) => *n as u64,
_ => 0,
},
_ => 0,
}
}
fn truncate_text(t: &str) -> String {
if t.len() <= MAX_EMBED_BYTES {
t.to_string()
} else {
let mut end = MAX_EMBED_BYTES;
while !t.is_char_boundary(end) {
end -= 1;
}
t[..end].to_string()
}
}
fn print_report(report: &ReindexReport, human: bool) {
if human {
let atoms = report
.knowledge_atoms_indexed
.map(|n| format!(", {n} knowledge atoms"))
.unwrap_or_default();
let sections = report
.knowledge_sections_indexed
.map(|n| format!(", {n} sections"))
.unwrap_or_default();
let status = if report.has_failures() {
"Reindex completed WITH FAILURES"
} else {
"Reindex complete"
};
let fts_errors = report.entities_fts_failed + report.notes_fts_failed;
println!(
"{status}: {} entities, {} notes{}{} ({} vector errors, {} FTS errors) in {}ms",
report.entities_processed,
report.notes_processed,
atoms,
sections,
report.errors_skipped,
fts_errors,
report.elapsed_ms
);
if report.entities_fts_failed > 0 {
println!(
"FTS backfill: {} entity upserts FAILED",
report.entities_fts_failed
);
}
if report.notes_fts_failed > 0 {
println!(
"FTS backfill: {} note upserts FAILED",
report.notes_fts_failed
);
}
if report.knowledge_pass_errored {
println!("Knowledge pass: FAILED (did not run to completion)");
} else if report.knowledge_atoms_failed > 0 {
println!(
"Knowledge pass: {} atom vector inserts FAILED",
report.knowledge_atoms_failed
);
}
if report.knowledge_sections_failed > 0 {
println!(
"Knowledge sections: {} section embed/write failures",
report.knowledge_sections_failed
);
}
if report.knowledge_ann_failed {
println!("Knowledge ANN: FAILED (snapshot not rebuilt/persisted)");
}
if !report.models_used.is_empty() {
println!("Models: {}", report.models_used.join(", "));
}
} else {
let json = serde_json::to_string(report).expect("serialize ReindexReport");
println!("{json}");
}
}
#[cfg(test)]
mod tests {
use super::*;
use crate::dbpath::resolve_db_override;
use clap::Parser;
use khive_storage::types::{SqlStatement, SqlValue};
use serial_test::serial;
#[tokio::test]
async fn test_reindex_invalidates_vamana_snapshots() {
let rt = KhiveRuntime::memory().expect("in-memory runtime");
let sql = rt.sql();
let mut w = sql.writer().await.expect("writer");
w.execute_script(
"CREATE TABLE IF NOT EXISTS retrieval_snapshots (\
namespace TEXT NOT NULL, \
index_type TEXT NOT NULL, \
snapshot BLOB NOT NULL, \
created_at INTEGER NOT NULL, \
PRIMARY KEY (namespace, index_type));"
.into(),
)
.await
.expect("create table");
for (ns, idx_type) in &[
("local::vamana::model-a", "vamana"),
("local::vamana::model-b", "vamana"),
("other::vamana::model-a", "vamana"),
("local::hnsw::model-a", "hnsw"),
] {
w.execute(SqlStatement {
sql: "INSERT INTO retrieval_snapshots \
(namespace, index_type, snapshot, created_at) \
VALUES (?1, ?2, ?3, 0)"
.into(),
params: vec![
SqlValue::Text(ns.to_string()),
SqlValue::Text(idx_type.to_string()),
SqlValue::Blob(b"{}".to_vec()),
],
label: None,
})
.await
.expect("insert row");
}
drop(w);
invalidate_vamana_snapshots(&rt, "local")
.await
.expect("invalidate");
let mut r = sql.reader().await.expect("reader");
let rows = r
.query_all(SqlStatement {
sql: "SELECT namespace FROM retrieval_snapshots ORDER BY namespace".into(),
params: vec![],
label: None,
})
.await
.expect("query");
let remaining: Vec<String> = rows
.iter()
.filter_map(|row| match row.get("namespace") {
Some(SqlValue::Text(s)) => Some(s.clone()),
_ => None,
})
.collect();
assert!(
remaining.contains(&"other::vamana::model-a".to_string()),
"other namespace must survive: {remaining:?}"
);
assert!(
remaining.contains(&"local::hnsw::model-a".to_string()),
"HNSW rows must survive: {remaining:?}"
);
assert!(
!remaining.contains(&"local::vamana::model-a".to_string()),
"local vamana model-a must be deleted: {remaining:?}"
);
assert!(
!remaining.contains(&"local::vamana::model-b".to_string()),
"local vamana model-b must be deleted: {remaining:?}"
);
}
fn report_with(errors: u64, k_failed: u64, k_errored: bool) -> ReindexReport {
ReindexReport {
entities_processed: 0,
notes_processed: 0,
knowledge_atoms_indexed: Some(0),
knowledge_sections_indexed: None,
knowledge_atoms_failed: k_failed,
knowledge_pass_errored: k_errored,
knowledge_ann_failed: false,
knowledge_sections_failed: 0,
models_used: vec![],
elapsed_ms: 0,
errors_skipped: errors,
entities_fts_failed: 0,
notes_fts_failed: 0,
}
}
#[test]
fn has_failures_flags_each_failure_source() {
assert!(!report_with(0, 0, false).has_failures());
assert!(
report_with(1, 0, false).has_failures(),
"entity/note errors"
);
assert!(
report_with(0, 1, false).has_failures(),
"knowledge atom fails"
);
assert!(
report_with(0, 0, true).has_failures(),
"knowledge pass error"
);
}
#[test]
fn has_failures_flags_knowledge_ann_failed() {
let report = ReindexReport {
entities_processed: 0,
notes_processed: 0,
knowledge_atoms_indexed: Some(10),
knowledge_sections_indexed: None,
knowledge_atoms_failed: 0,
knowledge_pass_errored: false,
knowledge_ann_failed: true,
knowledge_sections_failed: 0,
models_used: vec![],
elapsed_ms: 0,
errors_skipped: 0,
entities_fts_failed: 0,
notes_fts_failed: 0,
};
assert!(
report.has_failures(),
"knowledge_ann_failed alone must drive has_failures() = true"
);
assert!(
decide_result(report.has_failures(), false).is_err(),
"knowledge_ann_failed must fail closed (non-zero exit)"
);
assert!(
decide_result(report.has_failures(), true).is_ok(),
"best-effort downgrades knowledge_ann_failed to exit 0"
);
}
#[test]
fn has_failures_flags_knowledge_sections_failed() {
let report = ReindexReport {
entities_processed: 0,
notes_processed: 0,
knowledge_atoms_indexed: None,
knowledge_sections_indexed: Some(0),
knowledge_atoms_failed: 0,
knowledge_pass_errored: false,
knowledge_ann_failed: false,
knowledge_sections_failed: 3,
models_used: vec![],
elapsed_ms: 0,
errors_skipped: 0,
entities_fts_failed: 0,
notes_fts_failed: 0,
};
assert!(
report.has_failures(),
"knowledge_sections_failed > 0 alone must drive has_failures() = true"
);
assert!(
decide_result(report.has_failures(), false).is_err(),
"knowledge_sections_failed must fail closed (non-zero exit)"
);
assert!(
decide_result(report.has_failures(), true).is_ok(),
"best-effort downgrades knowledge_sections_failed to exit 0"
);
}
#[test]
fn decide_result_fails_closed_by_default() {
assert!(decide_result(false, false).is_ok(), "clean run exits 0");
assert!(
decide_result(true, false).is_err(),
"failures fail closed (non-zero exit)"
);
}
#[test]
fn decide_result_best_effort_downgrades_to_ok() {
assert!(
decide_result(true, true).is_ok(),
"best-effort downgrades failures to exit 0"
);
assert!(decide_result(false, true).is_ok());
}
#[test]
fn db_memory_sentinel_resolves_to_none() {
assert_eq!(resolve_db_override(Some(":memory:")), Some(None));
}
#[test]
fn db_explicit_path_resolves_to_some() {
assert_eq!(
resolve_db_override(Some("/tmp/kkernel-reindex-test.db")),
Some(Some(PathBuf::from("/tmp/kkernel-reindex-test.db")))
);
}
#[test]
fn db_absent_leaves_default() {
assert_eq!(resolve_db_override(None), None);
}
#[test]
#[serial]
fn khive_db_env_binds_to_db_arg() {
std::env::set_var("KHIVE_DB", "/tmp/kkernel-reindex-env.db");
let args = ReindexArgs::parse_from(["reindex"]);
std::env::remove_var("KHIVE_DB");
assert_eq!(args.db.as_deref(), Some("/tmp/kkernel-reindex-env.db"));
}
#[test]
#[serial]
fn khive_config_env_binds_to_config_arg() {
std::env::set_var("KHIVE_CONFIG", "/tmp/kkernel-reindex.toml");
let args = ReindexArgs::parse_from(["reindex"]);
std::env::remove_var("KHIVE_CONFIG");
assert_eq!(
args.config.as_deref(),
Some(std::path::Path::new("/tmp/kkernel-reindex.toml"))
);
}
#[test]
#[serial]
fn namespace_absent_honors_config_actor_id() {
use std::io::Write;
std::env::remove_var("KHIVE_NAMESPACE");
std::env::remove_var("KHIVE_EMBEDDING_MODEL");
std::env::remove_var("KHIVE_ADDITIONAL_EMBEDDING_MODELS");
let dir = tempfile::tempdir().expect("temp dir");
let config_path = dir.path().join("khive.toml");
let mut f = std::fs::File::create(&config_path).expect("create config");
f.write_all(b"[actor]\nid = \"lambda:prod\"\n")
.expect("write config");
let resolved = resolve_runtime_config(RuntimeConfigInputs {
db: Some(":memory:"),
config: Some(&config_path),
namespace: Namespace::parse("local").expect("ns"),
namespace_explicit: false,
no_embed: false,
packs: None,
brain_profile: None,
})
.expect("resolve config");
assert_eq!(
resolved.default_namespace.as_str(),
"lambda:prod",
"omitted --namespace must defer to config [actor] id"
);
let resolved_explicit = resolve_runtime_config(RuntimeConfigInputs {
db: Some(":memory:"),
config: Some(&config_path),
namespace: Namespace::parse("explicit-ns").expect("ns"),
namespace_explicit: true,
no_embed: false,
packs: None,
brain_profile: None,
})
.expect("resolve config explicit");
assert_eq!(
resolved_explicit.default_namespace.as_str(),
"explicit-ns",
"explicit --namespace must override config [actor] id"
);
}
#[test]
#[serial]
fn namespace_env_var_sets_explicit_flag() {
std::env::set_var("KHIVE_NAMESPACE", "env-ns");
let args = ReindexArgs::parse_from(["reindex"]);
std::env::remove_var("KHIVE_NAMESPACE");
assert_eq!(
args.namespace.as_deref(),
Some("env-ns"),
"KHIVE_NAMESPACE env var must bind to --namespace"
);
assert!(
args.namespace.is_some(),
"env var binding must make namespace Some (explicit)"
);
}
#[test]
fn namespace_absent_defaults_to_none() {
let args = ReindexArgs::parse_from(["reindex"]);
assert!(
args.namespace.is_none(),
"omitted --namespace must be None (not a String default)"
);
}
#[test]
fn has_failures_flags_notes_fts_failed() {
let report = ReindexReport {
entities_processed: 0,
notes_processed: 0,
knowledge_atoms_indexed: None,
knowledge_sections_indexed: None,
knowledge_atoms_failed: 0,
knowledge_pass_errored: false,
knowledge_ann_failed: false,
knowledge_sections_failed: 0,
models_used: vec![],
elapsed_ms: 0,
errors_skipped: 0,
entities_fts_failed: 0,
notes_fts_failed: 1,
};
assert!(
report.has_failures(),
"notes_fts_failed > 0 alone must drive has_failures() = true"
);
assert!(
decide_result(report.has_failures(), false).is_err(),
"notes_fts_failed must fail closed (non-zero exit)"
);
assert!(
decide_result(report.has_failures(), true).is_ok(),
"best-effort downgrades notes_fts_failed to exit 0"
);
}
#[test]
fn note_fts_document_parity_with_name() {
let mut note = Note::new("local", "memory", "the content body");
note.name = Some("my title".to_string());
let doc = note_fts_document(¬e);
assert_eq!(doc.subject_id, note.id);
assert_eq!(doc.namespace, "local");
assert_eq!(doc.title.as_deref(), Some("my title"));
assert_eq!(doc.body, "my title the content body");
assert_eq!(doc.kind, SubstrateKind::Note);
}
#[test]
fn note_fts_document_parity_without_name() {
let note = Note::new("local", "memory", "body only content");
let doc = note_fts_document(¬e);
assert!(doc.title.is_none());
assert_eq!(doc.body, "body only content");
}
#[tokio::test]
async fn fts_backfill_populates_pre_existing_notes() {
use khive_storage::types::TextFilter;
use khive_types::SubstrateKind;
let rt = KhiveRuntime::memory().expect("in-memory runtime");
let ns = Namespace::parse("local").expect("ns");
let token = rt.authorize(ns).expect("authorize");
let notes: Vec<Note> = (0..5)
.map(|i| {
Note::new(
"local",
"memory",
format!("zxqsentinel{i} backfill content"),
)
})
.collect();
let note_store = rt.notes(&token).expect("note store");
for note in ¬es {
note_store
.upsert_note(note.clone())
.await
.expect("upsert note");
}
let fts = rt.text_for_notes(&token).expect("FTS store");
let before = fts
.count(TextFilter {
kinds: vec![SubstrateKind::Note],
namespaces: vec!["local".to_string()],
ids: vec![],
})
.await
.expect("count before");
assert_eq!(before, 0, "FTS must be empty before backfill");
let errors = fts_backfill_notes_batch(&rt, &token, ¬es).await;
assert_eq!(errors, 0, "backfill must produce zero errors");
let after = fts
.count(TextFilter {
kinds: vec![SubstrateKind::Note],
namespaces: vec!["local".to_string()],
ids: vec![],
})
.await
.expect("count after");
assert_eq!(after, 5, "FTS must contain exactly N docs after backfill");
let hits = fts
.search(khive_storage::types::TextSearchRequest {
query: "zxqsentinel0".to_string(),
mode: khive_storage::types::TextQueryMode::Plain,
filter: None,
top_k: 10,
snippet_chars: 0,
})
.await
.expect("FTS search");
assert!(
hits.iter().any(|h| h.subject_id == notes[0].id),
"pre-existing note must be findable by FTS after backfill"
);
}
#[tokio::test]
async fn note_fts_document_matches_runtime_create_path() {
let rt = KhiveRuntime::memory().expect("in-memory runtime");
let ns = Namespace::parse("local").expect("ns");
let token = rt.authorize(ns).expect("authorize");
let props = serde_json::json!({"key": "value", "score": 42});
let note = rt
.create_note(
&token,
"observation",
Some("cross path title"),
"cross path content body",
None,
Some(props),
vec![],
)
.await
.expect("create_note");
let fts = rt.text_for_notes(&token).expect("FTS store");
let stored = fts
.get_document("local", note.id)
.await
.expect("get_document")
.expect("document must exist after create");
let expected = note_fts_document(¬e);
assert_eq!(stored.subject_id, expected.subject_id, "subject_id");
assert_eq!(stored.kind, expected.kind, "kind");
assert_eq!(stored.title, expected.title, "title");
assert_eq!(stored.body, expected.body, "body");
assert_eq!(stored.namespace, expected.namespace, "namespace");
assert_eq!(stored.tags, expected.tags, "tags");
assert_eq!(stored.metadata, expected.metadata, "metadata");
assert_eq!(
stored.updated_at.timestamp_micros(),
note.updated_at,
"updated_at must be derived from the note, not Utc::now()"
);
}
#[tokio::test]
async fn run_reindex_populates_fts_without_embedding_model() {
use khive_storage::types::TextFilter;
use khive_types::SubstrateKind;
let db_file = tempfile::NamedTempFile::new().expect("temp db file");
let db_path = db_file.path().to_str().expect("utf8 path").to_string();
{
let cfg = resolve_runtime_config(RuntimeConfigInputs {
db: Some(&db_path),
config: None,
namespace: Namespace::parse("local").expect("ns"),
namespace_explicit: true,
no_embed: true,
packs: None,
brain_profile: None,
})
.expect("resolve config for seed");
let rt = KhiveRuntime::new(cfg).expect("seed runtime");
let token = rt
.authorize(Namespace::parse("local").expect("ns"))
.expect("authorize");
let note_store = rt.notes(&token).expect("note store");
for i in 0..3usize {
note_store
.upsert_note(Note::new(
"local",
"observation",
format!("run-reindex-sentinel{i} body"),
))
.await
.expect("upsert seed note");
}
}
let args = ReindexArgs {
db: Some(db_path.clone()),
config: None,
model: None,
batch_size: 100,
keep_existing: false,
namespace: Some("local".to_string()),
knowledge_only: false,
no_knowledge: true,
best_effort: true,
no_sections: false,
sections_only: false,
human: false,
};
run_reindex(args).await.expect("run_reindex must succeed");
let cfg = resolve_runtime_config(RuntimeConfigInputs {
db: Some(&db_path),
config: None,
namespace: Namespace::parse("local").expect("ns"),
namespace_explicit: true,
no_embed: true,
packs: None,
brain_profile: None,
})
.expect("resolve config for verify");
let rt = KhiveRuntime::new(cfg).expect("verify runtime");
let token = rt
.authorize(Namespace::parse("local").expect("ns"))
.expect("authorize");
let fts = rt.text_for_notes(&token).expect("FTS store");
let count = fts
.count(TextFilter {
kinds: vec![SubstrateKind::Note],
namespaces: vec!["local".to_string()],
ids: vec![],
})
.await
.expect("fts count");
assert_eq!(
count, 3,
"run_reindex must populate FTS even when no embedding model is configured"
);
}
#[tokio::test]
async fn fts_backfill_runs_without_embedding_model() {
use khive_storage::types::TextFilter;
use khive_types::SubstrateKind;
let rt = KhiveRuntime::memory().expect("in-memory runtime");
let ns = Namespace::parse("local").expect("ns");
let token = rt.authorize(ns).expect("authorize");
let notes: Vec<Note> = (0..3)
.map(|i| {
Note::new(
"local",
"observation",
format!("nomodel-sentinel{i} content"),
)
})
.collect();
let note_store = rt.notes(&token).expect("note store");
for note in ¬es {
note_store.upsert_note(note.clone()).await.expect("upsert");
}
let errors = fts_backfill_notes_batch(&rt, &token, ¬es).await;
assert_eq!(
errors, 0,
"FTS backfill must succeed with no embedding model"
);
let fts = rt.text_for_notes(&token).expect("FTS store");
let count = fts
.count(TextFilter {
kinds: vec![SubstrateKind::Note],
namespaces: vec!["local".to_string()],
ids: vec![],
})
.await
.expect("count");
assert_eq!(
count, 3,
"FTS must be populated even when no embedding model is configured"
);
}
#[tokio::test]
async fn fts_backfill_is_idempotent() {
use khive_storage::types::TextFilter;
use khive_types::SubstrateKind;
let rt = KhiveRuntime::memory().expect("in-memory runtime");
let ns = Namespace::parse("local").expect("ns");
let token = rt.authorize(ns).expect("authorize");
let notes: Vec<Note> = (0..3)
.map(|i| Note::new("local", "memory", format!("idemnote{i} content")))
.collect();
let note_store = rt.notes(&token).expect("note store");
for note in ¬es {
note_store
.upsert_note(note.clone())
.await
.expect("upsert note");
}
let errors1 = fts_backfill_notes_batch(&rt, &token, ¬es).await;
let errors2 = fts_backfill_notes_batch(&rt, &token, ¬es).await;
assert_eq!(errors1, 0);
assert_eq!(errors2, 0);
let fts = rt.text_for_notes(&token).expect("FTS store");
let count = fts
.count(TextFilter {
kinds: vec![SubstrateKind::Note],
namespaces: vec!["local".to_string()],
ids: vec![],
})
.await
.expect("count");
assert_eq!(count, 3, "second backfill pass must not duplicate rows");
}
#[test]
fn entity_fts_document_parity_with_description() {
use khive_storage::entity::Entity;
let mut entity = Entity::new("local", "concept", "TestEntity");
entity = entity.with_description("detail text");
let doc = entity_fts_document(&entity);
assert_eq!(doc.subject_id, entity.id);
assert_eq!(doc.namespace, "local");
assert_eq!(doc.title.as_deref(), Some("TestEntity"));
assert_eq!(doc.body, "TestEntity detail text");
assert_eq!(doc.kind, SubstrateKind::Entity);
}
#[test]
fn entity_fts_document_parity_without_description() {
use khive_storage::entity::Entity;
let entity = Entity::new("local", "concept", "NameOnly");
let doc = entity_fts_document(&entity);
assert_eq!(doc.title.as_deref(), Some("NameOnly"));
assert_eq!(doc.body, "NameOnly");
}
#[tokio::test]
async fn fts_backfill_populates_pre_existing_entities() {
use khive_storage::entity::Entity;
use khive_storage::types::TextFilter;
let rt = KhiveRuntime::memory().expect("in-memory runtime");
let ns = Namespace::parse("local").expect("ns");
let token = rt.authorize(ns).expect("authorize");
let entities: Vec<Entity> = (0..5)
.map(|i| {
Entity::new("local", "concept", format!("zxqentitysentinel{i}"))
.with_description(format!("backfill entity description {i}"))
})
.collect();
let entity_store = rt.entities(&token).expect("entity store");
for entity in &entities {
entity_store
.upsert_entity(entity.clone())
.await
.expect("upsert entity");
}
let fts = rt.text(&token).expect("FTS store");
let before = fts
.count(TextFilter {
kinds: vec![SubstrateKind::Entity],
namespaces: vec!["local".to_string()],
ids: vec![],
})
.await
.expect("count before");
assert_eq!(before, 0, "FTS must be empty before backfill");
let errors = fts_backfill_entities_batch(&rt, &token, &entities).await;
assert_eq!(errors, 0, "backfill must produce zero errors");
let after = fts
.count(TextFilter {
kinds: vec![SubstrateKind::Entity],
namespaces: vec!["local".to_string()],
ids: vec![],
})
.await
.expect("count after");
assert_eq!(after, 5, "FTS must contain exactly N docs after backfill");
let hits = fts
.search(khive_storage::types::TextSearchRequest {
query: "zxqentitysentinel0".to_string(),
mode: khive_storage::types::TextQueryMode::Plain,
filter: None,
top_k: 10,
snippet_chars: 0,
})
.await
.expect("FTS search");
assert!(
hits.iter().any(|h| h.subject_id == entities[0].id),
"pre-existing entity must be findable by FTS after backfill"
);
}
#[tokio::test]
async fn run_reindex_populates_entity_fts_without_embedding_model() {
use khive_storage::entity::Entity;
use khive_storage::types::TextFilter;
let db_file = tempfile::NamedTempFile::new().expect("temp db file");
let db_path = db_file.path().to_str().expect("utf8 path").to_string();
{
let cfg = resolve_runtime_config(RuntimeConfigInputs {
db: Some(&db_path),
config: None,
namespace: Namespace::parse("local").expect("ns"),
namespace_explicit: true,
no_embed: true,
packs: None,
brain_profile: None,
})
.expect("resolve config for seed");
let rt = KhiveRuntime::new(cfg).expect("seed runtime");
let token = rt
.authorize(Namespace::parse("local").expect("ns"))
.expect("authorize");
let entity_store = rt.entities(&token).expect("entity store");
for i in 0..3usize {
entity_store
.upsert_entity(Entity::new(
"local",
"concept",
format!("reindex-entity-sentinel{i}"),
))
.await
.expect("upsert seed entity");
}
}
let args = ReindexArgs {
db: Some(db_path.clone()),
config: None,
model: None,
batch_size: 100,
keep_existing: false,
namespace: Some("local".to_string()),
knowledge_only: false,
no_knowledge: true,
best_effort: true,
no_sections: false,
sections_only: false,
human: false,
};
run_reindex(args).await.expect("run_reindex must succeed");
let cfg = resolve_runtime_config(RuntimeConfigInputs {
db: Some(&db_path),
config: None,
namespace: Namespace::parse("local").expect("ns"),
namespace_explicit: true,
no_embed: true,
packs: None,
brain_profile: None,
})
.expect("resolve config for verify");
let rt = KhiveRuntime::new(cfg).expect("verify runtime");
let token = rt
.authorize(Namespace::parse("local").expect("ns"))
.expect("authorize");
let fts = rt.text(&token).expect("entity FTS store");
let count = fts
.count(TextFilter {
kinds: vec![SubstrateKind::Entity],
namespaces: vec!["local".to_string()],
ids: vec![],
})
.await
.expect("fts count");
assert_eq!(
count, 3,
"run_reindex must populate entity FTS even when no embedding model is configured"
);
}
#[tokio::test]
async fn fts_backfill_entities_is_idempotent() {
use khive_storage::entity::Entity;
use khive_storage::types::TextFilter;
let rt = KhiveRuntime::memory().expect("in-memory runtime");
let ns = Namespace::parse("local").expect("ns");
let token = rt.authorize(ns).expect("authorize");
let entities: Vec<Entity> = (0..3)
.map(|i| Entity::new("local", "concept", format!("idem-entity{i}")))
.collect();
let entity_store = rt.entities(&token).expect("entity store");
for entity in &entities {
entity_store
.upsert_entity(entity.clone())
.await
.expect("upsert entity");
}
let errors1 = fts_backfill_entities_batch(&rt, &token, &entities).await;
let errors2 = fts_backfill_entities_batch(&rt, &token, &entities).await;
assert_eq!(errors1, 0);
assert_eq!(errors2, 0);
let fts = rt.text(&token).expect("FTS store");
let count = fts
.count(TextFilter {
kinds: vec![SubstrateKind::Entity],
namespaces: vec!["local".to_string()],
ids: vec![],
})
.await
.expect("count");
assert_eq!(
count, 3,
"second backfill pass must not duplicate entity rows"
);
}
#[test]
fn has_failures_flags_entities_fts_failed() {
let report = ReindexReport {
entities_processed: 0,
notes_processed: 0,
knowledge_atoms_indexed: None,
knowledge_sections_indexed: None,
knowledge_atoms_failed: 0,
knowledge_pass_errored: false,
knowledge_ann_failed: false,
knowledge_sections_failed: 0,
models_used: vec![],
elapsed_ms: 0,
errors_skipped: 0,
entities_fts_failed: 1,
notes_fts_failed: 0,
};
assert!(
report.has_failures(),
"entities_fts_failed > 0 alone must drive has_failures() = true"
);
assert!(
decide_result(report.has_failures(), false).is_err(),
"entities_fts_failed must fail closed (non-zero exit)"
);
assert!(
decide_result(report.has_failures(), true).is_ok(),
"best-effort downgrades entities_fts_failed to exit 0"
);
}
}