use std::sync::Arc;
use super::helpers;
use super::{MapCache, ServerState};
pub(super) async fn run_background_gc(state: Arc<ServerState>) {
if state.store.read().await.blobs_shared {
tracing::debug!("background blob GC skipped: blob cache is shared across git worktrees");
return;
}
let result = tokio::task::spawn_blocking(move || {
let store = state.store.blocking_read();
let referenced = crate::store_gc::collect_referenced_hashes(&store.basemind_dir)?;
crate::store_gc::gc_blobs(&referenced)
})
.await;
match result {
Ok(Ok(report)) if report.removed > 0 => tracing::info!(
removed = report.removed,
bytes_freed = report.bytes_freed,
"background blob GC reclaimed orphaned blobs"
),
Ok(Ok(_)) => tracing::debug!("background blob GC: nothing to reclaim"),
Ok(Err(error)) => tracing::warn!(%error, "background blob GC failed"),
Err(error) => tracing::warn!(%error, "background blob GC task panicked"),
}
}
pub(super) fn spawn_initial_scan(state: Arc<ServerState>) {
tracing::info!("empty index on startup; running initial scan in background");
#[cfg(all(feature = "comms", any(unix, windows)))]
if state.daemon_writer {
tokio::spawn(async move {
use std::sync::atomic::Ordering;
state.initial_scan_active.store(true, Ordering::Relaxed);
let started = std::time::Instant::now();
match super::daemon_forward::forward_rescan_and_refresh(&state, None, false, false).await {
Ok(report) => tracing::info!(
scanned = report.scanned,
updated = report.updated,
elapsed_ms = started.elapsed().as_millis() as u64,
"initial scan complete (forwarded to daemon; embeddings deferred)"
),
Err(error) => tracing::warn!(%error, "initial forwarded scan failed"),
}
state
.initial_scan_ms
.store(started.elapsed().as_millis() as u64, Ordering::Relaxed);
state.initial_scan_active.store(false, Ordering::Relaxed);
let embed_state = Arc::clone(&state);
tokio::spawn(async move {
let embed_started = std::time::Instant::now();
tracing::info!("background embedding pass starting (forwarded to daemon)");
match super::daemon_forward::forward_rescan_and_refresh(&embed_state, None, false, true).await {
Ok(report) => tracing::info!(
scanned = report.scanned,
updated = report.updated,
elapsed_ms = embed_started.elapsed().as_millis() as u64,
"background embedding pass complete (forwarded to daemon)"
),
Err(error) => tracing::warn!(%error, "background forwarded embedding pass failed"),
}
});
});
return;
}
tokio::spawn(async move {
use std::sync::atomic::Ordering;
state.initial_scan_active.store(true, Ordering::Relaxed);
let started = std::time::Instant::now();
match helpers::scan_and_refresh(Arc::clone(&state), None, crate::scanner::EmbedMode::Deferred).await {
Ok(report) => tracing::info!(
scanned = report.stats.scanned,
updated = report.stats.updated,
elapsed_ms = started.elapsed().as_millis() as u64,
"initial background scan complete (code-map + keyword lane; embeddings deferred)"
),
Err(error) => tracing::warn!(%error, "initial background scan failed"),
}
state
.initial_scan_ms
.store(started.elapsed().as_millis() as u64, Ordering::Relaxed);
state.initial_scan_active.store(false, Ordering::Relaxed);
let embed_state = Arc::clone(&state);
tokio::spawn(async move {
let embed_started = std::time::Instant::now();
tracing::info!("background embedding pass starting");
match helpers::scan_and_refresh(Arc::clone(&embed_state), None, crate::scanner::EmbedMode::Inline).await {
Ok(report) => tracing::info!(
scanned = report.stats.scanned,
updated = report.stats.updated,
elapsed_ms = embed_started.elapsed().as_millis() as u64,
"background embedding pass complete"
),
Err(error) => tracing::warn!(%error, "background embedding pass failed"),
}
run_background_gc(embed_state).await;
});
});
}
pub(super) fn spawn_cache_warm(state: Arc<ServerState>) {
tracing::info!("warming in-RAM code map in background (handshake already served)");
tokio::spawn(async move {
use std::sync::atomic::Ordering;
let started = std::time::Instant::now();
let build_state = Arc::clone(&state);
let built = tokio::task::spawn_blocking(move || {
let store = build_state.store.blocking_read();
MapCache::build(&store)
})
.await;
match built {
Ok(cache) => {
let files = cache.by_path.len();
state.cache.store(Arc::new(cache));
state.cache_generation.fetch_add(1, Ordering::Relaxed);
state
.cache_warm_ms
.store(started.elapsed().as_millis() as u64, Ordering::Relaxed);
state.cache_warming.store(false, Ordering::Relaxed);
state.cache_ready.notify_waiters();
tracing::info!(
files,
elapsed_ms = started.elapsed().as_millis() as u64,
"in-RAM code map warm complete"
);
}
Err(error) => {
state.cache_warming.store(false, Ordering::Relaxed);
state.cache_ready.notify_waiters();
tracing::error!(%error, "in-RAM code map warm task panicked; serving un-warmed cache");
}
}
});
}
fn refresh_batch(
handle: &tokio::runtime::Handle,
state: &Arc<ServerState>,
paths: Vec<std::path::PathBuf>,
) -> Result<(usize, usize, usize), String> {
#[cfg(all(feature = "comms", any(unix, windows)))]
if state.daemon_writer {
let report = handle
.block_on(super::daemon_forward::forward_rescan_and_refresh(
state,
Some(paths),
false,
true,
))
.map_err(|error| error.to_string())?;
return Ok((report.scanned, report.updated, report.removed));
}
let report = handle
.block_on(helpers::scan_and_refresh(
Arc::clone(state),
Some(paths),
crate::scanner::EmbedMode::Inline,
))
.map_err(|error| error.to_string())?;
Ok((report.stats.scanned, report.stats.updated, report.stats.removed))
}
pub(super) fn spawn_serve_watcher(state: Arc<ServerState>) {
let root = state.root.clone();
let config = Arc::clone(&state.config);
let handle = tokio::runtime::Handle::current();
let (_shutdown_tx, shutdown_rx) = tokio::sync::oneshot::channel::<()>();
std::thread::Builder::new()
.name("basemind-mcp-serve-watcher".to_string())
.spawn(move || {
let _keep_sender_alive = _shutdown_tx;
tracing::info!(root = %root.display(), "serve watcher armed (live incremental rescan)");
let result = crate::watcher::watch_paths(&root, &config, shutdown_rx, |paths, _kind| {
use std::sync::atomic::Ordering;
let refresh_state = Arc::clone(&state);
refresh_state.rescan_active.store(true, Ordering::Relaxed);
let outcome = refresh_batch(&handle, &refresh_state, paths);
refresh_state.rescan_active.store(false, Ordering::Relaxed);
match outcome {
Ok((scanned, updated, removed)) => {
tracing::debug!(scanned, updated, removed, "serve watcher: incremental rescan complete")
}
Err(error) => tracing::warn!(
%error,
"serve watcher: incremental rescan failed (watcher continues)"
),
}
});
if let Err(error) = result {
tracing::warn!(%error, "serve watcher exited with error");
}
tracing::info!("serve watcher: exiting");
})
.ok();
}
fn reopen_read_only(state: &ServerState, view: &str) -> Result<crate::store::Store, crate::store::StoreError> {
#[cfg(all(feature = "comms", any(unix, windows)))]
if state.daemon_writer {
return crate::store::Store::open_read_only_no_index(state.root.as_path(), view);
}
crate::store::Store::open_read_only(state.root.as_path(), view)
}
pub(super) fn spawn_view_watcher(state: Arc<ServerState>) {
let (basemind_dir, view) = {
let store = match state.store.try_read() {
Ok(g) => g,
Err(_) => return,
};
(store.basemind_dir.clone(), store.view.clone())
};
let view_dir = basemind_dir.join(crate::store::VIEWS_DIR).join(&view);
let target = view_dir.join(crate::store::INDEX_FILE);
std::thread::Builder::new()
.name("basemind-mcp-view-watcher".to_string())
.spawn(move || {
use notify_debouncer_full::new_debouncer;
use std::time::Duration;
let (tx, rx) = std::sync::mpsc::channel();
let mut debouncer = match new_debouncer(Duration::from_millis(150), None, tx) {
Ok(d) => d,
Err(e) => {
tracing::warn!(error = %e, "view watcher: failed to start debouncer");
return;
}
};
if let Err(e) = debouncer.watch(&view_dir, notify::RecursiveMode::NonRecursive) {
tracing::warn!(error = %e, dir = %view_dir.display(), "view watcher: failed to watch");
return;
}
tracing::info!(target = %target.display(), "view watcher armed");
while let Ok(result) = rx.recv() {
let events = match result {
Ok(e) => e,
Err(_) => continue,
};
let touches_index = events.iter().any(|de| de.event.paths.iter().any(|p| p == &target));
if !touches_index {
continue;
}
let view = state.store.try_read().map(|g| g.view.clone()).unwrap_or_default();
let new_store = match reopen_read_only(&state, &view) {
Ok(s) => s,
Err(e) => {
tracing::warn!(error = %e, "view watcher: store reopen failed");
continue;
}
};
let fingerprint = super::map_fingerprint::index_fingerprint(&new_store);
if fingerprint == state.cache.load().fingerprint {
tracing::debug!("view watcher: index rewritten but unchanged; keeping the current MapCache");
continue;
}
let new_cache = Arc::new(MapCache::build(&new_store));
tracing::info!(
files = new_cache.by_path.len(),
"view watcher: rebuilt MapCache from refreshed index"
);
state.cache.store(new_cache);
state
.cache_generation
.fetch_add(1, std::sync::atomic::Ordering::Relaxed);
}
tracing::info!("view watcher: channel closed; exiting");
})
.ok();
}