use serde::{Deserialize, Serialize};
use wm_core::{CoreError, Galaxy, Result};
use crate::memory::Memory;
use crate::search::{SearchEngine, sanitize_content_for_index};
use crate::store::MemoryStore;
pub const REINDEX_BATCH: usize = 2_000;
#[derive(Debug, Clone, Default, PartialEq, Eq, Serialize, Deserialize)]
pub struct GalaxyRebuildStats {
pub galaxy: String,
pub scanned: usize,
pub indexed: usize,
pub skipped: usize,
}
#[derive(Debug, Clone, Default, PartialEq, Eq, Serialize, Deserialize)]
pub struct IndexRebuildReport {
pub scanned: usize,
pub indexed: usize,
pub skipped: usize,
pub galaxies: Vec<GalaxyRebuildStats>,
}
#[derive(Debug, Clone, Default, PartialEq, Eq, Serialize, Deserialize)]
pub struct GalaxyConsistency {
pub galaxy: String,
pub lmdb_count: usize,
pub tantivy_count: usize,
pub drift: bool,
}
#[derive(Debug, Clone, Default, PartialEq, Eq, Serialize, Deserialize)]
pub struct ConsistencyReport {
pub galaxies: Vec<GalaxyConsistency>,
pub total_lmdb: usize,
pub total_tantivy: usize,
pub has_drift: bool,
}
#[must_use]
pub fn check_consistency(store: &MemoryStore, search: &SearchEngine) -> ConsistencyReport {
let mut report = ConsistencyReport::default();
for galaxy in Galaxy::memory_galaxies() {
let lmdb_count = store.count(galaxy).unwrap_or(0);
let tantivy_count = search.count_docs_in_galaxy(galaxy.db_name()).unwrap_or(0);
let drift = lmdb_count != tantivy_count;
report.total_lmdb += lmdb_count;
report.total_tantivy += tantivy_count;
if drift {
report.has_drift = true;
}
report.galaxies.push(GalaxyConsistency {
galaxy: galaxy.db_name().to_string(),
lmdb_count,
tantivy_count,
drift,
});
}
report
}
#[derive(Debug, Clone, Default, PartialEq, Eq, Serialize, Deserialize)]
pub struct GalaxyDriftClass {
pub galaxy: String,
pub lmdb_count: usize,
pub tantivy_count: usize,
pub skip_reserve: usize,
pub healable_gap: i64,
}
#[derive(Debug, Clone, Default, PartialEq, Eq, Serialize, Deserialize)]
pub struct DriftClassification {
pub galaxies: Vec<GalaxyDriftClass>,
pub skip_reserve_total: usize,
pub healable_total: usize,
}
#[must_use]
pub fn classify_drift(store: &MemoryStore, search: &SearchEngine) -> DriftClassification {
let mut out = DriftClassification::default();
for galaxy in Galaxy::memory_galaxies() {
let lmdb_count = store.count(galaxy).unwrap_or(0);
let tantivy_count = search.count_docs_in_galaxy(galaxy.db_name()).unwrap_or(0);
let mut skip_reserve = 0usize;
if lmdb_count > tantivy_count {
for mem in store.scan(galaxy, lmdb_count).unwrap_or_default() {
if sanitize_content_for_index(&mem.content).is_none() {
skip_reserve += 1;
}
}
}
let indexable = usize::try_into(lmdb_count - skip_reserve).unwrap_or(i64::MAX);
let indexed = usize::try_into(tantivy_count).unwrap_or(i64::MAX);
let healable_gap = indexable - indexed;
out.skip_reserve_total += skip_reserve;
out.healable_total += healable_gap.unsigned_abs() as usize;
out.galaxies.push(GalaxyDriftClass {
galaxy: galaxy.db_name().to_string(),
lmdb_count,
tantivy_count,
skip_reserve,
healable_gap,
});
}
out
}
pub fn rebuild_index(
store: &MemoryStore,
search: &SearchEngine,
galaxy_filter: &[String],
) -> Result<IndexRebuildReport> {
let mut report = IndexRebuildReport::default();
{
let mut writer = search.writer()?;
if galaxy_filter.is_empty() {
writer
.as_mut()
.ok_or_else(|| {
CoreError::Memory("Tantivy writer unavailable: index opened read-only".into())
})?
.delete_all_documents()
.map_err(|e| CoreError::Memory(format!("Tantivy delete_all_documents: {e}")))?;
} else {
for galaxy in Galaxy::all() {
if galaxy_filter.iter().any(|g| g == galaxy.db_name()) {
search.delete_by_galaxy(&mut writer, galaxy.db_name())?;
}
}
}
for galaxy in Galaxy::all() {
if !galaxy_filter.is_empty() && !galaxy_filter.iter().any(|g| g == galaxy.db_name()) {
continue;
}
let memories = store.scan_all(galaxy)?;
let mut stats = GalaxyRebuildStats {
galaxy: galaxy.db_name().to_string(),
..GalaxyRebuildStats::default()
};
for mem in &memories {
stats.scanned += 1;
if index_memory(search, &mut writer, galaxy, mem)?.is_some() {
stats.indexed += 1;
} else {
stats.skipped += 1;
}
}
report.scanned += stats.scanned;
report.indexed += stats.indexed;
report.skipped += stats.skipped;
report.galaxies.push(stats);
}
search.commit(&mut writer)?;
drop(writer);
}
Ok(report)
}
fn index_memory(
search: &SearchEngine,
writer: &mut Option<tantivy::IndexWriter>,
galaxy: Galaxy,
mem: &Memory,
) -> Result<Option<()>> {
let Some(content) = sanitize_content_for_index(&mem.content) else {
return Ok(None);
};
let timestamp = mem.metadata.created_at.timestamp();
let id = mem.metadata.id.to_string();
search.add_document(
writer,
&id,
galaxy.db_name(),
&content,
&mem.metadata.tags,
timestamp,
)?;
Ok(Some(()))
}
pub fn heal_index_drift(
store: &MemoryStore,
search: &SearchEngine,
) -> Result<Option<IndexRebuildReport>> {
let class = classify_drift(store, search);
let drifted: Vec<String> = class
.galaxies
.iter()
.filter(|g| g.healable_gap != 0)
.map(|g| g.galaxy.clone())
.collect();
if drifted.is_empty() {
return Ok(None);
}
rebuild_index(store, search, &drifted).map(Some)
}
#[derive(Debug, Clone, Default, PartialEq, Eq, Serialize, Deserialize)]
pub struct GalaxyContentRepairStats {
pub galaxy: String,
pub scanned: usize,
pub repaired: usize,
pub unrepairable: usize,
pub already_clean: usize,
}
#[derive(Debug, Clone, Default, PartialEq, Eq, Serialize, Deserialize)]
pub struct ContentRepairReport {
pub scanned: usize,
pub repaired: usize,
pub unrepairable: usize,
pub already_clean: usize,
pub galaxies: Vec<GalaxyContentRepairStats>,
}
fn clean_for_repair(content: &str) -> String {
content
.chars()
.map(|c| {
if c.is_control() && c != '\n' && c != '\t' && c != '\r' {
' '
} else {
c
}
})
.collect()
}
pub fn repair_content(
store: &MemoryStore,
search: &SearchEngine,
galaxies: &[Galaxy],
) -> Result<ContentRepairReport> {
let mut report = ContentRepairReport::default();
let mut writer = search.writer()?;
for galaxy in galaxies {
let mut stats = GalaxyContentRepairStats {
galaxy: galaxy.db_name().to_string(),
..Default::default()
};
for mem in store.scan_all(*galaxy)? {
stats.scanned += 1;
if sanitize_content_for_index(&mem.content).is_some() {
stats.already_clean += 1;
continue;
}
let total = mem.content.chars().count();
let printable = mem.content.chars().filter(|c| !c.is_control()).count();
let cleaned = clean_for_repair(&mem.content);
if total == 0
|| (printable as f32 / total as f32) < 0.5
|| sanitize_content_for_index(&cleaned).is_none()
{
stats.unrepairable += 1;
continue;
}
let mut repaired_mem = mem;
let old_hash = repaired_mem.metadata.content_hash.clone();
repaired_mem.content = cleaned;
repaired_mem.metadata.content_hash = crate::content_hash(&repaired_mem.content);
repaired_mem.metadata.revision_count =
repaired_mem.metadata.revision_count.saturating_add(1);
store.put(*galaxy, &repaired_mem)?;
store.record_revision(
*galaxy,
repaired_mem.metadata.id,
&old_hash,
&repaired_mem.metadata.content_hash,
crate::revision::RevisionActor {
session: None,
user: Some("wm-repair-content".to_string()),
compartment: None,
},
)?;
let id_str = repaired_mem.metadata.id.to_string();
search.delete_document(&mut writer, &id_str)?;
search.add_document(
&mut writer,
&id_str,
galaxy.db_name(),
&repaired_mem.content,
&repaired_mem.metadata.tags,
repaired_mem.metadata.created_at.timestamp(),
)?;
stats.repaired += 1;
}
report.scanned += stats.scanned;
report.repaired += stats.repaired;
report.unrepairable += stats.unrepairable;
report.already_clean += stats.already_clean;
report.galaxies.push(stats);
}
search.commit(&mut writer)?;
Ok(report)
}
#[must_use]
pub fn tantivy_path_for(store_path: &std::path::Path) -> std::path::PathBuf {
store_path.join("tantivy")
}
#[must_use]
pub fn missing_index_error(store_path: &std::path::Path) -> CoreError {
CoreError::Memory(format!(
"Tantivy index not found at {} — run 'wm serve' once to create it",
tantivy_path_for(store_path).display()
))
}
#[cfg(test)]
mod tests {
use super::*;
use crate::Memory;
use tempfile::tempdir;
fn setup() -> (tempfile::TempDir, MemoryStore, SearchEngine) {
let tmp = tempdir().unwrap();
let store = MemoryStore::open_default(tmp.path()).unwrap();
let tantivy_dir = tmp.path().join("tantivy");
std::fs::create_dir_all(&tantivy_dir).unwrap();
let search = SearchEngine::open(&tantivy_dir).unwrap();
(tmp, store, search)
}
fn put_and_index(store: &MemoryStore, search: &SearchEngine, galaxy: Galaxy, content: &str) {
let mem = Memory::new(galaxy, content.to_string());
let id = mem.metadata.id;
store.put(galaxy, &mem).unwrap();
let mut writer = search.writer().unwrap();
search
.add_document(
&mut writer,
&id.to_string(),
galaxy.db_name(),
content,
&mem.metadata.tags,
mem.metadata.created_at.timestamp(),
)
.unwrap();
search.commit(&mut writer).unwrap();
}
#[test]
fn rebuild_repopulates_index_from_lmdb() {
let (_tmp, store, search) = setup();
put_and_index(&store, &search, Galaxy::Codex, "rust memory one");
put_and_index(&store, &search, Galaxy::Codex, "python memory two");
put_and_index(&store, &search, Galaxy::Research, "research notes");
{
let mut writer = search.writer().unwrap();
search
.add_document(
&mut writer,
"99999999-9999-9999-9999-999999999999",
"codex",
"stale ghost document",
&[],
1000,
)
.unwrap();
search.commit(&mut writer).unwrap();
}
let ghost = search.search("ghost", 10).unwrap();
assert_eq!(ghost.len(), 1);
let report = rebuild_index(&store, &search, &[]).unwrap();
assert_eq!(report.indexed, 3);
assert_eq!(report.scanned, 3);
assert_eq!(report.galaxies.len(), Galaxy::COUNT);
let ghost = search.search("ghost", 10).unwrap();
assert!(ghost.is_empty(), "stale index entry must be purged");
let rust = search.search("rust memory one", 10).unwrap();
assert_eq!(rust.len(), 1);
assert_eq!(rust[0].content, "rust memory one");
}
#[test]
fn rebuild_skips_binary_garbage() {
let (_tmp, store, search) = setup();
put_and_index(&store, &search, Galaxy::Codex, "clean text entry");
let mem = Memory::new(Galaxy::Codex, "\u{00}\u{01}\u{02}raw bytes".to_string());
store.put(Galaxy::Codex, &mem).unwrap();
let report = rebuild_index(&store, &search, &[]).unwrap();
assert_eq!(report.indexed, 1, "garbage content must be skipped");
assert_eq!(report.skipped, 1);
let results = search.search("raw", 10).unwrap();
assert!(results.is_empty());
}
#[test]
fn rebuild_respects_galaxy_filter() {
let (_tmp, store, search) = setup();
put_and_index(&store, &search, Galaxy::Codex, "codex memory");
put_and_index(&store, &search, Galaxy::Research, "research memory");
let report = rebuild_index(&store, &search, &["codex".to_string()]).unwrap();
assert_eq!(report.indexed, 1);
assert_eq!(report.galaxies.len(), 1);
assert_eq!(report.galaxies[0].galaxy, "codex");
let codex = search
.search_in_galaxy("codex memory", Some(Galaxy::Codex), 10)
.unwrap();
assert_eq!(codex.len(), 1);
let research = search
.search_in_galaxy("research memory", Some(Galaxy::Research), 10)
.unwrap();
assert_eq!(
research.len(),
1,
"filtered rebuild must preserve documents in unselected galaxies"
);
}
#[test]
fn consistency_check_no_drift_when_indexed() {
let (_tmp, store, search) = setup();
put_and_index(&store, &search, Galaxy::Codex, "hello world");
put_and_index(&store, &search, Galaxy::Codex, "another memory");
let report = check_consistency(&store, &search);
assert!(!report.has_drift, "no drift expected when all indexed");
let codex = report
.galaxies
.iter()
.find(|g| g.galaxy == "codex")
.unwrap();
assert_eq!(codex.lmdb_count, 2);
assert_eq!(codex.tantivy_count, 2);
}
#[test]
fn consistency_check_detects_drift() {
let (_tmp, store, search) = setup();
let mem = Memory::new(Galaxy::Codex, "unindexed".to_string());
store.put(Galaxy::Codex, &mem).unwrap();
let report = check_consistency(&store, &search);
assert!(
report.has_drift,
"drift expected when LMDB has unindexed memory"
);
let codex = report
.galaxies
.iter()
.find(|g| g.galaxy == "codex")
.unwrap();
assert_eq!(codex.lmdb_count, 1);
assert_eq!(codex.tantivy_count, 0);
}
#[test]
fn heal_repairs_only_drifted_galaxies() {
let (_tmp, store, search) = setup();
store
.put(
Galaxy::Sessions,
&Memory::new(Galaxy::Sessions, "session needle".into()),
)
.unwrap();
store
.put(
Galaxy::Research,
&Memory::new(Galaxy::Research, "research needle".into()),
)
.unwrap();
put_and_index(&store, &search, Galaxy::Codex, "healthy codex entry");
let report = heal_index_drift(&store, &search)
.unwrap()
.expect("drift expected before heal");
let healed: Vec<_> = report.galaxies.iter().map(|g| g.galaxy.as_str()).collect();
assert!(healed.contains(&"sessions"));
assert!(healed.contains(&"research"));
assert!(
!healed.contains(&"codex"),
"healthy galaxy must be untouched"
);
assert_eq!(report.indexed, 2);
assert!(
heal_index_drift(&store, &search).unwrap().is_none(),
"second heal must be a no-op once consistent"
);
assert_eq!(
search
.search_in_galaxy("session needle", Some(Galaxy::Sessions), 10)
.unwrap()
.len(),
1
);
assert_eq!(
search
.search_in_galaxy("healthy codex", Some(Galaxy::Codex), 10)
.unwrap()
.len(),
1
);
}
#[test]
fn heal_noop_when_consistent() {
let (_tmp, store, search) = setup();
put_and_index(&store, &search, Galaxy::Codex, "indexed entry");
put_and_index(&store, &search, Galaxy::Dreams, "dream entry");
assert!(heal_index_drift(&store, &search).unwrap().is_none());
}
#[test]
fn index_health_tracks_successes() {
let (_tmp, store, search) = setup();
put_and_index(&store, &search, Galaxy::Codex, "test content");
let health = search.health().snapshot();
let successes = health
.get("successes")
.and_then(serde_json::Value::as_u64)
.unwrap_or(0);
assert!(successes > 0, "expected at least one success");
let failures = health
.get("failures")
.and_then(serde_json::Value::as_u64)
.unwrap_or(0);
assert_eq!(failures, 0);
assert_eq!(
health.get("degraded").and_then(serde_json::Value::as_bool),
Some(false)
);
}
#[test]
fn consistency_check_ignores_non_memory_galaxies() {
let (_tmp, store, search) = setup();
put_and_index(&store, &search, Galaxy::Codex, "indexed memory");
store.put_raw(Galaxy::Karma, b"key1", b"value1").unwrap();
let report = check_consistency(&store, &search);
assert!(
!report.has_drift,
"karma entries should not cause drift — non-memory galaxies are excluded"
);
let galaxy_names: Vec<_> = report.galaxies.iter().map(|g| g.galaxy.as_str()).collect();
assert!(
!galaxy_names.contains(&"karma"),
"karma should not appear in consistency report"
);
assert!(
!galaxy_names.contains(&"dharma"),
"dharma should not appear in consistency report"
);
}
#[test]
fn classify_separates_skip_reserve_from_healable_drift() {
let (tmp, store, search) = setup();
put_and_index(&store, &search, Galaxy::Codex, "clean doc one");
put_and_index(&store, &search, Galaxy::Codex, "clean doc two");
store
.put(
Galaxy::Codex,
&Memory::new(Galaxy::Codex, "bad \u{1}\u{2} doc".into()),
)
.unwrap();
search.commit(&mut search.writer().unwrap()).unwrap();
let class = classify_drift(&store, &search);
let codex = class.galaxies.iter().find(|g| g.galaxy == "codex").unwrap();
assert_eq!(codex.lmdb_count, 3);
assert_eq!(codex.tantivy_count, 2);
assert_eq!(codex.skip_reserve, 1);
assert_eq!(codex.healable_gap, 0, "skip reserve fully explains the gap");
assert_eq!(class.healable_total, 0);
drop(tmp);
}
#[test]
fn classify_flags_missing_indexable_docs_as_healable() {
let (tmp, store, search) = setup();
put_and_index(&store, &search, Galaxy::Codex, "indexed doc");
store
.put(
Galaxy::Codex,
&Memory::new(Galaxy::Codex, "unindexed clean doc".into()),
)
.unwrap();
search.commit(&mut search.writer().unwrap()).unwrap();
let class = classify_drift(&store, &search);
let codex = class.galaxies.iter().find(|g| g.galaxy == "codex").unwrap();
assert_eq!(codex.skip_reserve, 0);
assert_eq!(codex.healable_gap, 1);
assert_eq!(class.healable_total, 1);
drop(tmp);
}
#[test]
fn heal_ignores_pure_skip_reserve_and_heals_real_gaps() {
let (tmp, store, search) = setup();
put_and_index(&store, &search, Galaxy::Codex, "clean doc");
store
.put(
Galaxy::Codex,
&Memory::new(Galaxy::Codex, "gate\u{0} fails".into()),
)
.unwrap();
search.commit(&mut search.writer().unwrap()).unwrap();
let healed = heal_index_drift(&store, &search).unwrap();
assert!(
healed.is_none(),
"skip-reserve-only drift must not trigger a rebuild"
);
store
.put(
Galaxy::Codex,
&Memory::new(Galaxy::Codex, "genuinely missing doc".into()),
)
.unwrap();
let healed = heal_index_drift(&store, &search).unwrap();
assert!(healed.is_some(), "healable drift must trigger a rebuild");
assert_eq!(healed.unwrap().indexed, 2);
drop(tmp);
}
#[test]
fn repair_rewrites_in_place_and_indexes_clean_content() {
let (tmp, store, search) = setup();
let mut repairable = Memory::new(
Galaxy::Codex,
"kumquat\u{0} ratchet \u{1} repair end".into(),
);
let mut binary = Memory::new(Galaxy::Codex, "\u{1}\u{2}\u{3}\u{4}\u{5}\u{6}".into());
let mut clean = Memory::new(Galaxy::Codex, "perfectly fine prose".into());
let (id_r, id_b, id_c) = (
repairable.metadata.id,
binary.metadata.id,
clean.metadata.id,
);
for m in [&mut repairable, &mut binary, &mut clean] {
store.put(Galaxy::Codex, m).unwrap();
}
let report = repair_content(&store, &search, &[Galaxy::Codex]).unwrap();
assert_eq!(report.scanned, 3);
assert_eq!(report.repaired, 1, "{report:?}");
assert_eq!(report.unrepairable, 1, "{report:?}");
assert_eq!(report.already_clean, 1);
let row = store.get(Galaxy::Codex, id_r).unwrap().unwrap();
assert_eq!(row.content, "kumquat ratchet repair end");
assert_eq!(row.metadata.content_hash, crate::content_hash(&row.content));
assert!(sanitize_content_for_index(&row.content).is_some());
assert_eq!(row.metadata.revision_count, 1);
let chain = store.revisions(Galaxy::Codex, id_r).unwrap();
assert_eq!(chain.len(), 1);
assert_eq!(
chain[0].old_hash,
crate::content_hash("kumquat\u{0} ratchet \u{1} repair end")
);
assert_eq!(chain[0].new_hash, row.metadata.content_hash);
assert_eq!(chain[0].actor_user.as_deref(), Some("wm-repair-content"));
assert_eq!(chain[0].actor_session, None);
let verdict = store
.verify_revision_chain(Galaxy::Codex, id_r, &row.metadata.content_hash)
.unwrap();
assert!(verdict.valid, "{:?}", verdict.breaks);
assert!(store.revisions(Galaxy::Codex, id_b).unwrap().is_empty());
assert!(store.revisions(Galaxy::Codex, id_c).unwrap().is_empty());
let untouched = store.get(Galaxy::Codex, id_b).unwrap().unwrap();
assert_eq!(untouched.content, "\u{1}\u{2}\u{3}\u{4}\u{5}\u{6}");
let kept = store.get(Galaxy::Codex, id_c).unwrap().unwrap();
assert_eq!(kept.content, "perfectly fine prose");
let hits = search.search("kumquat ratchet repair", 10).unwrap();
assert!(
hits.iter().any(|h| h.memory_id == id_r.to_string()),
"repaired doc must be indexed: {hits:?}"
);
let again = repair_content(&store, &search, &[Galaxy::Codex]).unwrap();
assert_eq!(again.repaired, 0);
assert_eq!(again.already_clean, 2);
assert_eq!(
store.revisions(Galaxy::Codex, id_r).unwrap().len(),
1,
"idempotent re-run must not append"
);
drop(tmp);
}
}