use std::path::{Path, PathBuf};
use std::time::{Duration, Instant};
use trusty_common::bm25_client::{is_method_not_found, Bm25Client};
use trusty_common::memory_core::palace::Drawer;
use trusty_common::memory_core::retrieval::PalaceHandle;
use crate::AppState;
const OP_TIMEOUT: Duration = Duration::from_secs(10);
const PALACE_BUDGET: Duration = Duration::from_secs(120);
const COVERAGE_CHUNK: usize = 256;
pub const ENV_NO_BACKFILL: &str = "TRUSTY_BM25_NO_BACKFILL";
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum BackfillStatus {
Disabled,
DaemonUnavailable,
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.clone()));
}
}
PalaceDocs {
docs,
skipped_empty,
}
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
enum Coverage {
Missing(usize),
Unsupported,
Unreachable,
}
async fn probe_coverage(client: &Bm25Client, palace: &str, ids: &[String]) -> Coverage {
let mut missing = 0usize;
for chunk in ids.chunks(COVERAGE_CHUNK) {
match tokio::time::timeout(OP_TIMEOUT, client.missing_docs(chunk)).await {
Ok(Ok(cov)) => missing += cov.missing.len(),
Ok(Err(e)) if is_method_not_found(&e) => {
tracing::error!(
palace = %palace,
"bm25 backfill: daemon does not implement `missing_docs` — it predates \
0.2.0. Coverage CANNOT be verified for this palace; upgrade the daemon. \
Reporting the palace as incomplete rather than guessing: {e:#}"
);
return Coverage::Unsupported;
}
Ok(Err(e)) => {
tracing::warn!(palace = %palace, "bm25 backfill: coverage probe failed: {e:#}");
return Coverage::Unreachable;
}
Err(_) => {
tracing::warn!(palace = %palace, "bm25 backfill: coverage probe timed out");
return Coverage::Unreachable;
}
}
}
Coverage::Missing(missing)
}
pub async fn backfill_palace(
socket: &Path,
palace: &str,
palace_docs: PalaceDocs,
force: bool,
) -> BackfillReport {
let started = Instant::now();
let client = Bm25Client::new(socket.to_path_buf());
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(&client, 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(&client, 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::Unsupported => {}
Coverage::Unreachable => {
let mut report = BackfillReport::short_circuit(
palace,
BackfillStatus::DaemonUnavailable,
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 tokio::time::timeout(OP_TIMEOUT, client.index(doc_id, text)).await {
Ok(Ok(())) => indexed += 1,
Ok(Err(e)) => {
failed += 1;
tracing::warn!(palace = %palace, doc_id = %doc_id, "bm25 backfill index failed: {e:#}");
}
Err(_) => {
failed += 1;
tracing::warn!(palace = %palace, doc_id = %doc_id, "bm25 backfill index timed out");
}
}
}
let missing_after = match probe_coverage(&client, palace, &ids).await {
Coverage::Missing(n) => Some(n),
Coverage::Unsupported | Coverage::Unreachable => None,
};
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(&client, 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(client: &Bm25Client, palace: &str) -> Option<usize> {
match tokio::time::timeout(OP_TIMEOUT, client.stats()).await {
Ok(Ok(stats)) => Some(stats.doc_count),
Ok(Err(e)) => {
tracing::debug!(palace = %palace, "bm25 backfill: stats read failed: {e:#}");
None
}
Err(_) => 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 daemon holds documents for drawers this palace no longer has — \
stale lexical hits are possible (see #5053)"
);
}
}
async fn socket_for_palace(state: &AppState, palace: &str) -> Option<PathBuf> {
state.bm25_client.as_ref()?;
let supervisor = state.bm25_supervisor.as_ref()?;
let data_dir = state.data_root.join(palace).join("bm25");
match supervisor.ensure_running(palace, &data_dir).await {
Ok(socket) => Some(socket),
Err(e) => {
tracing::warn!(palace = %palace, "bm25 backfill: could not start daemon: {e:#}");
None
}
}
}
pub async fn backfill_state_palace(
state: &AppState,
handle: &PalaceHandle,
palace: &str,
force: bool,
) -> BackfillReport {
if state.bm25_client.is_none() {
return BackfillReport::short_circuit(palace, BackfillStatus::Disabled, 0);
}
let docs = palace_docs(handle);
let drawers_total = docs.drawers_total();
let Some(socket) = socket_for_palace(state, palace).await else {
let mut report =
BackfillReport::short_circuit(palace, BackfillStatus::DaemonUnavailable, drawers_total);
report.skipped_empty = docs.skipped_empty;
return report;
};
backfill_palace(&socket, 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 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 registry = std::sync::Arc::clone(&state.registry);
let root = state.data_root.clone();
let opened = tokio::task::spawn_blocking(move || registry.open_palace(&root, &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() {
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);
}
}
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_client.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;