use std::sync::Arc;
use std::time::Duration;
use crate::AppState;
pub const DEFAULT_REPAIR_INTERVAL_SECS: u64 = 300;
pub const ENV_REPAIR_INTERVAL_SECS: &str = "TRUSTY_BM25_REPAIR_INTERVAL_SECS";
pub type DirtyPalaces = Arc<dashmap::DashSet<String>>;
pub fn repair_interval() -> Option<Duration> {
match std::env::var(ENV_REPAIR_INTERVAL_SECS) {
Ok(raw) => match raw.trim().parse::<u64>() {
Ok(0) => None,
Ok(n) => Some(Duration::from_secs(n)),
Err(_) => {
tracing::warn!(
"{ENV_REPAIR_INTERVAL_SECS}={raw:?} is not an integer — \
using default {DEFAULT_REPAIR_INTERVAL_SECS}s"
);
Some(Duration::from_secs(DEFAULT_REPAIR_INTERVAL_SECS))
}
},
Err(_) => Some(Duration::from_secs(DEFAULT_REPAIR_INTERVAL_SECS)),
}
}
pub fn mark_dirty(state: &AppState, palace: &str) {
if state.bm25_dirty.insert(palace.to_string()) {
tracing::debug!(
palace = %palace,
"bm25: palace queued for coverage repair"
);
}
}
pub fn dirty_palaces(state: &AppState) -> Vec<String> {
state.bm25_dirty.iter().map(|e| e.key().clone()).collect()
}
pub async fn run_repair_pass(state: &AppState) -> (usize, usize) {
let queued: Vec<String> = state.bm25_dirty.iter().map(|e| e.key().clone()).collect();
if queued.is_empty() {
return (0, 0);
}
let root = state.data_root.clone();
let on_disk = match tokio::task::spawn_blocking(move || {
trusty_common::memory_core::registry::PalaceRegistry::list_palaces(&root)
})
.await
{
Ok(Ok(palaces)) => palaces
.into_iter()
.map(|p| p.id.0)
.collect::<std::collections::HashSet<String>>(),
Ok(Err(e)) => {
tracing::warn!("bm25 repair: could not enumerate palaces, deferring pass: {e:#}");
return (0, 0);
}
Err(e) => {
tracing::warn!("bm25 repair: enumeration task failed, deferring pass: {e}");
return (0, 0);
}
};
let mut repaired = 0usize;
for palace in &queued {
if !on_disk.contains(palace) {
state.bm25_dirty.remove(palace);
tracing::debug!(palace = %palace, "bm25 repair: palace no longer on disk — dropping");
continue;
}
let id = trusty_common::memory_core::palace::PalaceId::new(palace.clone());
let registry = Arc::clone(&state.registry);
let root = state.data_root.clone();
let handle = match tokio::task::spawn_blocking(move || registry.open_palace(&root, &id))
.await
{
Ok(Ok(h)) => h,
Ok(Err(e)) => {
tracing::warn!(palace = %palace, "bm25 repair: open failed, staying queued: {e:#}");
continue;
}
Err(e) => {
tracing::warn!(palace = %palace, "bm25 repair: open task failed, staying queued: {e}");
continue;
}
};
let report =
crate::bm25_backfill::backfill_state_palace(state, &handle, palace, false).await;
if report.fully_indexed() {
state.bm25_dirty.remove(palace);
repaired += 1;
tracing::info!(
palace = %palace,
indexed = report.indexed,
"bm25 repair: coverage restored"
);
} else {
tracing::warn!(
palace = %palace,
status = ?report.status,
missing_after = ?report.missing_after,
"bm25 repair: coverage still unverified — staying queued"
);
}
}
(queued.len(), repaired)
}
pub fn spawn_repair_sweep(state: &AppState) {
if state.bm25_client.is_none() {
tracing::debug!("bm25 repair: lane disabled — no repair sweep");
return;
}
let Some(interval) = repair_interval() else {
tracing::info!("bm25 repair: {ENV_REPAIR_INTERVAL_SECS}=0 — repair sweep disabled");
return;
};
let state = state.clone();
tokio::spawn(async move {
let mut ticker = tokio::time::interval(interval);
ticker.tick().await;
loop {
ticker.tick().await;
let (attempted, repaired) = run_repair_pass(&state).await;
if attempted > 0 {
tracing::info!(attempted, repaired, "bm25 repair: pass complete");
}
}
});
tracing::info!(
interval_secs = interval.as_secs(),
"bm25 repair: coverage repair sweep armed"
);
}
#[cfg(test)]
#[path = "bm25_repair_tests.rs"]
mod tests;