use std::time::{Duration, Instant};
use trusty_common::memory_core::palace::Drawer;
use trusty_common::memory_core::retrieval::PalaceHandle;
use crate::bm25_lane::Bm25Lane;
use crate::AppState;
const PALACE_BUDGET: Duration = Duration::from_secs(120);
pub const ENV_NO_BACKFILL: &str = "TRUSTY_BM25_NO_BACKFILL";
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum BackfillStatus {
Disabled,
IndexUnavailable,
AlreadyIndexed,
Completed,
Partial,
}
#[derive(Debug, Clone)]
pub struct BackfillReport {
pub palace: String,
pub status: BackfillStatus,
pub drawers_total: usize,
pub skipped_empty: usize,
pub indexed: usize,
pub failed: usize,
pub missing_after: Option<usize>,
pub final_doc_count: Option<usize>,
pub elapsed_ms: u64,
}
impl BackfillReport {
fn short_circuit(palace: &str, status: BackfillStatus, drawers_total: usize) -> Self {
Self {
palace: palace.to_string(),
status,
drawers_total,
skipped_empty: 0,
indexed: 0,
failed: 0,
missing_after: None,
final_doc_count: None,
elapsed_ms: 0,
}
}
pub fn fully_indexed(&self) -> bool {
self.missing_after == Some(0)
}
pub fn stale_doc_estimate(&self) -> Option<usize> {
let indexable = self.drawers_total.saturating_sub(self.skipped_empty);
self.final_doc_count.map(|n| n.saturating_sub(indexable))
}
}
#[derive(Debug, Clone, Default)]
pub struct PalaceDocs {
pub docs: Vec<(String, String)>,
pub skipped_empty: usize,
}
impl PalaceDocs {
pub fn drawers_total(&self) -> usize {
self.docs.len() + self.skipped_empty
}
pub fn from_pairs(docs: Vec<(String, String)>) -> Self {
Self {
docs,
skipped_empty: 0,
}
}
}
pub fn palace_docs(handle: &PalaceHandle) -> PalaceDocs {
let drawers = handle.drawers.read();
docs_from_drawers(&drawers)
}
pub fn docs_from_drawers(drawers: &[Drawer]) -> PalaceDocs {
let mut docs = Vec::with_capacity(drawers.len());
let mut skipped_empty = 0usize;
for d in drawers {
if d.content().trim().is_empty() {
skipped_empty += 1;
} else {
docs.push((d.id.to_string(), d.content().to_string()));
}
}
PalaceDocs {
docs,
skipped_empty,
}
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
enum Coverage {
Missing(usize),
Unreachable,
}
async fn probe_coverage(lane: &Bm25Lane, palace: &str, ids: &[String]) -> Coverage {
match lane.missing_docs(palace, ids).await {
Ok(cov) => Coverage::Missing(cov.missing.len()),
Err(e) => {
tracing::warn!(palace = %palace, "bm25 backfill: coverage probe failed: {e:#}");
Coverage::Unreachable
}
}
}
pub async fn backfill_palace(
lane: &Bm25Lane,
palace: &str,
palace_docs: PalaceDocs,
force: bool,
) -> BackfillReport {
let started = Instant::now();
let PalaceDocs {
docs,
skipped_empty,
} = palace_docs;
let drawers_total = docs.len() + skipped_empty;
let total = docs.len();
if total == 0 {
let mut report =
BackfillReport::short_circuit(palace, BackfillStatus::AlreadyIndexed, drawers_total);
report.skipped_empty = skipped_empty;
report.missing_after = Some(0);
return report;
}
let ids: Vec<String> = docs.iter().map(|(id, _)| id.clone()).collect();
if !force {
match probe_coverage(lane, palace, &ids).await {
Coverage::Missing(0) => {
tracing::debug!(
palace = %palace,
drawers = total,
"bm25 backfill: every drawer id already present — skipping"
);
let mut report = BackfillReport::short_circuit(
palace,
BackfillStatus::AlreadyIndexed,
drawers_total,
);
report.skipped_empty = skipped_empty;
report.missing_after = Some(0);
report.final_doc_count = read_doc_count(lane, palace).await;
report.elapsed_ms = started.elapsed().as_millis() as u64;
log_stale_docs(&report);
return report;
}
Coverage::Missing(n) => tracing::info!(
palace = %palace,
missing = n,
drawers = total,
"bm25 backfill: palace under-indexed — running"
),
Coverage::Unreachable => {
let mut report = BackfillReport::short_circuit(
palace,
BackfillStatus::IndexUnavailable,
drawers_total,
);
report.skipped_empty = skipped_empty;
return report;
}
}
}
let deadline = started + PALACE_BUDGET;
let mut indexed = 0usize;
let mut failed = 0usize;
let mut truncated = false;
for (doc_id, text) in &docs {
if Instant::now() >= deadline {
tracing::warn!(
palace = %palace,
indexed,
remaining = total - indexed - failed,
"bm25 backfill: time budget expired — reporting partial coverage"
);
truncated = true;
break;
}
match lane.index(palace, doc_id, text).await {
Ok(()) => indexed += 1,
Err(e) => {
failed += 1;
tracing::warn!(palace = %palace, doc_id = %doc_id, "bm25 backfill index failed: {e:#}");
}
}
}
let missing_after = match probe_coverage(lane, palace, &ids).await {
Coverage::Missing(n) => Some(n),
Coverage::Unreachable => None,
};
if let Err(e) = lane.flush(palace).await {
tracing::warn!(palace = %palace, "bm25 backfill: snapshot flush failed: {e:#}");
}
let status = if missing_after == Some(0) && failed == 0 && !truncated {
BackfillStatus::Completed
} else {
BackfillStatus::Partial
};
let report = BackfillReport {
palace: palace.to_string(),
status,
drawers_total,
skipped_empty,
indexed,
failed,
missing_after,
final_doc_count: read_doc_count(lane, palace).await,
elapsed_ms: started.elapsed().as_millis() as u64,
};
if !report.fully_indexed() {
tracing::error!(
palace = %palace,
?status,
indexed,
failed,
missing_after = ?missing_after,
"bm25 backfill did NOT establish coverage — this palace answers lexical \
queries from a partial corpus"
);
} else {
tracing::info!(
palace = %palace,
indexed,
elapsed_ms = report.elapsed_ms,
"bm25 backfill finished — coverage verified by drawer id"
);
}
log_stale_docs(&report);
report
}
async fn read_doc_count(lane: &Bm25Lane, palace: &str) -> Option<usize> {
match lane.stats(palace).await {
Ok(stats) => Some(stats.doc_count),
Err(e) => {
tracing::debug!(palace = %palace, "bm25 backfill: stats read failed: {e:#}");
None
}
}
}
fn log_stale_docs(report: &BackfillReport) {
if let Some(stale) = report.stale_doc_estimate().filter(|n| *n > 0) {
tracing::warn!(
palace = %report.palace,
stale,
doc_count = ?report.final_doc_count,
"bm25 index holds documents for drawers this palace no longer has — \
stale lexical hits are possible (see #5053)"
);
}
}
pub async fn backfill_state_palace(
state: &AppState,
handle: &PalaceHandle,
palace: &str,
force: bool,
) -> BackfillReport {
let Some(lane) = state.bm25.as_ref() else {
return BackfillReport::short_circuit(palace, BackfillStatus::Disabled, 0);
};
let docs = palace_docs(handle);
backfill_palace(lane, palace, docs, force).await
}
pub fn startup_backfill_opted_out() -> bool {
std::env::var(ENV_NO_BACKFILL).as_deref() == Ok("1")
}
#[derive(Debug, Clone, Copy, Default, PartialEq, Eq)]
pub struct SweepOutcome {
pub enumerated: Option<usize>,
pub swept: usize,
pub incomplete: usize,
pub unopenable: usize,
}
impl SweepOutcome {
pub fn all_verified(&self) -> bool {
self.enumerated.is_some() && self.incomplete == 0 && self.unopenable == 0
}
}
fn palace_ids_on_disk(data_root: &std::path::Path) -> anyhow::Result<Vec<String>> {
let mut ids = Vec::new();
for entry in std::fs::read_dir(data_root)? {
let entry = entry?;
let path = entry.path();
if !path.join("palace.json").is_file() {
continue;
}
match path.file_name().and_then(|n| n.to_str()) {
Some(name) => ids.push(name.to_string()),
None => anyhow::bail!(
"palace directory name is not valid UTF-8: {}",
path.display()
),
}
}
Ok(ids)
}
pub(crate) fn release_swept_palace(
state: &AppState,
id: &trusty_common::memory_core::palace::PalaceId,
was_resident_before: bool,
handle: std::sync::Arc<trusty_common::memory_core::PalaceHandle>,
permit: crate::startup_budget::StartupOpenPermit,
) {
drop(handle);
let data_dir = state.data_root.join(&id.0);
crate::startup_budget::release_after_sweep(
&state.registry,
id,
&data_dir,
was_resident_before,
std::time::Duration::from_secs(crate::startup_budget::DEFAULT_KEEP_RECENT_SECS),
);
drop(permit);
}
pub async fn run_startup_sweep(state: &AppState) -> SweepOutcome {
let started = Instant::now();
let root = state.data_root.clone();
let listed = tokio::task::spawn_blocking(move || palace_ids_on_disk(&root)).await;
let palaces = match listed {
Ok(Ok(p)) => p,
Ok(Err(e)) => {
tracing::error!(
"bm25 backfill: could not enumerate palaces on disk — the startup sweep \
verified NOTHING and no palace was queued for repair: {e:#}"
);
return SweepOutcome::default();
}
Err(e) => {
tracing::error!(
"bm25 backfill: palace enumeration task failed — the startup sweep \
verified NOTHING: {e}"
);
return SweepOutcome::default();
}
};
let mut out = SweepOutcome {
enumerated: Some(palaces.len()),
..Default::default()
};
for palace in palaces {
let id = trusty_common::memory_core::palace::PalaceId::new(palace.clone());
let permit = state.startup_gate.acquire().await;
let was_resident_before = state.registry.peek(&id).is_some();
let registry = std::sync::Arc::clone(&state.registry);
let root = state.data_root.clone();
let open_id = id.clone();
let opened =
tokio::task::spawn_blocking(move || registry.open_palace(&root, &open_id)).await;
let handle = match opened {
Ok(Ok(h)) => h,
Ok(Err(e)) => {
tracing::error!(palace = %palace, "bm25 backfill: could not open palace: {e:#}");
out.unopenable += 1;
crate::bm25_repair::mark_dirty(state, &palace);
continue;
}
Err(e) => {
tracing::error!(palace = %palace, "bm25 backfill: open task failed: {e}");
out.unopenable += 1;
crate::bm25_repair::mark_dirty(state, &palace);
continue;
}
};
if handle.drawers.read().is_empty() {
release_swept_palace(state, &id, was_resident_before, handle, permit);
continue;
}
let report = backfill_state_palace(state, &handle, &palace, false).await;
out.swept += 1;
if !report.fully_indexed() {
out.incomplete += 1;
crate::bm25_repair::mark_dirty(state, &palace);
}
release_swept_palace(state, &id, was_resident_before, handle, permit);
}
if out.all_verified() {
tracing::info!(
enumerated = ?out.enumerated,
swept = out.swept,
elapsed_ms = started.elapsed().as_millis() as u64,
"bm25 backfill: startup sweep complete, all coverage verified"
);
} else {
tracing::error!(
enumerated = ?out.enumerated,
swept = out.swept,
incomplete = out.incomplete,
unopenable = out.unopenable,
elapsed_ms = started.elapsed().as_millis() as u64,
"bm25 backfill: startup sweep left palaces without verified coverage — \
queued for repair"
);
}
out
}
pub fn spawn_startup_backfill(state: &AppState) {
if state.bm25.is_none() {
tracing::debug!("bm25 backfill: lane disabled — skipping startup sweep");
return;
}
if startup_backfill_opted_out() {
tracing::info!("bm25 backfill: {ENV_NO_BACKFILL}=1 — skipping startup sweep");
return;
}
let state = state.clone();
tokio::spawn(async move {
run_startup_sweep(&state).await;
});
}
#[cfg(test)]
#[path = "bm25_backfill_tests.rs"]
mod tests;