use crate::sql::sql;
use std::collections::{BTreeMap, 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::retrieval::EmbeddingTruncationReport;
use khive_runtime::{
entity_embedding_text, entity_fts_document, note_embedding_text, note_fts_document,
KhiveConfig, KhiveRuntime, Namespace,
};
use khive_storage::entity::Entity;
use khive_storage::error::StorageError;
use khive_storage::note::Note;
use khive_storage::types::VectorRecord;
use khive_storage::VectorStore;
use khive_types::SubstrateKind;
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 = "128")]
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, conflicts_with = "no_knowledge")]
pub rebuild_fts: bool,
#[arg(long)]
pub human: bool,
}
pub(crate) fn validate_declared_reindex_target(
db: Option<&str>,
config: Option<&std::path::Path>,
) -> Result<Option<khive_mcp::serve::ValidatedReindexTarget>> {
let discovery_anchor = khive_mcp::serve::config_discovery_db_anchor(db);
let loaded =
KhiveConfig::load_with_home_fallback_and_source(config, discovery_anchor.as_deref())
.context("load reindex khive config for backend-target validation")?;
let config_source = loaded.as_ref().map(|(_, source)| source.as_path());
let backends = loaded
.as_ref()
.map(|(config, _)| config.backends.as_slice())
.unwrap_or_default();
khive_mcp::serve::validate_reindex_db_target_with_source(db, backends, config_source)
}
pub(crate) fn open_validated_reindex_backend(
mut cfg: khive_runtime::RuntimeConfig,
validated: Option<&khive_mcp::serve::ValidatedReindexTarget>,
) -> Result<KhiveRuntime> {
if let Some(validated) = validated {
khive_mcp::serve::reverify_reindex_target_identity(validated)?;
cfg.db_path = Some(validated.path.clone());
}
KhiveRuntime::new(cfg).map_err(|e| anyhow::anyhow!("{e}"))
}
#[derive(Serialize)]
struct KnowledgeFtsRebuildReport {
indexes: Vec<String>,
elapsed_ms: u64,
integrity_ok: 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>,
#[serde(skip_serializing_if = "Option::is_none")]
knowledge_fts_rebuild: Option<KnowledgeFtsRebuildReport>,
knowledge_atoms_failed: u64,
knowledge_pass_errored: bool,
knowledge_ann_failed: bool,
knowledge_sections_failed: u64,
models_used: Vec<String>,
truncation_by_model: BTreeMap<String, EmbeddingTruncationReport>,
elapsed_ms: u64,
errors_skipped: u64,
entities_fts_failed: u64,
notes_fts_failed: u64,
epoch_bump_failed: bool,
vamana_snapshot_invalidation_failed: bool,
}
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
|| self.epoch_bump_failed
|| self.vamana_snapshot_invalidation_failed
}
}
fn entity_has_embedding_text(entity: &Entity) -> bool {
!entity.name.trim().is_empty()
|| entity
.description
.as_deref()
.is_some_and(|description| !description.trim().is_empty())
}
fn note_has_embedding_text(note: &Note) -> bool {
!note.content.trim().is_empty()
}
#[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,
truncation_by_model: &mut BTreeMap<String, EmbeddingTruncationReport>,
) -> 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)| t.clone()).collect();
match rt
.embed_document_batch_with_model_outcomes(model_name, &texts)
.await
{
Ok(outcomes) if outcomes.len() == subset.len() => {
let model_report = truncation_by_model.entry(model_name.clone()).or_default();
for outcome in &outcomes {
model_report.observe(outcome);
}
let expected = subset.len() as u64;
let now = chrono::Utc::now();
let records = subset
.iter()
.zip(outcomes)
.map(|((id, _), outcome)| VectorRecord {
subject_id: *id,
kind,
namespace: namespace.to_string(),
field: field.to_string(),
embedding_model: Some(model_name.clone()),
vectors: vec![outcome.vector],
updated_at: now,
})
.collect();
match vectors.insert_batch(records).await {
Ok(summary)
if summary.attempted == expected
&& summary.affected.saturating_add(summary.failed) == expected =>
{
if summary.failed > 0 {
tracing::warn!(
model = %model_name,
failed = summary.failed,
first_error = %summary.first_error,
"vector batch insert partially failed"
);
errors += summary.failed;
}
}
Ok(summary) => {
tracing::warn!(
model = %model_name,
expected,
attempted = summary.attempted,
affected = summary.affected,
failed = summary.failed,
"vector batch insert returned inconsistent accounting"
);
errors += expected;
}
Err(e) => {
tracing::warn!(model = %model_name, error = %e, "vector batch insert failed");
errors += expected;
}
}
}
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<()> {
run_reindex_with_setup(args, |cfg| cfg, |_| Ok(())).await
}
async fn run_reindex_with_setup(
args: ReindexArgs,
config_setup: impl FnOnce(khive_runtime::RuntimeConfig) -> khive_runtime::RuntimeConfig,
runtime_setup: impl FnOnce(&KhiveRuntime) -> Result<()>,
) -> Result<()> {
let validated_target =
validate_declared_reindex_target(args.db.as_deref(), args.config.as_deref())?;
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,
actor_explicit: false,
no_embed: false,
packs: None,
brain_profile: None,
})?;
let cfg = config_setup(cfg);
let resolved_ns = cfg.default_namespace.clone();
let rt = open_validated_reindex_backend(cfg, validated_target.as_ref())?;
runtime_setup(&rt)?;
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 rebuild_fts = args.rebuild_fts;
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;
let mut truncation_by_model = BTreeMap::new();
let mut epoch_bump_failed = false;
let mut vamana_snapshot_invalidation_failed = false;
if do_graph {
begin_reindex_epoch(&rt)
.await
.context("aborting reindex before any vector mutation")?;
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 embeddable = if model_names.is_empty() {
batch
.iter()
.filter(|entity| entity_has_embedding_text(entity))
.count()
} else {
let mut staged = Vec::with_capacity(n);
for entity in &batch {
if entity_has_embedding_text(entity) {
staged.push((entity.id, entity_embedding_text(entity)));
}
}
if !staged.is_empty() {
errors_skipped += embed_and_store_batch(
&rt,
&token,
&model_names,
&ns_str,
&staged,
SubstrateKind::Entity,
"entity.body",
drop_existing,
&mut truncation_by_model,
)
.await;
}
staged.len()
};
entities_processed += embeddable 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 embeddable = if model_names.is_empty() {
batch
.iter()
.filter(|note| note_has_embedding_text(note))
.count()
} else {
let mut staged = Vec::with_capacity(n);
for note in &batch {
if note_has_embedding_text(note) {
staged.push((note.id, note_embedding_text(note)));
}
}
if !staged.is_empty() {
errors_skipped += embed_and_store_batch(
&rt,
&token,
&model_names,
&ns_str,
&staged,
SubstrateKind::Note,
"note.content",
drop_existing,
&mut truncation_by_model,
)
.await;
}
staged.len()
};
notes_processed += embeddable 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");
vamana_snapshot_invalidation_failed = true;
}
purge_stale_memory_vamana_snapshots(&rt).await;
epoch_bump_failed = !invalidate_active_memory_vamana_snapshot(&rt).await;
sweep_stale_fts_partitions(&rt, &ns_str).await;
}
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;
let mut knowledge_fts_rebuild: Option<KnowledgeFtsRebuildReport> = None;
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 let Some(per_model) = v.get("truncation_by_model").and_then(|v| v.as_object()) {
for (model, value) in per_model {
if let Ok(report) =
serde_json::from_value::<EmbeddingTruncationReport>(value.clone())
{
truncation_by_model
.entry(model.clone())
.or_default()
.merge(report);
}
}
}
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();
}
}
if rebuild_fts && !knowledge_pass_errored {
match khive_pack_knowledge::rebuild_knowledge_fts_indexes(&rt).await {
Ok(fts) => knowledge_fts_rebuild = Some(fts_rebuild_report(&fts)),
Err(e) => {
tracing::error!(error = %e, "knowledge FTS rebuild failed");
eprintln!("\nerror: knowledge FTS rebuild failed: {e}");
knowledge_pass_errored = true;
}
}
}
let elapsed_ms = start.elapsed().as_millis() as u64;
let report = ReindexReport {
entities_processed,
notes_processed,
knowledge_atoms_indexed,
knowledge_sections_indexed,
knowledge_fts_rebuild,
knowledge_atoms_failed,
knowledge_pass_errored,
knowledge_ann_failed,
knowledge_sections_failed,
models_used: model_names,
truncation_by_model,
elapsed_ms,
errors_skipped,
entities_fts_failed,
notes_fts_failed,
epoch_bump_failed,
vamana_snapshot_invalidation_failed,
};
print_report(&report, args.human);
finish(&report, args.best_effort)
}
fn fts_rebuild_report(fts: &serde_json::Value) -> KnowledgeFtsRebuildReport {
KnowledgeFtsRebuildReport {
indexes: fts
.get("indexes")
.and_then(|v| v.as_array())
.map(|a| {
a.iter()
.filter_map(|v| v.as_str().map(str::to_string))
.collect()
})
.unwrap_or_default(),
elapsed_ms: fts.get("elapsed_ms").and_then(|n| n.as_u64()).unwrap_or(0),
integrity_ok: fts
.get("integrity_ok")
.and_then(|b| b.as_bool())
.unwrap_or(false),
}
}
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
}
fn escape_like(input: &str) -> String {
let mut out = String::with_capacity(input.len());
for c in input.chars() {
if matches!(c, '\\' | '%' | '_') {
out.push('\\');
}
out.push(c);
}
out
}
async fn invalidate_vamana_snapshots(rt: &KhiveRuntime, namespace: &str) -> anyhow::Result<()> {
use khive_storage::types::{SqlStatement, SqlValue};
let pattern = format!("{}::vamana::%", escape_like(namespace));
let sql = rt.sql();
let mut writer = sql
.writer()
.await
.context("open SQL writer for Vamana snapshot invalidation")?;
match writer
.execute(SqlStatement {
sql: sql!("retrieval_snapshots_delete_namespace").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 purge_stale_memory_vamana_snapshots(rt: &KhiveRuntime) {
use khive_storage::types::{SqlStatement, SqlValue};
let sql = rt.sql();
let Ok(mut writer) = sql.writer().await else {
return;
};
match writer
.execute(SqlStatement {
sql: sql!("retrieval_snapshots_delete_stale_memory_vamana").into(),
params: vec![SqlValue::Text("global::memory_vamana::*".into())],
label: Some("purge_stale_memory_vamana_snapshots".into()),
})
.await
{
Ok(deleted) => {
if deleted > 0 {
tracing::info!(deleted, "purged stale per-ns memory Vamana snapshot rows");
}
}
Err(e) => {
let msg = e.to_string();
if !msg.contains("no such table") {
tracing::warn!(error = %e, "failed to purge stale memory Vamana snapshots");
}
}
}
}
async fn begin_reindex_epoch(rt: &KhiveRuntime) -> Result<()> {
khive_pack_memory::ensure_ann_epoch_schema(rt)
.await
.map_err(|e| anyhow::anyhow!("{e}"))
.context("failed to ensure memory_ann_epoch schema before reindex")?;
khive_pack_memory::bump_memory_ann_epoch(rt)
.await
.map_err(|e| anyhow::anyhow!("{e}"))
.context("failed to durably mark reindex-in-progress epoch")?;
Ok(())
}
async fn invalidate_active_memory_vamana_snapshot(rt: &KhiveRuntime) -> bool {
use khive_storage::types::{SqlStatement, SqlValue};
let sql = rt.sql();
if let Ok(mut writer) = sql.writer().await {
match writer
.execute(SqlStatement {
sql: sql!("retrieval_snapshots_delete_active_memory_vamana").into(),
params: vec![SqlValue::Text("global::memory_vamana::%".into())],
label: Some("invalidate_active_memory_vamana_snapshot".into()),
})
.await
{
Ok(deleted) => {
if deleted > 0 {
tracing::info!(
deleted,
"invalidated active global memory Vamana snapshot after reindex"
);
}
}
Err(e) => {
let msg = e.to_string();
if !msg.contains("no such table") {
tracing::warn!(error = %e, "failed to invalidate active memory Vamana snapshot");
}
}
}
}
if let Err(e) = khive_pack_memory::bump_memory_ann_epoch(rt).await {
tracing::warn!(error = %e, "failed to bump durable memory ANN epoch after reindex");
return false;
}
true
}
async fn distinct_base_namespaces(rt: &KhiveRuntime) -> HashSet<String> {
use khive_storage::types::SqlStatement;
let sql = rt.sql();
let Ok(mut reader) = sql.reader().await else {
return HashSet::new();
};
let rows = reader
.query_all(SqlStatement {
sql: sql!("base_namespaces_list").into(),
params: vec![],
label: Some("distinct_base_namespaces".into()),
})
.await
.unwrap_or_default();
rows.into_iter()
.filter_map(|row| {
row.get("namespace").and_then(|v| {
if let khive_storage::types::SqlValue::Text(s) = v {
Some(s.clone())
} else {
None
}
})
})
.collect()
}
async fn sweep_stale_fts_partitions(rt: &KhiveRuntime, covered_ns: &str) {
use khive_storage::types::{SqlStatement, SqlValue};
let base_namespaces = distinct_base_namespaces(rt).await;
let uncovered: Vec<&str> = base_namespaces
.iter()
.filter(|ns| ns.as_str() != covered_ns)
.map(String::as_str)
.collect();
if !uncovered.is_empty() {
tracing::warn!(
covered = covered_ns,
uncovered = ?uncovered,
"skipping stale FTS partition sweep: base tables contain namespaces not \
covered by this reindex pass; run reindex for each namespace first, \
or normalize all rows to one namespace before sweeping"
);
return;
}
let canonical: &[&str] = &["fts_entities", "fts_notes", "fts_knowledge", "fts_sections"];
let shadow_suffixes: &[&str] = &["_data", "_idx", "_docsize", "_config", "_content"];
let sql = rt.sql();
let Ok(mut reader) = sql.reader().await else {
return;
};
let rows = reader
.query_all(SqlStatement {
sql: sql!("fts_partitions_list_stale").into(),
params: vec![],
label: Some("sweep_stale_fts_partitions_discover".into()),
})
.await;
let rows = match rows {
Ok(r) => r,
Err(e) => {
tracing::warn!(error = %e, "failed to discover stale FTS partition tables");
return;
}
};
let mut to_drop: Vec<String> = Vec::new();
for row in &rows {
let name = match row.get("name") {
Some(SqlValue::Text(s)) => s.clone(),
_ => continue,
};
if canonical.contains(&name.as_str()) {
continue;
}
if shadow_suffixes.iter().any(|suf| name.ends_with(suf)) {
continue;
}
to_drop.push(name);
}
drop(reader);
if to_drop.is_empty() {
return;
}
let Ok(mut writer) = sql.writer().await else {
return;
};
for table in &to_drop {
let ddl = format!(
concat!("DROP ", "TABLE IF EXISTS {}"),
quote_sqlite_identifier(table)
);
match writer
.execute(SqlStatement {
sql: ddl,
params: vec![],
label: Some("sweep_stale_fts_partitions_drop".into()),
})
.await
{
Ok(_) => {
tracing::info!(table, "dropped stale FTS partition table");
}
Err(e) => {
tracing::warn!(error = %e, table, "failed to drop stale FTS partition table");
}
}
}
}
fn quote_sqlite_identifier(identifier: &str) -> String {
format!("\"{}\"", identifier.replace('"', "\"\""))
}
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: sql!("notes_count").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 render_human_report(report: &ReindexReport) -> String {
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;
let mut output = format!(
"{status}: {} entities, {} notes{}{} ({} vector errors, {} FTS errors) in {}ms\n",
report.entities_processed,
report.notes_processed,
atoms,
sections,
report.errors_skipped,
fts_errors,
report.elapsed_ms
);
if report.entities_fts_failed > 0 {
output.push_str(&format!(
"FTS backfill: {} entity upserts FAILED\n",
report.entities_fts_failed
));
}
if report.notes_fts_failed > 0 {
output.push_str(&format!(
"FTS backfill: {} note upserts FAILED\n",
report.notes_fts_failed
));
}
if report.knowledge_pass_errored {
output.push_str("Knowledge pass: FAILED (did not run to completion)\n");
} else if report.knowledge_atoms_failed > 0 {
output.push_str(&format!(
"Knowledge pass: {} atom vector inserts FAILED\n",
report.knowledge_atoms_failed
));
}
if report.knowledge_sections_failed > 0 {
output.push_str(&format!(
"Knowledge sections: {} section embed/write failures\n",
report.knowledge_sections_failed
));
}
if report.vamana_snapshot_invalidation_failed {
output.push_str(
"Vamana snapshot invalidation: FAILED (snapshots may be stale; prior writes remain committed)\n",
);
}
if report.knowledge_ann_failed {
output.push_str("Knowledge ANN: FAILED (snapshot not rebuilt/persisted)\n");
}
if let Some(fts) = &report.knowledge_fts_rebuild {
output.push_str(&format!(
"Knowledge FTS rebuild: {} in {}ms, integrity {}\n",
fts.indexes.join(", "),
fts.elapsed_ms,
if fts.integrity_ok { "OK" } else { "FAILED" }
));
}
if !report.models_used.is_empty() {
output.push_str(&format!("Models: {}\n", report.models_used.join(", ")));
}
for (model, truncation) in &report.truncation_by_model {
let input_label = if truncation.truncated == 1 {
"input"
} else {
"inputs"
};
output.push_str(&format!(
"Embedding truncation ({model}): {} {input_label} truncated, {} bytes discarded\n",
truncation.truncated, truncation.discarded_bytes
));
}
output
}
fn print_report(report: &ReindexReport, human: bool) {
if human {
print!("{}", render_human_report(report));
} 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;
async fn run_reindex_without_embeddings(args: ReindexArgs) -> Result<()> {
run_reindex_with_setup(
args,
|mut cfg| {
cfg.embedding_model = None;
cfg.additional_embedding_models.clear();
cfg
},
|runtime| {
assert!(
runtime.registered_embedding_model_names().is_empty(),
"FTS-only reindex fixture must have zero configured models"
);
Ok(())
},
)
.await
}
struct FixedReindexEmbedder {
name: String,
dimensions: usize,
}
#[async_trait::async_trait]
impl khive_runtime::EmbedderProvider for FixedReindexEmbedder {
fn name(&self) -> &str {
&self.name
}
fn dimensions(&self) -> usize {
self.dimensions
}
async fn build(
&self,
) -> Result<std::sync::Arc<dyn lattice_embed::EmbeddingService>, khive_runtime::RuntimeError>
{
Ok(std::sync::Arc::new(FixedReindexEmbeddingService(
self.dimensions,
)))
}
}
struct FixedReindexEmbeddingService(usize);
#[async_trait::async_trait]
impl lattice_embed::EmbeddingService for FixedReindexEmbeddingService {
async fn embed(
&self,
texts: &[String],
_model: lattice_embed::EmbeddingModel,
) -> Result<Vec<Vec<f32>>, lattice_embed::EmbedError> {
Ok(texts.iter().map(|_| vec![1.0_f32; self.0]).collect())
}
fn supports_model(&self, _model: lattice_embed::EmbeddingModel) -> bool {
true
}
fn name(&self) -> &'static str {
"reindex-test"
}
}
async fn run_reindex_offline(args: ReindexArgs) -> Result<()> {
run_reindex_with_setup(
args,
|cfg| cfg,
|runtime| {
for name in runtime.registered_embedding_model_names() {
let dimensions = runtime.resolve_embedding_model(Some(&name))?.dimensions();
runtime.register_embedder(FixedReindexEmbedder { name, dimensions });
}
Ok(())
},
)
.await
}
fn write_empty_test_config(dir: &std::path::Path) -> PathBuf {
let path = dir.join("empty-khive-config.toml");
std::fs::write(&path, "").expect("write isolated empty config");
path
}
fn write_declared_backend_test_config(dir: &std::path::Path) -> (PathBuf, PathBuf, PathBuf) {
let main = dir.join("main.db");
let knowledge = dir.join("knowledge.db");
let config = dir.join("khive.toml");
std::fs::write(
&config,
format!(
r#"
[[backends]]
name = "main"
kind = "sqlite"
path = "{}"
[[backends]]
name = "knowledge"
kind = "sqlite"
path = "{}"
"#,
main.display(),
knowledge.display(),
),
)
.expect("write declared-backend config");
(config, main, knowledge)
}
fn write_read_only_backend_test_config(dir: &std::path::Path) -> (PathBuf, PathBuf, PathBuf) {
let main = dir.join("main.db");
let archive = dir.join("archive.db");
let config = dir.join("khive.toml");
std::fs::write(
&config,
format!(
r#"
[[backends]]
name = "main"
kind = "sqlite"
path = "{}"
[[backends]]
name = "archive"
kind = "sqlite"
path = "{}"
read_only = true
"#,
main.display(),
archive.display(),
),
)
.expect("write read-only-backend config");
(config, main, archive)
}
#[test]
fn allocation_free_embedding_eligibility_matches_canonical_text() {
let entities = [
Entity::new("eligibility", "concept", "named"),
Entity::new("eligibility", "concept", "").with_description("description only"),
Entity::new("eligibility", "concept", " ").with_description("\t"),
];
for entity in &entities {
assert_eq!(
entity_has_embedding_text(entity),
!entity_embedding_text(entity).trim().is_empty()
);
}
let notes = [
Note::new("eligibility", "observation", "content"),
Note::new("eligibility", "observation", " \n\t"),
];
for note in ¬es {
assert_eq!(
note_has_embedding_text(note),
!note_embedding_text(note).trim().is_empty()
);
}
}
#[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:?}"
);
}
#[tokio::test]
async fn test_reindex_invalidate_does_not_cross_underscore_namespace() {
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 in &["a_b::vamana::model-a", "aXb::vamana::model-a"] {
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("vamana".to_string()),
SqlValue::Blob(b"{}".to_vec()),
],
label: None,
})
.await
.expect("insert row");
}
drop(w);
invalidate_vamana_snapshots(&rt, "a_b")
.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(&"aXb::vamana::model-a".to_string()),
"unrelated namespace 'aXb' must survive invalidating 'a_b': {remaining:?}"
);
assert!(
!remaining.contains(&"a_b::vamana::model-a".to_string()),
"'a_b' own snapshot must still be deleted: {remaining:?}"
);
}
#[tokio::test]
async fn test_reindex_invalidates_active_memory_vamana_snapshot() {
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 &[
("global::memory_vamana::model-a", "memory_vamana"),
("local::vamana::model-a", "vamana"),
("local::memory_vamana::model-a", "memory_vamana"),
] {
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_active_memory_vamana_snapshot(&rt).await;
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(&"global::memory_vamana::model-a".to_string()),
"the active global memory Vamana snapshot must be deleted: {remaining:?}"
);
assert!(
remaining.contains(&"local::vamana::model-a".to_string()),
"unrelated knowledge Vamana rows must survive: {remaining:?}"
);
assert!(
remaining.contains(&"local::memory_vamana::model-a".to_string()),
"legacy per-namespace memory Vamana rows are purge_stale_memory_vamana_snapshots's \
job, not this function's: {remaining:?}"
);
}
#[tokio::test]
async fn test_purge_stale_memory_vamana_snapshots_keeps_current_key() {
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 &[
("global::memory_vamana::model-a", "memory_vamana"),
("local::memory_vamana::model-a", "memory_vamana"),
("tenant-a::memory_vamana::model-b", "memory_vamana"),
("local::vamana::model-a", "vamana"),
] {
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);
purge_stale_memory_vamana_snapshots(&rt).await;
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(&"global::memory_vamana::model-a".to_string()),
"current-key global memory Vamana snapshot must be retained: {remaining:?}"
);
assert!(
!remaining.contains(&"local::memory_vamana::model-a".to_string()),
"legacy per-namespace memory Vamana snapshot must be purged: {remaining:?}"
);
assert!(
!remaining.contains(&"tenant-a::memory_vamana::model-b".to_string()),
"legacy per-namespace memory Vamana snapshot must be purged: {remaining:?}"
);
assert!(
remaining.contains(&"local::vamana::model-a".to_string()),
"unrelated knowledge Vamana rows must survive: {remaining:?}"
);
}
#[tokio::test]
async fn test_purge_stale_memory_vamana_snapshots_is_case_sensitive() {
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 &[
("global::memory_vamana::model-a", "memory_vamana"),
("GLOBAL::memory_vamana::model-a", "memory_vamana"),
] {
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);
purge_stale_memory_vamana_snapshots(&rt).await;
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(&"global::memory_vamana::model-a".to_string()),
"current-key lowercase global memory Vamana snapshot must be retained: {remaining:?}"
);
assert!(
!remaining.contains(&"GLOBAL::memory_vamana::model-a".to_string()),
"legacy uppercase GLOBAL::memory_vamana snapshot must be purged, not mistaken for \
the retained lowercase key: {remaining:?}"
);
}
#[tokio::test]
async fn stale_fts_sweep_quotes_malicious_table_name_and_preserves_entities() {
let rt = KhiveRuntime::memory().expect("in-memory runtime");
let ns = Namespace::parse("local").expect("ns");
let token = rt.authorize(ns).expect("authorize");
rt.create_entity(&token, "concept", None, "seed", None, None, vec![])
.await
.expect("seed entity");
let sql = rt.sql();
let malicious = "fts_entities_x\"; DROP TABLE entities; --";
{
let mut w = sql.writer().await.expect("writer");
let ddl = format!(
"CREATE TABLE {} (rowid INTEGER)",
quote_sqlite_identifier(malicious)
);
w.execute(SqlStatement {
sql: ddl,
params: vec![],
label: None,
})
.await
.expect("create malicious stale table");
}
sweep_stale_fts_partitions(&rt, "local").await;
let mut r = sql.reader().await.expect("reader");
let rows = r
.query_all(SqlStatement {
sql: "SELECT COUNT(*) AS c FROM entities".into(),
params: vec![],
label: None,
})
.await
.expect("entities table must still exist and be queryable");
let count = rows
.first()
.and_then(|row| row.get("c"))
.map(|v| matches!(v, SqlValue::Integer(n) if *n >= 1))
.unwrap_or(false);
assert!(
count,
"entities table must survive the sweep with its seeded row intact"
);
let survivors = r
.query_all(SqlStatement {
sql: "SELECT name FROM sqlite_master WHERE name = ?1".into(),
params: vec![SqlValue::Text(malicious.to_string())],
label: None,
})
.await
.expect("query sqlite_master");
assert!(
survivors.is_empty(),
"malicious stale table should have been dropped"
);
}
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_fts_rebuild: None,
knowledge_atoms_failed: k_failed,
knowledge_pass_errored: k_errored,
knowledge_ann_failed: false,
knowledge_sections_failed: 0,
models_used: vec![],
truncation_by_model: BTreeMap::new(),
elapsed_ms: 0,
errors_skipped: errors,
entities_fts_failed: 0,
notes_fts_failed: 0,
epoch_bump_failed: false,
vamana_snapshot_invalidation_failed: false,
}
}
#[test]
fn report_serializes_per_model_truncation_accounting() {
let mut report = report_with(0, 0, false);
report.truncation_by_model.insert(
"strict-model".to_string(),
EmbeddingTruncationReport {
truncated: 2,
discarded_bytes: 17,
},
);
let json = serde_json::to_value(report).expect("serialize report");
assert_eq!(json["truncation_by_model"]["strict-model"]["truncated"], 2);
assert_eq!(
json["truncation_by_model"]["strict-model"]["discarded_bytes"],
17
);
}
#[test]
fn human_report_renders_per_model_truncation_in_sorted_order() {
let mut report = report_with(0, 0, false);
report.truncation_by_model.insert(
"zeta-model".to_string(),
EmbeddingTruncationReport {
truncated: 3,
discarded_bytes: 29,
},
);
report.truncation_by_model.insert(
"alpha-model".to_string(),
EmbeddingTruncationReport {
truncated: 1,
discarded_bytes: 7,
},
);
let rendered = render_human_report(&report);
let truncation_lines: Vec<&str> = rendered
.lines()
.filter(|line| line.starts_with("Embedding truncation"))
.collect();
assert_eq!(
truncation_lines,
[
"Embedding truncation (alpha-model): 1 input truncated, 7 bytes discarded",
"Embedding truncation (zeta-model): 3 inputs truncated, 29 bytes discarded",
]
);
}
#[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_fts_rebuild: None,
knowledge_atoms_failed: 0,
knowledge_pass_errored: false,
knowledge_ann_failed: true,
knowledge_sections_failed: 0,
models_used: vec![],
truncation_by_model: BTreeMap::new(),
elapsed_ms: 0,
errors_skipped: 0,
entities_fts_failed: 0,
notes_fts_failed: 0,
epoch_bump_failed: false,
vamana_snapshot_invalidation_failed: false,
};
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_fts_rebuild: None,
knowledge_atoms_failed: 0,
knowledge_pass_errored: false,
knowledge_ann_failed: false,
knowledge_sections_failed: 3,
models_used: vec![],
truncation_by_model: BTreeMap::new(),
elapsed_ms: 0,
errors_skipped: 0,
entities_fts_failed: 0,
notes_fts_failed: 0,
epoch_bump_failed: false,
vamana_snapshot_invalidation_failed: false,
};
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 rebuild_fts_conflicts_with_no_knowledge() {
let err = ReindexArgs::try_parse_from(["reindex", "--rebuild-fts", "--no-knowledge"])
.expect_err("--rebuild-fts with --no-knowledge must be rejected");
assert_eq!(err.kind(), clap::error::ErrorKind::ArgumentConflict);
let ok = ReindexArgs::try_parse_from(["reindex", "--rebuild-fts"])
.expect("--rebuild-fts alone parses");
assert!(ok.rebuild_fts);
}
#[test]
fn rebuild_fts_is_off_unless_requested() {
let default_run = ReindexArgs::try_parse_from(["reindex"]).expect("bare reindex parses");
assert!(!default_run.rebuild_fts);
let keep_existing_run = ReindexArgs::try_parse_from(["reindex", "--keep-existing"])
.expect("keep-existing run parses");
assert!(!keep_existing_run.rebuild_fts);
}
#[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 declared_backends_require_an_explicit_reindex_target() {
let dir = tempfile::tempdir().expect("temp dir");
let (config, _, _) = write_declared_backend_test_config(dir.path());
let error = validate_declared_reindex_target(None, Some(&config))
.expect_err("a topology-backed reindex without a target must fail closed");
let message = error.to_string();
assert!(message.contains("requires an explicit persistent --db / KHIVE_DB target"));
assert!(message.contains(&config.display().to_string()));
}
#[test]
#[serial]
fn declared_secondary_backend_is_a_valid_reindex_target() {
let dir = tempfile::tempdir().expect("temp dir");
let (config, _, knowledge) = write_declared_backend_test_config(dir.path());
validate_declared_reindex_target(knowledge.to_str(), Some(&config))
.expect("reindex may target any explicitly declared SQLite backend");
}
#[test]
#[serial]
fn undeclared_reindex_target_is_rejected_with_config_source() {
let dir = tempfile::tempdir().expect("temp dir");
let (config, _, _) = write_declared_backend_test_config(dir.path());
let wrong = dir.path().join("typo.db");
let error = validate_declared_reindex_target(wrong.to_str(), Some(&config))
.expect_err("an undeclared target must never be reindexed");
let message = error.to_string();
assert!(message.contains("is not a path declared in [[backends]]"));
assert!(message.contains(&config.display().to_string()));
}
#[test]
#[serial]
fn read_only_declared_backend_is_refused_as_reindex_target() {
let dir = tempfile::tempdir().expect("temp dir");
let (config, main, archive) = write_read_only_backend_test_config(dir.path());
let error = validate_declared_reindex_target(archive.to_str(), Some(&config))
.expect_err("a backend declared read_only must never be reindexed");
let message = error.to_string();
assert!(message.contains(&archive.display().to_string()));
assert!(message.contains("read_only"));
validate_declared_reindex_target(main.to_str(), Some(&config))
.expect("a writable declared backend remains a valid reindex target");
let wrong = dir.path().join("typo.db");
validate_declared_reindex_target(wrong.to_str(), Some(&config))
.expect_err("an undeclared target must never be reindexed");
}
#[test]
#[serial]
fn single_backend_reindex_keeps_ordinary_db_override_behavior() {
let dir = tempfile::tempdir().expect("temp dir");
let config = write_empty_test_config(dir.path());
let target = dir.path().join("ordinary.db");
let validated = validate_declared_reindex_target(target.to_str(), Some(&config))
.expect("without [[backends]], --db remains an ordinary target");
assert!(
validated.is_none(),
"no [[backends]] declared: there is no validated identity to bind, so the \
ordinary single-backend --db path must be used unchanged"
);
}
fn resolve_reindex_test_config(
db: Option<&str>,
config: &std::path::Path,
) -> khive_runtime::RuntimeConfig {
resolve_runtime_config(RuntimeConfigInputs {
db,
config: Some(config),
namespace: Namespace::parse("local").expect("ns"),
namespace_explicit: true,
actor_explicit: false,
no_embed: true,
packs: None,
brain_profile: None,
})
.expect("resolve reindex runtime config")
}
#[test]
#[serial]
fn nothing_changed_between_validation_and_open_succeeds() {
let dir = tempfile::tempdir().expect("temp dir");
let a_db = dir.path().join("a.db");
let config = dir.path().join("khive.toml");
std::fs::write(
&config,
format!(
"[[backends]]\nname = \"main\"\nkind = \"sqlite\"\npath = \"{}\"\n",
a_db.display()
),
)
.expect("write config");
let validated = validate_declared_reindex_target(a_db.to_str(), Some(&config))
.expect("a.db is the declared main backend")
.expect("backends declared: a validated target must be returned");
let cfg = resolve_reindex_test_config(a_db.to_str(), &config);
open_validated_reindex_backend(cfg, Some(&validated))
.expect("an unchanged declared target must open normally (control)");
}
#[test]
#[serial]
fn declared_secondary_backend_open_targets_its_own_file_not_main() {
let dir = tempfile::tempdir().expect("temp dir");
let (config, main_db, knowledge_db) = write_declared_backend_test_config(dir.path());
let validated = validate_declared_reindex_target(knowledge_db.to_str(), Some(&config))
.expect("knowledge.db is a declared secondary backend")
.expect("backends declared: a validated target must be returned");
assert_eq!(
validated.path.file_name(),
knowledge_db.file_name(),
"the validated target must be knowledge.db, not main.db"
);
let cfg = resolve_reindex_test_config(knowledge_db.to_str(), &config);
open_validated_reindex_backend(cfg, Some(&validated))
.expect("declared secondary backend must open");
assert!(
knowledge_db.exists(),
"reindex must create/open the targeted secondary backend"
);
assert!(
!main_db.exists(),
"no override normalization on the reindex path may redirect the secondary \
target to main (control)"
);
}
#[test]
#[serial]
#[cfg(unix)]
fn symlink_retargeted_after_validation_does_not_redirect_the_open() {
let dir = tempfile::tempdir().expect("temp dir");
let a_db = dir.path().join("a.db");
let b_db = dir.path().join("b.db");
let link_db = dir.path().join("link.db");
std::fs::write(&a_db, b"").expect("create a.db");
std::os::unix::fs::symlink(&a_db, &link_db).expect("create symlink");
let config = dir.path().join("khive.toml");
std::fs::write(
&config,
format!(
"[[backends]]\nname = \"main\"\nkind = \"sqlite\"\npath = \"{}\"\n",
link_db.display()
),
)
.expect("write config");
let validated = validate_declared_reindex_target(link_db.to_str(), Some(&config))
.expect("link.db resolves through the symlink to the declared backend")
.expect("backends declared: a validated target must be returned");
assert_eq!(
validated.path,
a_db.canonicalize().expect("canonicalize a.db"),
"validation must resolve the symlink to a.db's own canonical path"
);
std::fs::remove_file(&link_db).expect("remove symlink");
std::os::unix::fs::symlink(&b_db, &link_db).expect("retarget symlink");
let pre_fix_cfg = resolve_reindex_test_config(link_db.to_str(), &config);
KhiveRuntime::new(pre_fix_cfg)
.expect("red before the fix: the retargeted symlink opens without complaint");
assert!(
b_db.exists(),
"red before the fix: resolving the raw --db string at open time follows the \
retargeted symlink and creates the undeclared b.db"
);
std::fs::remove_file(&b_db).expect("reset the undeclared file for the fixed path");
let cfg = resolve_reindex_test_config(link_db.to_str(), &config);
open_validated_reindex_backend(cfg, Some(&validated))
.expect("open the declared backend through its validated identity");
assert!(
!b_db.exists(),
"the retargeted undeclared file must never be touched"
);
}
#[test]
#[serial]
fn file_replaced_in_place_after_validation_is_refused_at_open() {
let dir = tempfile::tempdir().expect("temp dir");
let a_db = dir.path().join("a.db");
let other_db = dir.path().join("other.db");
std::fs::write(&a_db, b"declared backend").expect("create a.db");
std::fs::write(&other_db, b"a different file entirely").expect("create other.db");
let config = dir.path().join("khive.toml");
std::fs::write(
&config,
format!(
"[[backends]]\nname = \"main\"\nkind = \"sqlite\"\npath = \"{}\"\n",
a_db.display()
),
)
.expect("write config");
let validated = validate_declared_reindex_target(a_db.to_str(), Some(&config))
.expect("a.db is the declared main backend")
.expect("backends declared: a validated target must be returned");
std::fs::rename(&other_db, &a_db).expect("swap a.db's contents in place");
let cfg = resolve_reindex_test_config(a_db.to_str(), &config);
let error = open_validated_reindex_backend(cfg, Some(&validated))
.map(|_| ())
.expect_err("a file swapped in place after validation must be refused, not opened");
let message = error.to_string();
assert!(message.contains("changed identity between validation and open"));
assert!(message.contains(&a_db.display().to_string()));
}
#[test]
#[serial]
fn khive_db_env_binds_to_db_arg() {
if crate::test_process::run_in_child() {
return;
}
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() {
if crate::test_process::run_in_child() {
return;
}
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_defers_to_local_not_config_actor_id() {
if crate::test_process::run_in_child() {
return;
}
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,
actor_explicit: false,
no_embed: false,
packs: None,
brain_profile: None,
})
.expect("resolve config");
assert_eq!(
resolved.default_namespace.as_str(),
"local",
"omitted --namespace must stay local; config [actor] id does NOT set \
default_namespace (ADR-007 Rev 4 Rule 0)"
);
let resolved_explicit = resolve_runtime_config(RuntimeConfigInputs {
db: Some(":memory:"),
config: Some(&config_path),
namespace: Namespace::parse("explicit-ns").expect("ns"),
namespace_explicit: true,
actor_explicit: false,
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() {
if crate::test_process::run_in_child() {
return;
}
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]
#[serial]
fn namespace_absent_defaults_to_none() {
if crate::test_process::run_in_child() {
return;
}
std::env::remove_var("KHIVE_NAMESPACE");
let args = ReindexArgs::parse_from(["reindex"]);
assert!(
args.namespace.is_none(),
"omitted --namespace must be None (not a String default)"
);
}
#[tokio::test]
async fn embed_and_store_batch_preserves_stale_vector_on_embed_failure() {
use async_trait::async_trait;
use khive_runtime::{EmbedderProvider, RuntimeConfig, RuntimeError};
use lattice_embed::{EmbedError, EmbeddingModel, EmbeddingService};
use std::sync::Arc;
struct FailingStubService;
#[async_trait]
impl EmbeddingService for FailingStubService {
async fn embed(
&self,
_texts: &[String],
_model: EmbeddingModel,
) -> Result<Vec<Vec<f32>>, EmbedError> {
Err(EmbedError::ModelInitialization(
"simulated transient embed failure".into(),
))
}
fn supports_model(&self, _model: EmbeddingModel) -> bool {
true
}
fn name(&self) -> &'static str {
"stub-failing-embed"
}
}
struct StubProvider {
model_name: &'static str,
dims: usize,
}
#[async_trait]
impl EmbedderProvider for StubProvider {
fn name(&self) -> &str {
self.model_name
}
fn dimensions(&self) -> usize {
self.dims
}
async fn build(&self) -> Result<Arc<dyn EmbeddingService>, RuntimeError> {
Ok(Arc::new(FailingStubService))
}
}
const MODEL: &str = "stub-model-embed-fail";
const DIMS: usize = 4;
const NS: &str = "local";
let rt = KhiveRuntime::new(RuntimeConfig {
db_path: None,
embedding_model: None,
additional_embedding_models: vec![],
..RuntimeConfig::default()
})
.expect("runtime");
rt.register_embedder(StubProvider {
model_name: MODEL,
dims: DIMS,
});
let ns = Namespace::parse(NS).expect("ns");
let token = rt.authorize(ns).expect("authorize");
let store = rt.vectors_for_model(&token, MODEL).expect("store");
let subject_id = Uuid::new_v4();
let stale_vec = vec![0.1_f32, 0.2, 0.3, 0.4];
store
.insert_batch(vec![VectorRecord {
subject_id,
kind: SubstrateKind::Note,
namespace: NS.to_string(),
field: "note.content".to_string(),
embedding_model: Some(MODEL.to_string()),
vectors: vec![stale_vec.clone()],
updated_at: chrono::Utc::now(),
}])
.await
.expect("stale insert_batch");
assert_eq!(store.count().await.expect("count before"), 1);
let staged = vec![(subject_id, "some note content".to_string())];
let errors = embed_and_store_batch(
&rt,
&token,
&[MODEL.to_string()],
NS,
&staged,
SubstrateKind::Note,
"note.content",
true,
&mut BTreeMap::new(),
)
.await;
assert_eq!(
errors, 1,
"the forced embed failure must count as one error"
);
let after = store.count().await.expect("count after");
assert_eq!(
after, 1,
"an embed failure must leave the prior vector in place, not absent"
);
assert!(
store
.batch_exists(&[subject_id], NS)
.await
.expect("batch_exists after failure")
.contains(&subject_id),
"stale subject must still resolve to a row after the embed failure"
);
let hits = store
.search(khive_storage::types::VectorSearchRequest {
query_vectors: vec![stale_vec],
top_k: 1,
namespace: Some(NS.to_string()),
kind: Some(SubstrateKind::Note),
embedding_model: None,
filter: None,
backend_hints: None,
})
.await
.expect("search after failure");
assert_eq!(hits.len(), 1, "stale vector must still be searchable");
assert_eq!(hits[0].subject_id, subject_id);
assert!(
hits[0].score.to_f64() > 0.999,
"surviving row must be the original stale vector, not a partial write"
);
}
#[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_fts_rebuild: None,
knowledge_atoms_failed: 0,
knowledge_pass_errored: false,
knowledge_ann_failed: false,
knowledge_sections_failed: 0,
models_used: vec![],
truncation_by_model: BTreeMap::new(),
elapsed_ms: 0,
errors_skipped: 0,
entities_fts_failed: 0,
notes_fts_failed: 1,
epoch_bump_failed: false,
vamana_snapshot_invalidation_failed: false,
};
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],
record_kinds: vec![],
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],
record_kinds: vec![],
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() {
if crate::test_process::run_in_child() {
return;
}
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 config_dir = tempfile::tempdir().expect("config temp dir");
let config = write_empty_test_config(config_dir.path());
{
let cfg = resolve_runtime_config(RuntimeConfigInputs {
db: Some(&db_path),
config: Some(&config),
namespace: Namespace::parse("local").expect("ns"),
namespace_explicit: true,
actor_explicit: false,
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: Some(config.clone()),
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,
rebuild_fts: false,
human: false,
};
run_reindex_without_embeddings(args)
.await
.expect("run_reindex must succeed");
let cfg = resolve_runtime_config(RuntimeConfigInputs {
db: Some(&db_path),
config: Some(&config),
namespace: Namespace::parse("local").expect("ns"),
namespace_explicit: true,
actor_explicit: false,
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],
record_kinds: vec![],
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],
record_kinds: vec![],
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],
record_kinds: vec![],
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],
record_kinds: vec![],
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],
record_kinds: vec![],
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() {
if crate::test_process::run_in_child() {
return;
}
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 config_dir = tempfile::tempdir().expect("config temp dir");
let config = write_empty_test_config(config_dir.path());
{
let cfg = resolve_runtime_config(RuntimeConfigInputs {
db: Some(&db_path),
config: Some(&config),
namespace: Namespace::parse("local").expect("ns"),
namespace_explicit: true,
actor_explicit: false,
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: Some(config.clone()),
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,
rebuild_fts: false,
human: false,
};
run_reindex_without_embeddings(args)
.await
.expect("run_reindex must succeed");
let cfg = resolve_runtime_config(RuntimeConfigInputs {
db: Some(&db_path),
config: Some(&config),
namespace: Namespace::parse("local").expect("ns"),
namespace_explicit: true,
actor_explicit: false,
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],
record_kinds: vec![],
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"
);
}
async fn seed_desynced_knowledge_fts(db_path: &str, config: &std::path::Path) {
let cfg = resolve_runtime_config(RuntimeConfigInputs {
db: Some(db_path),
config: Some(config),
namespace: Namespace::parse("local").expect("ns"),
namespace_explicit: true,
actor_explicit: false,
no_embed: true,
packs: None,
brain_profile: None,
})
.expect("resolve config for seed");
let rt = KhiveRuntime::new(cfg).expect("seed runtime");
let mut writer = rt.sql().writer().await.expect("knowledge writer");
writer
.execute_batch(vec![
SqlStatement {
sql: "INSERT INTO knowledge_atoms \
(id, namespace, slug, name, content, created_at, updated_at) \
VALUES ('9de50000-0000-4000-8000-000000000001', 'local', \
'reindex-fts-scope', 'Reindex FTS Scope', \
'scopeable lexical atom document', 1, 1)"
.into(),
params: vec![],
label: Some("test.reindex_fts_scope.atom".into()),
},
SqlStatement {
sql: "INSERT INTO fts_knowledge \
(fts_knowledge, rowid, id, namespace, slug, name, content) \
SELECT 'delete', rowid, id, namespace, slug, name, content \
FROM knowledge_atoms \
WHERE id = '9de50000-0000-4000-8000-000000000001'"
.into(),
params: vec![],
label: Some("test.reindex_fts_scope.desync".into()),
},
])
.await
.expect("seed and desynchronize fts_knowledge");
}
async fn knowledge_fts_repaired(db_path: &str, config: &std::path::Path) -> bool {
let cfg = resolve_runtime_config(RuntimeConfigInputs {
db: Some(db_path),
config: Some(config),
namespace: Namespace::parse("local").expect("ns"),
namespace_explicit: true,
actor_explicit: false,
no_embed: true,
packs: None,
brain_profile: None,
})
.expect("resolve config for verify");
let rt = KhiveRuntime::new(cfg).expect("verify runtime");
let mut reader = rt.sql().reader().await.expect("knowledge reader");
let row = reader
.query_row(SqlStatement {
sql: "SELECT count(*) AS n FROM fts_knowledge \
WHERE fts_knowledge MATCH 'scopeable'"
.into(),
params: vec![],
label: Some("test.reindex_fts_scope.verify".into()),
})
.await
.expect("query fts_knowledge")
.expect("count row");
matches!(row.get("n"), Some(SqlValue::Integer(1)))
}
#[tokio::test]
async fn run_reindex_scoped_run_does_not_rebuild_fts() {
if crate::test_process::run_in_child() {
return;
}
let db_file = tempfile::NamedTempFile::new().expect("temp db file");
let db_path = db_file.path().to_str().expect("utf8 path").to_string();
let config_dir = tempfile::tempdir().expect("config temp dir");
let config = write_empty_test_config(config_dir.path());
seed_desynced_knowledge_fts(&db_path, &config).await;
let args = ReindexArgs {
db: Some(db_path.clone()),
config: Some(config.clone()),
model: None,
batch_size: 100,
keep_existing: false,
namespace: Some("local".to_string()), knowledge_only: false,
no_knowledge: false,
best_effort: true,
no_sections: false,
sections_only: false,
rebuild_fts: false,
human: false,
};
run_reindex_offline(args)
.await
.expect("run_reindex must succeed");
assert!(
!knowledge_fts_repaired(&db_path, &config).await,
"a namespace-scoped run must NOT rebuild the global knowledge FTS indexes"
);
}
#[tokio::test]
async fn run_reindex_without_the_flag_does_not_rebuild_fts() {
if crate::test_process::run_in_child() {
return;
}
let db_file = tempfile::NamedTempFile::new().expect("temp db file");
let db_path = db_file.path().to_str().expect("utf8 path").to_string();
let config_dir = tempfile::tempdir().expect("config temp dir");
let config = write_empty_test_config(config_dir.path());
seed_desynced_knowledge_fts(&db_path, &config).await;
let args = ReindexArgs {
db: Some(db_path.clone()),
config: Some(config.clone()),
model: None,
batch_size: 100,
keep_existing: false,
namespace: None, knowledge_only: false,
no_knowledge: false,
best_effort: true,
no_sections: false,
sections_only: false,
rebuild_fts: false,
human: false,
};
run_reindex_offline(args)
.await
.expect("run_reindex must succeed");
assert!(
!knowledge_fts_repaired(&db_path, &config).await,
"a run without --rebuild-fts must not rebuild the global knowledge FTS indexes"
);
}
#[tokio::test]
async fn run_reindex_with_the_flag_rebuilds_fts() {
if crate::test_process::run_in_child() {
return;
}
let db_file = tempfile::NamedTempFile::new().expect("temp db file");
let db_path = db_file.path().to_str().expect("utf8 path").to_string();
let config_dir = tempfile::tempdir().expect("config temp dir");
let config = write_empty_test_config(config_dir.path());
seed_desynced_knowledge_fts(&db_path, &config).await;
let args = ReindexArgs {
db: Some(db_path.clone()),
config: Some(config.clone()),
model: None,
batch_size: 100,
keep_existing: false,
namespace: None,
knowledge_only: false,
no_knowledge: false,
best_effort: true,
no_sections: false,
sections_only: false,
rebuild_fts: true,
human: false,
};
run_reindex_offline(args)
.await
.expect("run_reindex must succeed");
assert!(
knowledge_fts_repaired(&db_path, &config).await,
"--rebuild-fts must rebuild and repair the global knowledge FTS indexes"
);
}
#[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],
record_kinds: vec![],
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_fts_rebuild: None,
knowledge_atoms_failed: 0,
knowledge_pass_errored: false,
knowledge_ann_failed: false,
knowledge_sections_failed: 0,
models_used: vec![],
truncation_by_model: BTreeMap::new(),
elapsed_ms: 0,
errors_skipped: 0,
entities_fts_failed: 1,
notes_fts_failed: 0,
epoch_bump_failed: false,
vamana_snapshot_invalidation_failed: false,
};
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"
);
}
#[test]
fn snapshot_invalidation_failure_is_reported_and_controls_exit() {
let mut report = report_with(0, 0, false);
assert!(finish(&report, false).is_ok());
assert_eq!(
serde_json::to_value(&report).unwrap()["vamana_snapshot_invalidation_failed"],
false
);
report.vamana_snapshot_invalidation_failed = true;
assert!(finish(&report, false).is_err());
assert!(finish(&report, true).is_ok());
assert_eq!(
serde_json::to_value(&report).unwrap()["vamana_snapshot_invalidation_failed"],
true
);
let human = render_human_report(&report);
assert!(human.contains("Reindex completed WITH FAILURES"));
assert!(human.contains("Vamana snapshot invalidation: FAILED"));
assert!(human.contains("prior writes remain committed"));
}
fn snapshot_reindex_args(dir: &std::path::Path, best_effort: bool) -> ReindexArgs {
ReindexArgs {
db: Some(dir.join("reindex.db").to_str().unwrap().to_owned()),
config: Some(write_empty_test_config(dir)),
model: None,
batch_size: 100,
keep_existing: false,
namespace: Some("local".into()),
knowledge_only: false,
no_knowledge: true,
best_effort,
no_sections: false,
sections_only: false,
rebuild_fts: false,
human: false,
}
}
fn snapshot_test_runtime(args: &ReindexArgs) -> KhiveRuntime {
let cfg = resolve_runtime_config(RuntimeConfigInputs {
db: args.db.as_deref(),
config: args.config.as_deref(),
namespace: Namespace::local(),
namespace_explicit: true,
actor_explicit: false,
no_embed: true,
packs: None,
brain_profile: None,
})
.expect("resolve owned test database");
KhiveRuntime::new(cfg).expect("owned test runtime")
}
#[tokio::test]
async fn reindex_fts_fixture_clears_primary_and_additional_models() {
if crate::test_process::run_in_child() {
return;
}
let dir = tempfile::tempdir().expect("reindex fixture");
let args = snapshot_reindex_args(dir.path(), false);
std::fs::write(
args.config.as_ref().unwrap(),
r#"
[[engines]]
name = "primary"
model = "bge-small-en-v1.5"
default = true
[[engines]]
name = "additional"
model = "paraphrase"
default = false
"#,
)
.expect("write two-engine fixture");
let configured = resolve_runtime_config(RuntimeConfigInputs {
db: args.db.as_deref(),
config: args.config.as_deref(),
namespace: Namespace::local(),
namespace_explicit: true,
actor_explicit: false,
no_embed: false,
packs: None,
brain_profile: None,
})
.expect("resolve fixture engine precondition");
assert!(configured.embedding_model.is_some());
assert_eq!(configured.additional_embedding_models.len(), 1);
run_reindex_without_embeddings(args)
.await
.expect("FTS-only fixture must clear both configured engines");
}
async fn snapshot_test_count(rt: &KhiveRuntime, query: &str) -> i64 {
let mut reader = rt.sql().reader().await.expect("reader");
let row = reader
.query_row(SqlStatement {
sql: query.into(),
params: vec![],
label: Some("test.reindex_snapshot.count".into()),
})
.await
.expect("query")
.expect("count row");
match row.get("n") {
Some(SqlValue::Integer(n)) => *n,
other => panic!("expected integer count, got {other:?}"),
}
}
async fn seed_snapshot_test_note(rt: &KhiveRuntime) {
let token = rt.authorize(Namespace::local()).expect("authorize");
rt.notes(&token)
.expect("notes")
.upsert_note(Note::new("local", "observation", ""))
.await
.expect("seed note without indexing");
assert_eq!(
snapshot_test_count(rt, "SELECT count(*) AS n FROM fts_notes").await,
0
);
}
#[tokio::test]
async fn run_reindex_snapshot_failure_fails_closed_and_preserves_committed_work() {
if crate::test_process::run_in_child() {
return;
}
let dir = tempfile::tempdir().expect("owned database directory");
let args = snapshot_reindex_args(dir.path(), false);
let rt = snapshot_test_runtime(&args);
seed_snapshot_test_note(&rt).await;
{
let mut writer = rt.sql().writer().await.expect("writer");
writer
.execute_script(
"CREATE TABLE 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));
INSERT INTO retrieval_snapshots VALUES
('local::vamana::test-model', 'vamana', X'00', 0),
('other::vamana::test-model', 'vamana', X'00', 0),
('local::hnsw::test-model', 'hnsw', X'00', 0),
('global::memory_vamana::test-model', 'memory_vamana', X'00', 0);
CREATE TRIGGER refuse_namespace_snapshot_delete
BEFORE DELETE ON retrieval_snapshots
WHEN OLD.namespace = 'local::vamana::test-model'
BEGIN SELECT RAISE(ABORT, 'injected namespace snapshot failure'); END;"
.into(),
)
.await
.expect("seed snapshots and rejecting trigger");
}
let injected = invalidate_vamana_snapshots(&rt, "local")
.await
.expect_err("the owned trigger must reject the real invalidation statement");
assert!(injected
.to_string()
.contains("injected namespace snapshot failure"));
let error = run_reindex_offline(args)
.await
.expect_err("snapshot failure must fail the command");
assert!(error
.to_string()
.contains("reindex completed with failures"));
assert_eq!(
snapshot_test_count(&rt, "SELECT count(*) AS n FROM fts_notes").await,
1,
"the earlier FTS write remains committed"
);
assert_eq!(
snapshot_test_count(&rt, "SELECT epoch AS n FROM memory_ann_epoch").await,
2,
"snapshot failure must not suppress the completion epoch"
);
assert_eq!(
snapshot_test_count(&rt, "SELECT count(*) AS n FROM retrieval_snapshots").await,
3,
"only the independent active-memory snapshot was removed"
);
run_reindex_offline(snapshot_reindex_args(dir.path(), true))
.await
.expect("explicit best effort allows partial completion");
assert_eq!(
snapshot_test_count(&rt, "SELECT count(*) AS n FROM retrieval_snapshots").await,
3,
"best effort must not bypass the injected DELETE refusal"
);
{
let mut writer = rt.sql().writer().await.expect("writer");
writer
.execute_script("DROP TRIGGER refuse_namespace_snapshot_delete;".into())
.await
.expect("remove injected failure");
}
run_reindex_offline(snapshot_reindex_args(dir.path(), false))
.await
.expect("same fixture succeeds once invalidation can commit");
assert_eq!(
snapshot_test_count(&rt, "SELECT count(*) AS n FROM retrieval_snapshots").await,
2,
"unrelated namespace and HNSW snapshots survive"
);
assert_eq!(
snapshot_test_count(
&rt,
"SELECT count(*) AS n FROM retrieval_snapshots WHERE namespace = 'local::vamana::test-model'"
)
.await,
0
);
}
#[tokio::test]
async fn run_reindex_snapshot_missing_table_remains_successful() {
if crate::test_process::run_in_child() {
return;
}
let dir = tempfile::tempdir().expect("owned database directory");
let args = snapshot_reindex_args(dir.path(), false);
let rt = snapshot_test_runtime(&args);
seed_snapshot_test_note(&rt).await;
let table_count = "SELECT count(*) AS n FROM sqlite_master WHERE type = 'table' AND name = 'retrieval_snapshots'";
assert_eq!(snapshot_test_count(&rt, table_count).await, 0);
run_reindex_offline(args)
.await
.expect("missing snapshots are a successful no-op");
assert_eq!(snapshot_test_count(&rt, table_count).await, 0);
assert_eq!(
snapshot_test_count(&rt, "SELECT count(*) AS n FROM fts_notes").await,
1
);
assert_eq!(
snapshot_test_count(&rt, "SELECT epoch AS n FROM memory_ann_epoch").await,
2
);
}
}