use super::db::Store;
use std::collections::HashMap;
use std::sync::Mutex as StdMutex;
use std::sync::atomic::{AtomicBool, Ordering};
use std::time::SystemTime;
static LAST_INDEXED: StdMutex<Option<HashMap<String, SystemTime>>> = StdMutex::new(None);
static REFRESHING: AtomicBool = AtomicBool::new(false);
pub async fn refresh_stale_brain_files() -> usize {
if REFRESHING
.compare_exchange(false, true, Ordering::AcqRel, Ordering::Acquire)
.is_err()
{
tracing::debug!("Memory freshness: refresh already in flight, skipping");
return 0;
}
let refreshed = refresh_inner().await;
REFRESHING.store(false, Ordering::Release);
refreshed
}
async fn refresh_inner() -> usize {
let home = crate::config::opencrabs_home();
let mut refreshed = 0usize;
for &name in super::BRAIN_FILES {
let path = home.join(name);
let disk_mtime = match std::fs::metadata(&path).and_then(|m| m.modified()) {
Ok(t) => t,
Err(_) => continue,
};
let seen = {
let guard = LAST_INDEXED.lock().unwrap_or_else(|e| e.into_inner());
guard.as_ref().and_then(|m| m.get(name).copied())
};
let changed = match seen {
Some(t) => disk_mtime > t,
None => true,
};
if changed && index_one(&path, name).await {
refreshed += 1;
let mut guard = LAST_INDEXED.lock().unwrap_or_else(|e| e.into_inner());
guard
.get_or_insert_with(HashMap::new)
.insert(name.to_string(), disk_mtime);
}
}
if refreshed > 0 {
tracing::info!("Memory freshness: reindexed {refreshed} stale brain file(s)");
}
refreshed
}
async fn index_one(path: &std::path::Path, name: &str) -> bool {
let store = match super::get_store() {
Ok(s) => s,
Err(e) => {
tracing::warn!("Memory freshness: store unavailable for {name}: {e}");
return false;
}
};
match super::index_file_fts_only(store, path).await {
Ok(()) => true,
Err(e) => {
tracing::warn!("Memory freshness: failed to reindex {name}: {e}");
false
}
}
}
static EXTERNAL_LAST_INDEXED: StdMutex<Option<HashMap<String, SystemTime>>> = StdMutex::new(None);
static EXTERNAL_REFRESHING: AtomicBool = AtomicBool::new(false);
pub async fn refresh_stale_external(paths: &[String]) -> usize {
if paths.is_empty() {
return 0;
}
if EXTERNAL_REFRESHING
.compare_exchange(false, true, Ordering::AcqRel, Ordering::Acquire)
.is_err()
{
tracing::debug!("Memory external freshness: refresh already in flight, skipping");
return 0;
}
let refreshed = refresh_external_inner(paths).await;
EXTERNAL_REFRESHING.store(false, Ordering::Release);
refreshed
}
async fn refresh_external_inner(paths: &[String]) -> usize {
let store = match super::get_store() {
Ok(s) => s,
Err(e) => {
tracing::warn!("Memory external freshness: store unavailable: {e}");
return 0;
}
};
let mut refreshed = 0usize;
for key in paths {
let path = std::path::Path::new(key);
if !path.is_absolute() {
continue;
}
let disk_mtime = match std::fs::metadata(path).and_then(|m| m.modified()) {
Ok(t) => t,
Err(_) => continue,
};
let seen = {
let guard = EXTERNAL_LAST_INDEXED
.lock()
.unwrap_or_else(|e| e.into_inner());
guard.as_ref().and_then(|m| m.get(key).copied())
};
let changed = match seen {
Some(t) => disk_mtime > t,
None => true,
};
if !changed {
continue;
}
if index_external_one(store, key, path).await {
refreshed += 1;
let mut guard = EXTERNAL_LAST_INDEXED
.lock()
.unwrap_or_else(|e| e.into_inner());
guard
.get_or_insert_with(HashMap::new)
.insert(key.clone(), disk_mtime);
}
}
if refreshed > 0 {
tracing::info!("Memory external freshness: reindexed {refreshed} stale external file(s)");
}
refreshed
}
async fn index_external_one(
store: &'static StdMutex<Store>,
key: &str,
path: &std::path::Path,
) -> bool {
let body = match tokio::fs::read_to_string(path).await {
Ok(b) => b,
Err(e) => {
tracing::warn!("Memory external freshness: unreadable {key}: {e}");
return false;
}
};
let key = key.to_string();
let log_key = key.clone();
let result = tokio::task::spawn_blocking(move || {
let s = store
.lock()
.map_err(|e| format!("Store lock poisoned: {e}"))?;
super::index::index_file_sync_keyed(&s, super::COLLECTION_EXTERNAL, &key, &body)
})
.await;
match result {
Ok(Ok(_)) => true,
Ok(Err(e)) => {
tracing::warn!("Memory external freshness: failed to reindex {log_key}: {e}");
false
}
Err(e) => {
tracing::warn!("Memory external freshness: join failed for {log_key}: {e}");
false
}
}
}
static SWEEP_SPAWNED: AtomicBool = AtomicBool::new(false);
static SWEEP_RUNNING: AtomicBool = AtomicBool::new(false);
pub fn spawn_external_sweep() {
if SWEEP_SPAWNED.swap(true, Ordering::AcqRel) {
return;
}
tokio::spawn(async {
loop {
let secs = super::sweep_interval_secs().max(10);
tokio::time::sleep(std::time::Duration::from_secs(secs)).await;
if SWEEP_RUNNING
.compare_exchange(false, true, Ordering::AcqRel, Ordering::Acquire)
.is_err()
{
continue;
}
sweep_once().await;
SWEEP_RUNNING.store(false, Ordering::Release);
}
});
}
async fn sweep_once() {
let store = match super::get_store() {
Ok(s) => s,
Err(e) => {
tracing::warn!("Memory external sweep: store unavailable: {e}");
return;
}
};
let result = tokio::task::spawn_blocking(move || {
let s = match store.lock() {
Ok(s) => s,
Err(e) => return Err(format!("Store lock poisoned: {e}")),
};
Ok::<_, String>(super::external_sweep::sweep_external(&s))
})
.await;
match result {
Ok(Ok(report)) => report.log(),
Ok(Err(e)) => tracing::warn!("Memory external sweep failed: {e}"),
Err(e) => tracing::warn!("Memory external sweep join failed: {e}"),
}
}