use std::collections::HashMap;
use std::hash::Hash;
use std::path::{Path, PathBuf};
use std::sync::Arc;
use std::sync::OnceLock;
use std::sync::atomic::{AtomicU32, AtomicU64, AtomicUsize, Ordering};
use std::time::{Duration, Instant, SystemTime};
use tokio::sync::{Mutex, OwnedMutexGuard, OwnedSemaphorePermit, Semaphore};
use tokio::task::JoinHandle;
use tokio_util::sync::CancellationToken;
use super::pressure::IoPressureProbe;
#[cfg(not(target_os = "linux"))]
use super::pressure::NoPressureSignal;
#[cfg(target_os = "linux")]
use super::pressure::PsiProbe;
use hallouminate_adapters::{
EmbedBatch, Embedder, FastembedCrossencoder, LanceStore, StoreLockOwner,
};
use hallouminate_config::Config;
use hallouminate_domain::common::{HallouminateError, expand_tilde};
use hallouminate_domain::corpus::{Tokenizer, load_tokenizer, missing_roots};
use hallouminate_domain::indexer::{HandlerRegistry, SearchHit, index_corpus};
use hallouminate_domain::search::{Crossencoder, canonical_crossencoder_model};
use super::ladder::LadderAction;
use super::maintenance::{DeferReason, maintenance_loop};
const CHUNK_BUDGET_TOKENS: usize = 384;
pub(crate) const STALE_BACKUP_MAX_AGE: Duration = Duration::from_secs(30 * 24 * 60 * 60);
static PROCESS_START: OnceLock<Instant> = OnceLock::new();
fn monotonic_secs() -> u64 {
PROCESS_START.get_or_init(Instant::now).elapsed().as_secs()
}
async fn init_embedder(
model: &str,
quantized: bool,
cache_dir: PathBuf,
) -> anyhow::Result<Box<dyn EmbedBatch>> {
let model = model.to_owned();
tokio::task::spawn_blocking(move || Embedder::try_new(&model, quantized, &cache_dir))
.await
.map_err(|e| anyhow::anyhow!("embedder initialization task failed: {e}"))?
.map(|embedder| Box::new(embedder) as Box<dyn EmbedBatch>)
.map_err(|e| anyhow::anyhow!("init embedder: {e}"))
}
fn is_idle(last_use_secs: u64, now_secs: u64, idle_secs: u64) -> bool {
now_secs.saturating_sub(last_use_secs) >= idle_secs
}
struct KeyedLockMap<K> {
inner: Mutex<HashMap<K, Arc<Mutex<()>>>>,
}
impl<K> Default for KeyedLockMap<K> {
fn default() -> Self {
Self {
inner: Mutex::new(HashMap::new()),
}
}
}
impl<K: Eq + Hash> KeyedLockMap<K> {
async fn lock<Q>(&self, key: &Q) -> OwnedMutexGuard<()>
where
Q: ToOwned<Owned = K> + ?Sized,
{
let mutex = {
let mut map = self.inner.lock().await;
map.entry(key.to_owned())
.or_insert_with(|| Arc::new(Mutex::new(())))
.clone()
};
mutex.lock_owned().await
}
}
#[derive(Clone, PartialEq, Eq, Hash)]
struct ResourceKey {
ground_dir: PathBuf,
model: String,
quantized: bool,
enabled: bool,
}
impl ResourceKey {
fn from_config(cfg: &Config) -> Self {
Self {
ground_dir: expand_tilde(&cfg.storage.ground_dir),
model: cfg.embeddings.model.clone(),
quantized: cfg.embeddings.quantized,
enabled: cfg.embeddings.enabled,
}
}
}
pub struct RequestResources {
pub store: Arc<LanceStore>,
pub tokenizer: Tokenizer,
pub embeddings_enabled: bool,
pub ground_dir: PathBuf,
}
#[derive(Clone)]
pub struct DaemonState {
inner: Arc<DaemonStateInner>,
}
struct DaemonStateInner {
baseline: Config,
baseline_xdg_path: Option<PathBuf>,
baseline_resources: Arc<RequestResources>,
resources: Mutex<HashMap<ResourceKey, Arc<RequestResources>>>,
resource_build_locks: KeyedLockMap<ResourceKey>,
store_lock_owner: StoreLockOwner,
corpus_locks: KeyedLockMap<String>,
write_lane: Arc<Semaphore>,
crossencoders: Arc<Mutex<HashMap<String, FastembedCrossencoder>>>,
last_activity_secs: Arc<AtomicU64>,
active_connections: Arc<AtomicUsize>,
external_last_activity_secs: Arc<AtomicU64>,
internal_last_activity_secs: Arc<AtomicU64>,
active_external_connections: Arc<AtomicUsize>,
active_internal_connections: Arc<AtomicUsize>,
defer_count: AtomicU32,
watcher_counters: WatcherCounters,
last_ladder_trip: std::sync::Mutex<Option<LadderTrip>>,
supervisor: Arc<super::supervisor::Supervisor>,
heartbeat: Arc<super::heartbeat::HeartbeatRegistry>,
maintenance_task: Mutex<Option<JoinHandle<()>>>,
shutdown: CancellationToken,
}
pub struct MutationGuard {
_permit: OwnedSemaphorePermit,
_corpus: OwnedMutexGuard<()>,
}
impl MutationGuard {
pub(super) fn new(permit: OwnedSemaphorePermit, corpus: OwnedMutexGuard<()>) -> Self {
Self {
_permit: permit,
_corpus: corpus,
}
}
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum WorkClass {
External,
Internal,
}
pub struct ConnectionGuard {
active: Arc<AtomicUsize>,
class_active: Arc<AtomicUsize>,
}
impl Drop for ConnectionGuard {
fn drop(&mut self) {
self.active.fetch_sub(1, Ordering::SeqCst);
self.class_active.fetch_sub(1, Ordering::SeqCst);
}
}
#[derive(Debug, Default)]
struct WatcherCounters {
events: AtomicU64,
reindexes: AtomicU64,
noop_reindexes: AtomicU64,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub(crate) struct LadderTrip {
pub(crate) action: LadderAction,
pub(crate) at_secs: u64,
}
impl DaemonState {
pub async fn open(cfg: Config, xdg_path: Option<PathBuf>) -> anyhow::Result<Self> {
Self::open_with_owner(cfg, xdg_path, StoreLockOwner::for_process()).await
}
pub(crate) async fn open_with_socket(
cfg: Config,
xdg_path: Option<PathBuf>,
socket: PathBuf,
) -> anyhow::Result<Self> {
Self::open_with_owner(cfg, xdg_path, StoreLockOwner::for_daemon(socket)).await
}
async fn open_with_owner(
cfg: Config,
xdg_path: Option<PathBuf>,
store_lock_owner: StoreLockOwner,
) -> anyhow::Result<Self> {
let ground_dir = expand_tilde(&cfg.storage.ground_dir);
if let Some(parent) = ground_dir.parent()
&& !parent.as_os_str().is_empty()
{
tokio::fs::create_dir_all(parent)
.await
.map_err(|e| anyhow::anyhow!("create ground dir parent: {e}"))?;
}
let stale_version = match LanceStore::validate_existing_metadata(
&ground_dir,
&cfg.embeddings.model,
cfg.embeddings.quantized,
cfg.embeddings.enabled,
) {
Ok(()) => None,
Err(HallouminateError::StoreSchemaStale {
found, expected, ..
}) => {
tracing::warn!(
target: "hallouminate::daemon",
%found,
%expected,
"ground store schema v{found} < expected v{expected}; rebuilding from source",
);
move_stale_store(&ground_dir, found).await?;
Some(found)
}
Err(e) => {
return Err(anyhow::anyhow!(
"validate ground dir {}: {e}",
ground_dir.display()
));
}
};
let cache_dir = expand_tilde(&cfg.embeddings.cache_dir);
let embedder: Option<Box<dyn EmbedBatch>> = if cfg.embeddings.enabled {
match init_embedder(
&cfg.embeddings.model,
cfg.embeddings.quantized,
cache_dir.clone(),
)
.await
{
Ok(embedder) => Some(embedder),
Err(e) => {
tracing::warn!(
target: "hallouminate::daemon",
model = %cfg.embeddings.model,
error = %e,
"embedder unavailable at startup; the next request will retry initialization",
);
None
}
}
} else {
None
};
let tokenizer = load_tokenizer(&cfg.embeddings.model)
.map_err(|e| anyhow::anyhow!("load tokenizer for {}: {e}", cfg.embeddings.model))?;
let build_result: anyhow::Result<LanceStore> = async {
let store = LanceStore::open_or_create_with_owner(
&ground_dir,
&cfg.embeddings.model,
cfg.embeddings.quantized,
cfg.embeddings.enabled,
embedder,
store_lock_owner.clone(),
)
.await
.map_err(|e| anyhow::anyhow!("open ground dir {}: {e}", ground_dir.display()))?;
if stale_version.is_some() {
let registry = HandlerRegistry::new(tokenizer.clone(), CHUNK_BUDGET_TOKENS);
for corpus in cfg
.effective_corpora()
.map_err(|e| anyhow::anyhow!("rebuild: list corpora: {e}"))?
{
let missing = missing_roots(&corpus);
if !missing.is_empty() {
tracing::warn!(
target: "hallouminate::daemon",
corpus = %corpus.name,
"rebuild: corpus root missing; skipped",
);
continue;
}
let stats = index_corpus(&corpus, &store, ®istry)
.await
.map_err(|e| anyhow::anyhow!("rebuild: index {}: {e}", corpus.name))?;
tracing::info!(
target: "hallouminate::daemon",
corpus = %corpus.name,
files = stats.files_upserted,
chunks = stats.chunks_inserted,
"rebuild: reindexed",
);
}
}
Ok(store)
}
.await;
let store = match build_result {
Ok(store) => store,
Err(e) => {
if let Some(found) = stale_version
&& ground_dir.exists()
{
let _ = tokio::fs::remove_dir_all(&ground_dir).await;
tracing::warn!(
target: "hallouminate::daemon",
"rebuild failed; removed partial ground dir so next boot retries. \
Backup preserved at {}.bak-v{found}",
ground_dir.display(),
);
}
return Err(e);
}
};
if let Err(e) =
prune_stale_backups(&ground_dir, SystemTime::now(), STALE_BACKUP_MAX_AGE).await
{
tracing::warn!(
target: "hallouminate::daemon",
error = %e,
"failed to prune stale ground store backups",
);
}
let mut crossencoders: HashMap<String, FastembedCrossencoder> = HashMap::new();
if let Some(model) = cfg.search.crossencoder.as_deref() {
match canonical_crossencoder_model(model)
.map_err(anyhow::Error::from)
.and_then(|canonical| {
FastembedCrossencoder::try_new(canonical, &cache_dir)
.map(|c| (canonical, c))
.map_err(anyhow::Error::from)
}) {
Ok((canonical, c)) => {
crossencoders.insert(canonical.to_string(), c);
}
Err(e) => {
tracing::warn!(
target: "hallouminate::daemon",
model = %model,
error = %e,
"crossencoder unavailable at startup; ground will skip rerank until reload",
);
}
}
}
let shutdown = CancellationToken::new();
let crossencoders_arc = Arc::new(Mutex::new(crossencoders));
let last_activity = Arc::new(AtomicU64::new(monotonic_secs()));
let store = Arc::new(store);
let write_lane = Arc::new(Semaphore::new(1));
if cfg.embeddings.idle_evict_secs
!= hallouminate_config::EmbeddingsConfig::default().idle_evict_secs
{
tracing::warn!(
target: "hallouminate::daemon",
idle_evict_secs = cfg.embeddings.idle_evict_secs,
"embeddings.idle_evict_secs is deprecated and does nothing; \
set [daemon].idle_exit_secs to control idle-exit instead",
);
}
let baseline_key = ResourceKey {
ground_dir: ground_dir.clone(),
model: cfg.embeddings.model.clone(),
quantized: cfg.embeddings.quantized,
enabled: cfg.embeddings.enabled,
};
let baseline_resources = Arc::new(RequestResources {
store,
tokenizer,
embeddings_enabled: cfg.embeddings.enabled,
ground_dir,
});
let mut resources_map = HashMap::new();
resources_map.insert(baseline_key, Arc::clone(&baseline_resources));
let maintenance_interval_secs = cfg.daemon.maintenance_interval_secs;
let restart_cap = cfg.daemon.restart_intensity_cap;
let restart_window = Duration::from_secs(cfg.daemon.restart_intensity_window_secs);
let heartbeat = Arc::new(super::heartbeat::HeartbeatRegistry::default());
let ladder = super::ladder::Ladder {
warn_at: 3,
act_at: 5,
action: LadderAction::WatchdogTrip,
};
let state = DaemonState {
inner: Arc::new_cyclic(|weak: &std::sync::Weak<DaemonStateInner>| {
let weak = weak.clone();
let escalate: super::supervisor::EscalationHook = Arc::new(move |task, action| {
tracing::error!(
target: "hallouminate::daemon",
task = ?task,
action = ?action,
"supervised task exceeded the restart intensity cap",
);
if let Some(inner) = weak.upgrade() {
*inner
.last_ladder_trip
.lock()
.expect("ladder trip mutex poisoned") = Some(LadderTrip {
action,
at_secs: monotonic_secs(),
});
}
});
DaemonStateInner {
baseline: cfg,
baseline_xdg_path: xdg_path,
baseline_resources,
resources: Mutex::new(resources_map),
resource_build_locks: KeyedLockMap::default(),
corpus_locks: KeyedLockMap::default(),
store_lock_owner,
write_lane,
crossencoders: crossencoders_arc,
last_activity_secs: last_activity,
active_connections: Arc::new(AtomicUsize::new(0)),
external_last_activity_secs: Arc::new(AtomicU64::new(monotonic_secs())),
internal_last_activity_secs: Arc::new(AtomicU64::new(monotonic_secs())),
active_external_connections: Arc::new(AtomicUsize::new(0)),
active_internal_connections: Arc::new(AtomicUsize::new(0)),
defer_count: AtomicU32::new(0),
watcher_counters: WatcherCounters::default(),
last_ladder_trip: std::sync::Mutex::new(None),
supervisor: Arc::new(super::supervisor::Supervisor::new(
restart_cap,
restart_window,
ladder,
escalate,
shutdown.clone(),
)),
heartbeat,
maintenance_task: Mutex::new(None),
shutdown,
}
}),
};
if maintenance_interval_secs == 0 {
tracing::info!(
target: "hallouminate::daemon",
"automatic maintenance disabled (daemon.maintenance_interval_secs = 0)",
);
} else {
let loop_state = state.clone();
let cancel = loop_state.shutdown_token().clone();
let interval = Duration::from_secs(maintenance_interval_secs);
#[cfg(target_os = "linux")]
let probe: Arc<dyn IoPressureProbe> = Arc::new(PsiProbe);
#[cfg(not(target_os = "linux"))]
let probe: Arc<dyn IoPressureProbe> = Arc::new(NoPressureSignal);
let maintenance_task =
state
.inner
.supervisor
.spawn(super::heartbeat::TaskName::Maintenance, move || {
maintenance_loop(
loop_state.clone(),
cancel.clone(),
interval,
probe.clone(),
)
});
*state.inner.maintenance_task.lock().await = Some(maintenance_task);
}
Ok(state)
}
pub fn shutdown_token(&self) -> &CancellationToken {
&self.inner.shutdown
}
pub(crate) async fn take_maintenance_task(&self) -> Option<JoinHandle<()>> {
self.inner.maintenance_task.lock().await.take()
}
pub(crate) fn supervisor(&self) -> &Arc<super::supervisor::Supervisor> {
&self.inner.supervisor
}
pub(crate) fn heartbeat(&self) -> &Arc<super::heartbeat::HeartbeatRegistry> {
&self.inner.heartbeat
}
pub fn baseline_xdg_path(&self) -> Option<&Path> {
self.inner.baseline_xdg_path.as_deref()
}
pub fn baseline(&self) -> &Config {
&self.inner.baseline
}
pub fn store(&self) -> Arc<LanceStore> {
self.inner.baseline_resources.store.clone()
}
pub fn ground_dir(&self) -> &std::path::Path {
&self.inner.baseline_resources.ground_dir
}
pub fn embeddings_enabled(&self) -> bool {
self.inner.baseline_resources.embeddings_enabled
}
pub async fn resources_for(&self, cfg: &Config) -> anyhow::Result<Arc<RequestResources>> {
self.resources_for_with_initializer(cfg, |model, quantized, cache_dir| async move {
init_embedder(&model, quantized, cache_dir).await
})
.await
}
async fn resources_for_with_initializer<F, Fut>(
&self,
cfg: &Config,
initialize: F,
) -> anyhow::Result<Arc<RequestResources>>
where
F: FnOnce(String, bool, PathBuf) -> Fut,
Fut: std::future::Future<Output = anyhow::Result<Box<dyn EmbedBatch>>>,
{
let key = ResourceKey::from_config(cfg);
if let Some(existing) = self.inner.resources.lock().await.get(&key)
&& (!cfg.embeddings.enabled || existing.store.embedder_available())
{
return Ok(Arc::clone(existing));
}
let _build = self.inner.resource_build_locks.lock(&key).await;
if let Some(existing) = self.inner.resources.lock().await.get(&key).cloned() {
if cfg.embeddings.enabled && !existing.store.embedder_available() {
let cache_dir = expand_tilde(&cfg.embeddings.cache_dir);
let retry = match initialize(
cfg.embeddings.model.clone(),
cfg.embeddings.quantized,
cache_dir,
)
.await
{
Ok(embedder) => existing
.store
.install_embedder(embedder)
.map_err(anyhow::Error::from),
Err(error) => Err(error),
};
if let Err(error) = retry {
tracing::warn!(
target: "hallouminate::daemon",
model = %cfg.embeddings.model,
error = %error,
"embedder retry failed; serving cached resources without embeddings",
);
}
}
return Ok(existing);
}
let ground_dir = key.ground_dir.clone();
if let Some(parent) = ground_dir.parent()
&& !parent.as_os_str().is_empty()
{
tokio::fs::create_dir_all(parent)
.await
.map_err(|e| anyhow::anyhow!("create ground dir parent: {e}"))?;
}
let cache_dir = expand_tilde(&cfg.embeddings.cache_dir);
let embedder = if cfg.embeddings.enabled {
Some(
initialize(
cfg.embeddings.model.clone(),
cfg.embeddings.quantized,
cache_dir,
)
.await?,
)
} else {
None
};
let tokenizer = load_tokenizer(&cfg.embeddings.model)
.map_err(|e| anyhow::anyhow!("load tokenizer for {}: {e}", cfg.embeddings.model))?;
let store = LanceStore::open_or_create_with_owner(
&ground_dir,
&cfg.embeddings.model,
cfg.embeddings.quantized,
cfg.embeddings.enabled,
embedder,
self.inner.store_lock_owner.clone(),
)
.await
.map_err(|e| anyhow::anyhow!("open ground dir {}: {e}", ground_dir.display()))?;
let resources = Arc::new(RequestResources {
store: Arc::new(store),
tokenizer,
embeddings_enabled: cfg.embeddings.enabled,
ground_dir: ground_dir.clone(),
});
self.inner
.resources
.lock()
.await
.insert(key, Arc::clone(&resources));
Ok(resources)
}
pub async fn crossencoder(
&self,
model_name: Option<&str>,
) -> anyhow::Result<Option<CrossencoderGuard>> {
let Some(model_name) = model_name else {
return Ok(None);
};
let canonical = canonical_crossencoder_model(model_name)?;
let mut guard = Arc::clone(&self.inner.crossencoders).lock_owned().await;
if !guard.contains_key(canonical) {
let cache_dir = expand_tilde(&self.inner.baseline.embeddings.cache_dir);
let model = FastembedCrossencoder::try_new(canonical, &cache_dir)
.map_err(|e| anyhow::anyhow!("init crossencoder ({canonical}): {e}"))?;
guard.insert(canonical.to_string(), model);
}
self.inner
.last_activity_secs
.store(monotonic_secs(), Ordering::Relaxed);
Ok(Some(CrossencoderGuard {
guard,
key: canonical.to_string(),
last_use_secs: Arc::clone(&self.inner.last_activity_secs),
}))
}
pub fn touch_activity(&self, class: WorkClass) {
let now = monotonic_secs();
self.inner.last_activity_secs.store(now, Ordering::Relaxed);
match class {
WorkClass::External => self
.inner
.external_last_activity_secs
.store(now, Ordering::Relaxed),
WorkClass::Internal => self
.inner
.internal_last_activity_secs
.store(now, Ordering::Relaxed),
}
}
pub fn enter_connection(&self, class: WorkClass) -> ConnectionGuard {
self.inner.active_connections.fetch_add(1, Ordering::SeqCst);
let class_counter = match class {
WorkClass::External => &self.inner.active_external_connections,
WorkClass::Internal => &self.inner.active_internal_connections,
};
class_counter.fetch_add(1, Ordering::SeqCst);
ConnectionGuard {
active: Arc::clone(&self.inner.active_connections),
class_active: Arc::clone(class_counter),
}
}
pub(crate) fn defer_count(&self) -> u32 {
self.inner.defer_count.load(Ordering::Relaxed)
}
pub(super) fn reset_defer_count(&self) {
self.inner.defer_count.store(0, Ordering::Relaxed);
}
pub(super) fn increment_defer_count(&self) -> u32 {
self.inner.defer_count.fetch_add(1, Ordering::Relaxed) + 1
}
pub(crate) fn record_watcher_events(&self, count: u64) {
self.inner
.watcher_counters
.events
.fetch_add(count, Ordering::Relaxed);
}
pub(crate) fn record_watcher_reindex(&self, noop: bool) {
self.inner
.watcher_counters
.reindexes
.fetch_add(1, Ordering::Relaxed);
if noop {
self.inner
.watcher_counters
.noop_reindexes
.fetch_add(1, Ordering::Relaxed);
}
}
pub(crate) fn watcher_counters_snapshot(&self) -> (u64, u64, u64) {
(
self.inner.watcher_counters.events.load(Ordering::Relaxed),
self.inner
.watcher_counters
.reindexes
.load(Ordering::Relaxed),
self.inner
.watcher_counters
.noop_reindexes
.load(Ordering::Relaxed),
)
}
pub(crate) fn record_ladder_trip(&self, action: LadderAction) {
let at_secs = monotonic_secs();
*self
.inner
.last_ladder_trip
.lock()
.expect("ladder trip mutex poisoned") = Some(LadderTrip { action, at_secs });
}
pub(crate) fn last_ladder_trip(&self) -> Option<LadderTrip> {
*self
.inner
.last_ladder_trip
.lock()
.expect("ladder trip mutex poisoned")
}
pub(super) fn maintenance_defer_reason(
&self,
probe: &dyn IoPressureProbe,
) -> Option<DeferReason> {
self.maintenance_defer_reason_at(probe, monotonic_secs())
}
fn maintenance_defer_reason_at(
&self,
probe: &dyn IoPressureProbe,
now_secs: u64,
) -> Option<DeferReason> {
if self
.inner
.active_external_connections
.load(Ordering::SeqCst)
!= 0
{
return Some(DeferReason::Active);
}
if !is_idle(
self.inner
.external_last_activity_secs
.load(Ordering::Relaxed),
now_secs,
60,
) {
return Some(DeferReason::Active);
}
if probe.elevated() {
return Some(DeferReason::IoPressure);
}
None
}
pub(crate) fn should_idle_exit(&self, idle_secs: u64) -> bool {
self.should_idle_exit_at(idle_secs, monotonic_secs())
}
fn should_idle_exit_at(&self, idle_secs: u64, now_secs: u64) -> bool {
if idle_secs == 0 {
return false;
}
if self.inner.active_connections.load(Ordering::SeqCst) != 0 {
return false;
}
is_idle(
self.inner.last_activity_secs.load(Ordering::Relaxed),
now_secs,
idle_secs,
)
}
pub(crate) fn secs_until_idle(&self, idle_exit_secs: u64) -> u64 {
self.secs_until_idle_at(idle_exit_secs, monotonic_secs())
}
fn secs_until_idle_at(&self, idle_exit_secs: u64, now_secs: u64) -> u64 {
let elapsed =
now_secs.saturating_sub(self.inner.last_activity_secs.load(Ordering::Relaxed));
idle_exit_secs.saturating_sub(elapsed)
}
#[cfg(test)]
pub(crate) fn last_activity_secs(&self) -> u64 {
self.inner.last_activity_secs.load(Ordering::Relaxed)
}
#[cfg(test)]
pub(crate) fn set_last_activity_secs_for_test(&self, secs: u64) {
self.inner.last_activity_secs.store(secs, Ordering::Relaxed);
self.inner
.external_last_activity_secs
.store(secs, Ordering::Relaxed);
self.inner
.internal_last_activity_secs
.store(secs, Ordering::Relaxed);
}
pub fn make_registry(&self) -> HandlerRegistry {
HandlerRegistry::new(
self.inner.baseline_resources.tokenizer.clone(),
CHUNK_BUDGET_TOKENS,
)
}
pub async fn lock_corpus(&self, corpus: &str) -> OwnedMutexGuard<()> {
self.inner.corpus_locks.lock(corpus).await
}
pub fn write_lane(&self) -> Arc<Semaphore> {
self.inner.write_lane.clone()
}
pub async fn acquire_mutation_guard(
&self,
corpus: &str,
) -> Result<MutationGuard, &'static str> {
super::backpressure::acquire(self, corpus).await
}
}
async fn move_stale_store(ground_dir: &Path, found_version: u32) -> anyhow::Result<()> {
let bak = ground_dir.with_file_name(format!(
"{}.bak-v{found_version}",
ground_dir
.file_name()
.and_then(|s| s.to_str())
.unwrap_or("ground"),
));
if bak.exists() {
tokio::fs::remove_dir_all(&bak).await?;
}
tokio::fs::rename(ground_dir, &bak).await?;
let stamp_target = bak.clone();
tokio::task::spawn_blocking(move || {
std::fs::File::open(&stamp_target)?.set_modified(SystemTime::now())
})
.await??;
tracing::info!(
target: "hallouminate::daemon",
backup = %bak.display(),
"moved stale ground store aside; recoverable until pruned",
);
Ok(())
}
async fn prune_stale_backups(
ground_dir: &Path,
now: SystemTime,
max_age: Duration,
) -> anyhow::Result<()> {
let parent = ground_dir
.parent()
.filter(|p| !p.as_os_str().is_empty())
.unwrap_or_else(|| Path::new("."));
let prefix = format!(
"{}.bak-v",
ground_dir
.file_name()
.and_then(|s| s.to_str())
.unwrap_or("ground"),
);
let mut entries = match tokio::fs::read_dir(parent).await {
Ok(entries) => entries,
Err(e) if e.kind() == std::io::ErrorKind::NotFound => return Ok(()),
Err(e) => return Err(e.into()),
};
while let Some(entry) = entries.next_entry().await? {
let file_name = entry.file_name();
let Some(name) = file_name.to_str() else {
continue;
};
let Some(suffix) = name.strip_prefix(&prefix) else {
continue;
};
if suffix.is_empty() || !suffix.chars().all(|c| c.is_ascii_digit()) {
continue;
}
let metadata = match entry.metadata().await {
Ok(metadata) => metadata,
Err(e) => {
tracing::warn!(
target: "hallouminate::daemon",
entry = %name,
error = %e,
"skipping stale-backup entry: failed to read metadata",
);
continue;
}
};
if !metadata.is_dir() {
continue;
}
let modified = match metadata.modified() {
Ok(modified) => modified,
Err(e) => {
tracing::warn!(
target: "hallouminate::daemon",
entry = %name,
error = %e,
"skipping stale-backup entry: failed to read modified time",
);
continue;
}
};
let age = now.duration_since(modified).unwrap_or(Duration::ZERO);
if age < max_age {
continue;
}
let path = entry.path();
if let Err(e) = tokio::fs::remove_dir_all(&path).await {
tracing::warn!(
target: "hallouminate::daemon",
entry = %name,
error = %e,
"failed to remove stale ground store backup",
);
continue;
}
tracing::info!(
target: "hallouminate::daemon",
backup = %path.display(),
age_days = age.as_secs() / 86_400,
"pruned stale ground store backup",
);
}
Ok(())
}
pub struct CrossencoderGuard {
guard: OwnedMutexGuard<HashMap<String, FastembedCrossencoder>>,
key: String,
last_use_secs: Arc<AtomicU64>,
}
impl std::ops::Deref for CrossencoderGuard {
type Target = FastembedCrossencoder;
fn deref(&self) -> &FastembedCrossencoder {
self.guard.get(&self.key).expect("crossencoder loaded")
}
}
impl std::ops::DerefMut for CrossencoderGuard {
fn deref_mut(&mut self) -> &mut FastembedCrossencoder {
self.guard.get_mut(&self.key).expect("crossencoder loaded")
}
}
impl Drop for CrossencoderGuard {
fn drop(&mut self) {
self.last_use_secs
.store(monotonic_secs(), Ordering::Relaxed);
}
}
impl Crossencoder for CrossencoderGuard {
fn rerank(
&mut self,
query: &str,
hits: &mut [SearchHit],
) -> hallouminate_domain::common::Result<()> {
(**self).rerank(query, hits)
}
}
impl std::fmt::Debug for DaemonState {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
f.debug_struct("DaemonState")
.field("ground_dir", &self.inner.baseline_resources.ground_dir)
.field("model", &self.inner.baseline.embeddings.model)
.finish()
}
}
#[cfg(test)]
mod tests {
use super::super::maintenance::{MaintenanceTick, jittered_sleep_secs};
use super::*;
use hallouminate_adapters::{EMBEDDING_DIM, EmbedRole, MaintenanceStats};
use std::fmt;
use tracing::Subscriber;
use tracing::field::{Field, Visit};
use tracing_subscriber::layer::{Context, SubscriberExt};
use tracing_subscriber::{Layer, Registry};
#[derive(Clone, Debug, Default)]
struct CapturedEvent {
strings: HashMap<String, String>,
numbers: HashMap<String, u64>,
}
#[derive(Clone, Default)]
struct EventCapture(Arc<std::sync::Mutex<Vec<CapturedEvent>>>);
impl EventCapture {
fn maintenance_events(&self) -> Vec<CapturedEvent> {
let events = self.0.lock().expect("capture lock");
let mut maintenance = Vec::new();
for event in events.iter() {
if event.strings.contains_key("maintenance_event") {
maintenance.push(event.clone());
}
}
maintenance
}
}
impl<S> Layer<S> for EventCapture
where
S: Subscriber,
{
fn on_event(&self, event: &tracing::Event<'_>, _ctx: Context<'_, S>) {
let mut captured = CapturedEvent::default();
event.record(&mut captured);
self.0.lock().expect("capture lock").push(captured);
}
}
impl Visit for CapturedEvent {
fn record_u64(&mut self, field: &Field, value: u64) {
self.numbers.insert(field.name().to_owned(), value);
}
fn record_i64(&mut self, field: &Field, value: i64) {
let value = u64::try_from(value).expect("maintenance numeric fields are non-negative");
self.numbers.insert(field.name().to_owned(), value);
}
fn record_str(&mut self, field: &Field, value: &str) {
self.strings
.insert(field.name().to_owned(), value.to_owned());
}
fn record_debug(&mut self, field: &Field, value: &dyn fmt::Debug) {
self.strings
.insert(field.name().to_owned(), format!("{value:?}"));
}
}
fn event_stages(events: &[CapturedEvent]) -> Vec<&str> {
let mut stages = Vec::new();
for event in events {
stages.push(
event
.strings
.get("maintenance_event")
.expect("maintenance event stage")
.as_str(),
);
}
stages
}
fn assert_correlated(events: &[CapturedEvent]) {
let id = events[0]
.numbers
.get("maintenance_id")
.expect("first maintenance id");
for event in events {
assert_eq!(
event.numbers.get("maintenance_id"),
Some(id),
"every lifecycle event must carry the same maintenance id",
);
}
}
fn assert_terminal_durations(event: &CapturedEvent) {
for field in ["queue_wait_ms", "maintenance_ms", "total_ms"] {
assert!(
event.numbers.contains_key(field),
"terminal event must carry numeric {field}: {event:?}",
);
}
}
#[tokio::test]
async fn baseline_returns_the_configured_config() {
let tmp = tempfile::tempdir().expect("tempdir");
let mut cfg = Config::default();
cfg.storage.ground_dir = tmp.path().to_string_lossy().into_owned();
let expected_model = cfg.embeddings.model.clone();
let state = DaemonState::open(cfg, None)
.await
.expect("open daemon state");
assert_eq!(state.baseline().embeddings.model, expected_model);
}
#[tokio::test]
async fn should_idle_exit_is_false_when_disabled() {
let tmp = tempfile::tempdir().expect("tempdir");
let mut cfg = Config::default();
cfg.embeddings.enabled = false;
cfg.storage.ground_dir = tmp.path().to_string_lossy().into_owned();
let state = DaemonState::open(cfg, None).await.expect("open");
assert!(
!state.should_idle_exit_at(0, u64::MAX),
"idle_secs=0 disables idle-exit; must never fire",
);
}
#[tokio::test]
async fn should_idle_exit_is_false_when_recently_active() {
let tmp = tempfile::tempdir().expect("tempdir");
let mut cfg = Config::default();
cfg.embeddings.enabled = false;
cfg.storage.ground_dir = tmp.path().to_string_lossy().into_owned();
let state = DaemonState::open(cfg, None).await.expect("open");
let last = state.last_activity_secs();
assert!(
!state.should_idle_exit_at(300, last + 1),
"1 s elapsed < 300 s idle; must not exit",
);
}
#[tokio::test]
async fn should_idle_exit_is_true_when_idle_and_no_connections() {
let tmp = tempfile::tempdir().expect("tempdir");
let mut cfg = Config::default();
cfg.embeddings.enabled = false;
cfg.storage.ground_dir = tmp.path().to_string_lossy().into_owned();
let state = DaemonState::open(cfg, None).await.expect("open");
let last = state.last_activity_secs();
assert!(
state.should_idle_exit_at(300, last + 300),
"elapsed == idle_secs (>= threshold); must exit",
);
assert!(
state.should_idle_exit_at(300, last + 301),
"elapsed > idle_secs; must exit",
);
}
#[tokio::test]
async fn should_idle_exit_is_false_one_second_below_threshold() {
let tmp = tempfile::tempdir().expect("tempdir");
let mut cfg = Config::default();
cfg.embeddings.enabled = false;
cfg.storage.ground_dir = tmp.path().to_string_lossy().into_owned();
let state = DaemonState::open(cfg, None).await.expect("open");
let last = state.last_activity_secs();
assert!(
!state.should_idle_exit_at(300, last + 299),
"elapsed = idle_secs - 1 (< threshold); must not exit",
);
}
#[tokio::test]
async fn secs_until_idle_counts_down_to_the_deadline() {
let tmp = tempfile::tempdir().expect("tempdir");
let mut cfg = Config::default();
cfg.embeddings.enabled = false;
cfg.storage.ground_dir = tmp.path().to_string_lossy().into_owned();
let state = DaemonState::open(cfg, None).await.expect("open");
let last = state.last_activity_secs();
assert_eq!(
state.secs_until_idle_at(300, last),
300,
"no time elapsed since activity; the full window remains",
);
assert_eq!(
state.secs_until_idle_at(300, last + 100),
200,
"100 s elapsed of a 300 s window; 200 s remain",
);
assert_eq!(
state.secs_until_idle_at(300, last + 300),
0,
"window exactly elapsed; deadline reached",
);
assert_eq!(
state.secs_until_idle_at(300, last + 10_000),
0,
"well past the window; saturates to zero, never underflows",
);
}
#[tokio::test]
async fn active_connection_defers_idle_exit_even_when_clock_is_idle() {
let tmp = tempfile::tempdir().expect("tempdir");
let mut cfg = Config::default();
cfg.embeddings.enabled = false;
cfg.storage.ground_dir = tmp.path().to_string_lossy().into_owned();
let state = DaemonState::open(cfg, None).await.expect("open");
let last = state.last_activity_secs();
let guard = state.enter_connection(WorkClass::External);
assert!(
!state.should_idle_exit_at(300, last + 10_000),
"an active connection must defer idle-exit even when long idle",
);
drop(guard);
assert!(
state.should_idle_exit_at(300, last + 10_000),
"once the connection count returns to zero, idle-exit fires",
);
}
#[tokio::test]
async fn internal_connection_defers_idle_exit_even_when_clock_is_idle() {
let state = test_state().await;
let last = state.last_activity_secs();
let guard = state.enter_connection(WorkClass::Internal);
assert!(
!state.should_idle_exit_at(300, last + 10_000),
"an active Internal connection must defer idle-exit even when long idle",
);
drop(guard);
assert!(
state.should_idle_exit_at(300, last + 10_000),
"once the Internal guard drops, idle-exit fires",
);
}
#[tokio::test]
async fn internal_touch_activity_defers_idle_exit() {
let state = test_state().await;
state.set_last_activity_secs_for_test(1);
assert!(
state.should_idle_exit_at(300, 1000),
"clock stale at 1 s; now=1000 is well past idle",
);
state.touch_activity(WorkClass::Internal);
assert!(
!state.should_idle_exit_at(300, state.last_activity_secs() + 1),
"Internal activity must stamp the idle clock (both classes gate idle-exit)",
);
}
#[tokio::test]
async fn maintenance_tick_emits_correlated_structured_lifecycle() {
let capture = EventCapture::default();
let subscriber = Registry::default().with(capture.clone());
let _guard = tracing::subscriber::set_default(subscriber);
let tmp = tempfile::tempdir().expect("tempdir");
let mut cfg = Config::default();
cfg.embeddings.enabled = false;
cfg.storage.ground_dir = tmp.path().to_string_lossy().into_owned();
let state = DaemonState::open(cfg, None).await.expect("open");
state.set_last_activity_secs_for_test(u64::MAX / 2);
let idle_clock_before_tick = state.last_activity_secs();
let tick = state.run_maintenance_tick(true).await;
assert_eq!(tick, MaintenanceTick::Continue);
assert_eq!(
state.last_activity_secs(),
idle_clock_before_tick,
"maintenance must not stamp the idle-activity clock (ADR-002)",
);
let events = capture.maintenance_events();
assert_eq!(
event_stages(&events),
[
"started",
"write_lane_acquired",
"compaction_started",
"compaction_finished",
"prune_started",
"prune_finished",
"finished",
],
);
assert_correlated(&events);
let terminal = events.last().expect("terminal maintenance event");
assert_eq!(
terminal.strings.get("outcome").map(String::as_str),
Some("success"),
);
assert_terminal_durations(terminal);
for field in [
"fragments_removed",
"fragments_added",
"old_versions_pruned",
] {
assert!(
terminal.numbers.contains_key(field),
"success event must retain numeric {field}: {terminal:?}",
);
}
for field in ["fragments_removed", "fragments_added"] {
assert!(
events[3].numbers.contains_key(field),
"compaction event must retain numeric {field}: {:?}",
events[3],
);
}
assert!(
events[5].numbers.contains_key("old_versions_pruned"),
"prune event must retain numeric old_versions_pruned: {:?}",
events[5],
);
assert!(
!terminal.numbers.contains_key("bytes_read")
&& !terminal.numbers.contains_key("bytes_written"),
"unmeasured byte counters must not be emitted",
);
}
#[tokio::test]
async fn maintenance_tick_emits_failure_outcome() {
let capture = EventCapture::default();
let subscriber = Registry::default().with(capture.clone());
let _guard = tracing::subscriber::set_default(subscriber);
let state = test_state().await;
let tick = state
.run_maintenance_tick_with(true, |_| async {
Err(HallouminateError::Config(
"forced maintenance failure".to_owned(),
))
})
.await;
assert_eq!(tick, MaintenanceTick::Continue);
let events = capture.maintenance_events();
assert_eq!(
event_stages(&events),
["started", "write_lane_acquired", "finished"],
);
assert_correlated(&events);
let terminal = events.last().expect("terminal maintenance event");
assert_eq!(
terminal.strings.get("outcome").map(String::as_str),
Some("failure"),
);
assert_terminal_durations(terminal);
assert!(terminal.strings.contains_key("error"));
}
#[tokio::test]
async fn maintenance_tick_emits_cancellation_outcome_when_aborted_while_queued() {
let capture = EventCapture::default();
let subscriber = Registry::default().with(capture.clone());
let _guard = tracing::subscriber::set_default(subscriber);
let state = test_state().await;
let lane = state.write_lane();
let _permit = lane.acquire().await.expect("write lane");
let tick_state = state.clone();
let task = tokio::spawn(async move { tick_state.run_maintenance_tick(true).await });
for _ in 0..100 {
let events = capture.maintenance_events();
if event_stages(&events) == ["started"] {
break;
}
tokio::task::yield_now().await;
}
task.abort();
let join = task.await;
assert!(join.expect_err("aborted task").is_cancelled());
let events = capture.maintenance_events();
assert_eq!(event_stages(&events), ["started", "finished"]);
assert_correlated(&events);
let terminal = events.last().expect("terminal maintenance event");
assert_eq!(
terminal.strings.get("outcome").map(String::as_str),
Some("cancelled"),
);
assert_terminal_durations(terminal);
}
#[tokio::test]
async fn maintenance_tick_emits_shutdown_outcome_when_cancelled_while_queued() {
let capture = EventCapture::default();
let subscriber = Registry::default().with(capture.clone());
let _guard = tracing::subscriber::set_default(subscriber);
let state = test_state().await;
let lane = state.write_lane();
let _permit = lane.acquire().await.expect("write lane");
let tick_state = state.clone();
let task = tokio::spawn(async move { tick_state.run_maintenance_tick(true).await });
for _ in 0..100 {
let events = capture.maintenance_events();
if event_stages(&events) == ["started"] {
break;
}
tokio::task::yield_now().await;
}
state.shutdown_token().cancel();
let tick = task.await.expect("maintenance task");
assert_eq!(tick, MaintenanceTick::Stop);
let events = capture.maintenance_events();
assert_eq!(event_stages(&events), ["started", "finished"]);
assert_correlated(&events);
let terminal = events.last().expect("terminal maintenance event");
assert_eq!(
terminal.strings.get("outcome").map(String::as_str),
Some("shutdown"),
);
assert_terminal_durations(terminal);
}
#[tokio::test]
async fn maintenance_tick_emits_shutdown_outcome_when_cancelled_while_running() {
let capture = EventCapture::default();
let subscriber = Registry::default().with(capture.clone());
let _guard = tracing::subscriber::set_default(subscriber);
let state = test_state().await;
let lane = state.write_lane();
let (started_tx, started_rx) = tokio::sync::oneshot::channel();
let (finish_tx, finish_rx) = tokio::sync::oneshot::channel();
let tick_state = state.clone();
let task = tokio::spawn(async move {
tick_state
.run_maintenance_tick_with(true, |_| async move {
started_tx.send(()).expect("maintenance started");
finish_rx.await.expect("finish maintenance");
Ok(MaintenanceStats {
fragments_removed: None,
fragments_added: None,
old_versions_pruned: None,
})
})
.await
});
started_rx.await.expect("maintenance start signal");
state.shutdown_token().cancel();
tokio::task::yield_now().await;
assert!(
!task.is_finished(),
"shutdown must not cancel in-flight maintenance",
);
assert!(
lane.try_acquire().is_err(),
"write lane must remain held until in-flight maintenance completes",
);
finish_tx.send(()).expect("finish signal");
let tick = task.await.expect("maintenance task");
assert_eq!(tick, MaintenanceTick::Stop);
let events = capture.maintenance_events();
assert_eq!(
event_stages(&events),
["started", "write_lane_acquired", "finished"],
);
assert_correlated(&events);
let terminal = events.last().expect("terminal maintenance event");
assert_eq!(
terminal.strings.get("outcome").map(String::as_str),
Some("shutdown"),
);
assert_terminal_durations(terminal);
}
async fn test_state() -> DaemonState {
let tmp = tempfile::tempdir().expect("tempdir");
let mut cfg = Config::default();
cfg.embeddings.enabled = false;
cfg.storage.ground_dir = tmp.path().to_string_lossy().into_owned();
DaemonState::open(cfg, None).await.expect("open")
}
#[tokio::test]
async fn touch_activity_advances_the_idle_clock() {
let tmp = tempfile::tempdir().expect("tempdir");
let mut cfg = Config::default();
cfg.embeddings.enabled = false;
cfg.storage.ground_dir = tmp.path().to_string_lossy().into_owned();
let state = DaemonState::open(cfg, None).await.expect("open");
state.inner.last_activity_secs.store(1, Ordering::Relaxed);
assert!(
state.should_idle_exit_at(300, 1000),
"clock stale at 1 s; now=1000 is well past idle",
);
state.touch_activity(WorkClass::External);
assert!(
!state.should_idle_exit_at(300, state.last_activity_secs() + 1),
"touch_activity must reset the clock so a fresh now is not idle",
);
}
#[tokio::test]
async fn enter_connection_and_touch_activity_track_per_class_state() {
let state = test_state().await;
let baseline = state.last_activity_secs();
let ext = state.enter_connection(WorkClass::External);
assert_eq!(
state
.inner
.active_external_connections
.load(Ordering::SeqCst),
1
);
assert_eq!(
state
.inner
.active_internal_connections
.load(Ordering::SeqCst),
0
);
assert_eq!(
state.inner.active_connections.load(Ordering::SeqCst),
1,
"aggregate must still count both classes",
);
let int = state.enter_connection(WorkClass::Internal);
assert_eq!(
state
.inner
.active_internal_connections
.load(Ordering::SeqCst),
1
);
assert_eq!(state.inner.active_connections.load(Ordering::SeqCst), 2);
state.touch_activity(WorkClass::External);
assert!(
state
.inner
.external_last_activity_secs
.load(Ordering::Relaxed)
>= baseline
);
assert_eq!(
state
.inner
.internal_last_activity_secs
.load(Ordering::Relaxed),
baseline,
"External touch must not stamp the Internal clock",
);
state.touch_activity(WorkClass::Internal);
assert!(
state
.inner
.internal_last_activity_secs
.load(Ordering::Relaxed)
>= baseline
);
drop(ext);
assert_eq!(
state
.inner
.active_external_connections
.load(Ordering::SeqCst),
0
);
assert_eq!(
state.inner.active_connections.load(Ordering::SeqCst),
1,
"aggregate reflects only the dropped guard",
);
drop(int);
assert_eq!(
state
.inner
.active_internal_connections
.load(Ordering::SeqCst),
0
);
assert_eq!(state.inner.active_connections.load(Ordering::SeqCst), 0);
assert!(!state.should_idle_exit_at(300, baseline + 1));
}
#[test]
fn every_enter_connection_call_site_declares_a_work_class() {
fn rs_files(dir: &Path, out: &mut Vec<PathBuf>) {
for entry in std::fs::read_dir(dir).expect("read_dir under src") {
let path = entry.expect("dir entry").path();
if path.is_dir() {
rs_files(&path, out);
} else if path.extension().is_some_and(|e| e == "rs") {
out.push(path);
}
}
}
let src = Path::new(env!("CARGO_MANIFEST_DIR")).join("src");
let mut files = Vec::new();
rs_files(&src, &mut files);
let needle = ".enter_connection";
let mut call_sites = 0usize;
for path in files {
let text =
std::fs::read_to_string(&path).unwrap_or_else(|e| panic!("read {path:?}: {e}"));
let mut start = 0;
while let Some(found) = text[start..].find(needle) {
let after = &text[start + found + needle.len()..];
start += found + needle.len();
let Some(args) = after.trim_start().strip_prefix('(') else {
continue;
};
call_sites += 1;
assert!(
args.trim_start().starts_with("WorkClass::"),
"{path:?}: `enter_connection` call site must pass an explicit \
`WorkClass::` literal (ADR daemon-rework-002); found: {:?}",
after.chars().take(40).collect::<String>(),
);
}
}
assert!(
call_sites >= 4,
"expected at least the four production call sites (server, watcher, \
boot catch-up, maintenance tick); scan found {call_sites} -- did the \
scan root move?",
);
}
#[tokio::test]
async fn defer_count_resets_and_increments() {
let state = test_state().await;
assert_eq!(state.defer_count(), 0);
assert_eq!(state.increment_defer_count(), 1);
assert_eq!(state.increment_defer_count(), 2);
assert_eq!(state.defer_count(), 2);
state.reset_defer_count();
assert_eq!(state.defer_count(), 0);
}
#[tokio::test]
async fn watcher_counters_track_events_and_reindex_outcomes() {
let state = test_state().await;
state.record_watcher_events(3);
state.record_watcher_reindex(false);
state.record_watcher_reindex(true);
assert_eq!(state.watcher_counters_snapshot(), (3, 2, 1));
}
#[tokio::test]
async fn ladder_trip_storage_records_and_overwrites_the_latest_action() {
let state = test_state().await;
assert_eq!(state.last_ladder_trip(), None);
state.record_ladder_trip(LadderAction::ForceMaintenance);
assert_eq!(
state.last_ladder_trip().expect("trip recorded").action,
LadderAction::ForceMaintenance,
);
state.record_ladder_trip(LadderAction::WatchdogTrip);
assert_eq!(
state.last_ladder_trip().expect("trip recorded").action,
LadderAction::WatchdogTrip,
);
}
#[test]
fn idle_clock_is_monotonic_not_wall_clock() {
let wall_clock_secs = SystemTime::now()
.duration_since(std::time::UNIX_EPOCH)
.expect("system clock after epoch")
.as_secs();
let monotonic = monotonic_secs();
assert!(
monotonic < wall_clock_secs / 2,
"idle clock must be process-relative (Instant-based), not wall-clock \
Unix seconds: monotonic={monotonic}, wall_clock={wall_clock_secs}",
);
}
#[tokio::test]
async fn crossencoder_guard_updates_last_use_on_drop() {
let last_use_secs = Arc::new(AtomicU64::new(1));
let before_drop = monotonic_secs();
drop(CrossencoderGuard {
guard: Arc::new(Mutex::new(HashMap::new())).lock_owned().await,
key: String::new(),
last_use_secs: Arc::clone(&last_use_secs),
});
let observed = last_use_secs.load(Ordering::Relaxed);
assert!(
observed >= before_drop,
"drop should stamp crossencoder use at or after guard lifetime start: observed {observed}, before {before_drop}",
);
}
#[tokio::test]
async fn prune_stale_backups_removes_dirs_at_or_past_max_age() {
let tmp = tempfile::tempdir().expect("tempdir");
let ground_dir = tmp.path().join("ground");
tokio::fs::create_dir_all(&ground_dir)
.await
.expect("create ground dir");
let stale = tmp.path().join("ground.bak-v1");
let unrelated = tmp.path().join("other-dir");
tokio::fs::create_dir_all(&stale)
.await
.expect("create stale backup");
tokio::fs::create_dir_all(&unrelated)
.await
.expect("create unrelated dir");
let max_age = Duration::from_secs(1);
let now = SystemTime::now() + Duration::from_secs(10);
prune_stale_backups(&ground_dir, now, max_age)
.await
.expect("prune");
assert!(!stale.exists(), "backup past max_age must be pruned");
assert!(
unrelated.exists(),
"dirs that don't match the backup prefix must be left alone",
);
}
#[tokio::test]
async fn prune_stale_backups_keeps_dirs_within_max_age() {
let tmp = tempfile::tempdir().expect("tempdir");
let ground_dir = tmp.path().join("ground");
tokio::fs::create_dir_all(&ground_dir)
.await
.expect("create ground dir");
let fresh = tmp.path().join("ground.bak-v2");
tokio::fs::create_dir_all(&fresh)
.await
.expect("create fresh backup");
prune_stale_backups(&ground_dir, SystemTime::now(), STALE_BACKUP_MAX_AGE)
.await
.expect("prune");
assert!(fresh.exists(), "backup younger than max_age must survive");
}
#[tokio::test]
async fn prune_stale_backups_keeps_non_numeric_suffix_even_if_stale() {
let tmp = tempfile::tempdir().expect("tempdir");
let ground_dir = tmp.path().join("ground");
tokio::fs::create_dir_all(&ground_dir)
.await
.expect("create ground dir");
let lookalike = tmp.path().join("ground.bak-vault");
tokio::fs::create_dir_all(&lookalike)
.await
.expect("create lookalike dir");
let max_age = Duration::from_secs(1);
let now = SystemTime::now() + Duration::from_secs(31 * 24 * 60 * 60);
prune_stale_backups(&ground_dir, now, max_age)
.await
.expect("prune");
assert!(
lookalike.exists(),
"non-numeric-suffix dir must survive even when older than max_age",
);
}
#[tokio::test]
async fn move_stale_store_stamps_backup_mtime_to_now() {
let tmp = tempfile::tempdir().expect("tempdir");
let ground_dir = tmp.path().join("ground");
tokio::fs::create_dir_all(&ground_dir)
.await
.expect("create ground dir");
let old_mtime = SystemTime::now() - Duration::from_secs(60 * 24 * 60 * 60);
let dir = ground_dir.clone();
tokio::task::spawn_blocking(move || std::fs::File::open(&dir)?.set_modified(old_mtime))
.await
.expect("join")
.expect("backdate ground dir mtime");
move_stale_store(&ground_dir, 1).await.expect("move");
let bak = tmp.path().join("ground.bak-v1");
prune_stale_backups(&ground_dir, SystemTime::now(), STALE_BACKUP_MAX_AGE)
.await
.expect("prune");
assert!(
bak.exists(),
"backup just moved aside must survive its own boot's prune, \
even though the source dir's mtime was 60d old",
);
}
#[tokio::test]
async fn resources_for_keys_on_ground_dir() {
let tmp_a = tempfile::tempdir().expect("tempdir a");
let tmp_b = tempfile::tempdir().expect("tempdir b");
let mut cfg_a = Config::default();
cfg_a.embeddings.enabled = false;
cfg_a.storage.ground_dir = tmp_a.path().to_string_lossy().into_owned();
let mut cfg_b = cfg_a.clone();
cfg_b.storage.ground_dir = tmp_b.path().to_string_lossy().into_owned();
let state = DaemonState::open(cfg_a.clone(), None)
.await
.expect("open daemon state");
let res_a1 = state
.resources_for(&cfg_a)
.await
.expect("resources_for cfg_a (first call)");
let res_a2 = state
.resources_for(&cfg_a)
.await
.expect("resources_for cfg_a (second call)");
assert!(
Arc::ptr_eq(&res_a1, &res_a2),
"same config must resolve the same cached RequestResources Arc, \
not rebuild or reopen a second store",
);
let res_b = state
.resources_for(&cfg_b)
.await
.expect("resources_for cfg_b");
assert!(
!Arc::ptr_eq(&res_a1, &res_b),
"a different ground_dir must key a distinct RequestResources entry",
);
assert_eq!(
res_a1.ground_dir,
tmp_a.path(),
"cfg_a's resources must be rooted at tmp_a's ground_dir",
);
assert_eq!(
res_b.ground_dir,
tmp_b.path(),
"cfg_b's resources must be rooted at tmp_b's ground_dir, not cfg_a's",
);
assert!(
tmp_b.path().join("meta.toml").exists(),
"resources_for must open (and initialize) a store at the new \
ground_dir on first use",
);
}
#[tokio::test]
async fn resources_for_builds_once_under_concurrent_same_key_calls() {
let tmp = tempfile::tempdir().expect("tempdir");
let mut cfg = Config::default();
cfg.embeddings.enabled = false;
cfg.storage.ground_dir = tmp.path().to_string_lossy().into_owned();
let state = DaemonState::open(cfg.clone(), None)
.await
.expect("open daemon state");
let tmp_race = tempfile::tempdir().expect("tempdir race");
cfg.storage.ground_dir = tmp_race.path().to_string_lossy().into_owned();
let mut handles = Vec::new();
for _ in 0..16 {
let state = state.clone();
let cfg = cfg.clone();
handles.push(tokio::spawn(async move {
state.resources_for(&cfg).await.expect("resources_for")
}));
}
let mut resolved = Vec::new();
for h in handles {
resolved.push(h.await.expect("join resources_for task"));
}
let first = &resolved[0];
for (i, res) in resolved.iter().enumerate() {
assert!(
Arc::ptr_eq(first, res),
"racer {i} resolved a different RequestResources Arc — the \
store was opened more than once for one ground dir",
);
}
}
struct ZeroEmbedder;
impl EmbedBatch for ZeroEmbedder {
fn embed_batch(
&mut self,
texts: &[String],
_role: EmbedRole,
) -> hallouminate_domain::common::Result<Vec<[f32; EMBEDDING_DIM]>> {
Ok(vec![[0.0; EMBEDDING_DIM]; texts.len()])
}
}
#[tokio::test]
async fn cached_enabled_resource_retries_transient_embedder_failure() {
let baseline_dir = tempfile::tempdir().expect("baseline tempdir");
let mut baseline = Config::default();
baseline.embeddings.enabled = false;
baseline.storage.ground_dir = baseline_dir.path().to_string_lossy().into_owned();
let state = DaemonState::open(baseline, None)
.await
.expect("open daemon state");
let retry_dir = tempfile::tempdir().expect("retry tempdir");
let mut cfg = Config::default();
cfg.embeddings.enabled = true;
cfg.storage.ground_dir = retry_dir.path().to_string_lossy().into_owned();
let store = LanceStore::open_or_create(
retry_dir.path(),
&cfg.embeddings.model,
cfg.embeddings.quantized,
true,
None,
)
.await
.expect("open enabled store without embedder");
let cached = Arc::new(RequestResources {
store: Arc::new(store),
tokenizer: load_tokenizer(&cfg.embeddings.model).expect("tokenizer"),
embeddings_enabled: true,
ground_dir: retry_dir.path().to_path_buf(),
});
state
.inner
.resources
.lock()
.await
.insert(ResourceKey::from_config(&cfg), Arc::clone(&cached));
let first = state
.resources_for_with_initializer(&cfg, |_, _, _| async {
Err(anyhow::anyhow!("transient initialization failure"))
})
.await
.expect("cached resources remain usable after retry failure");
assert!(Arc::ptr_eq(&cached, &first));
assert!(!first.store.embedder_available());
let retried = state
.resources_for_with_initializer(&cfg, |_, _, _| async {
Ok(Box::new(ZeroEmbedder) as Box<dyn EmbedBatch>)
})
.await
.expect("second request retries initialization");
assert!(Arc::ptr_eq(&cached, &retried));
assert!(retried.store.embedder_available());
let normal = state
.resources_for_with_initializer(&cfg, |_, _, _| async {
panic!("ready cache hit must not initialize again")
})
.await
.expect("normal request reuses repaired resource");
assert!(Arc::ptr_eq(&retried, &normal));
}
struct TestProbe(std::sync::atomic::AtomicBool);
impl TestProbe {
fn new(v: bool) -> Self {
Self(std::sync::atomic::AtomicBool::new(v))
}
fn set(&self, v: bool) {
self.0.store(v, Ordering::SeqCst);
}
}
impl IoPressureProbe for TestProbe {
fn elevated(&self) -> bool {
self.0.load(Ordering::SeqCst)
}
}
#[tokio::test]
async fn defer_reason_idle_and_clear_defers_nothing() {
let state = test_state().await;
state.set_last_activity_secs_for_test(0);
let probe = TestProbe::new(false);
assert_eq!(state.maintenance_defer_reason_at(&probe, 1000), None);
}
#[tokio::test]
async fn defer_reason_external_connection_defers_as_active() {
let state = test_state().await;
state.set_last_activity_secs_for_test(0);
let _conn = state.enter_connection(WorkClass::External);
let probe = TestProbe::new(false);
assert_eq!(
state.maintenance_defer_reason_at(&probe, 1000),
Some(DeferReason::Active),
);
}
#[tokio::test]
async fn defer_reason_recent_activity_defers_as_active() {
let state = test_state().await;
state.set_last_activity_secs_for_test(970);
let probe = TestProbe::new(false);
assert_eq!(
state.maintenance_defer_reason_at(&probe, 1000),
Some(DeferReason::Active),
);
}
#[tokio::test]
async fn defer_reason_external_activity_defers_as_active() {
let state = test_state().await;
state.set_last_activity_secs_for_test(0);
state
.inner
.external_last_activity_secs
.store(970, Ordering::Relaxed);
state.inner.last_activity_secs.store(970, Ordering::Relaxed);
let probe = TestProbe::new(false);
assert_eq!(
state.maintenance_defer_reason_at(&probe, 1000),
Some(DeferReason::Active),
);
}
#[tokio::test]
async fn defer_reason_internal_connection_does_not_defer_maintenance() {
let state = test_state().await;
state.set_last_activity_secs_for_test(0);
let _conn = state.enter_connection(WorkClass::Internal);
let probe = TestProbe::new(false);
assert_eq!(state.maintenance_defer_reason_at(&probe, 1000), None);
}
#[tokio::test]
async fn defer_reason_internal_connection_with_io_pressure_defers_as_io_pressure() {
let state = test_state().await;
state.set_last_activity_secs_for_test(0);
let _conn = state.enter_connection(WorkClass::Internal);
let probe = TestProbe::new(true);
assert_eq!(
state.maintenance_defer_reason_at(&probe, 1000),
Some(DeferReason::IoPressure),
);
}
#[tokio::test]
async fn defer_reason_internal_activity_does_not_defer_maintenance() {
let state = test_state().await;
state.set_last_activity_secs_for_test(0);
state
.inner
.internal_last_activity_secs
.store(970, Ordering::Relaxed);
state.inner.last_activity_secs.store(970, Ordering::Relaxed);
let probe = TestProbe::new(false);
assert_eq!(state.maintenance_defer_reason_at(&probe, 1000), None);
}
#[tokio::test]
async fn internal_housekeeping_alone_is_maintenance_eligible_but_defers_idle_exit() {
let state = test_state().await;
state.set_last_activity_secs_for_test(0);
let conn = state.enter_connection(WorkClass::Internal);
state
.inner
.internal_last_activity_secs
.store(970, Ordering::Relaxed);
state.inner.last_activity_secs.store(970, Ordering::Relaxed);
let probe = TestProbe::new(false);
assert_eq!(
state.maintenance_defer_reason_at(&probe, 1000),
None,
"internal-only housekeeping must leave maintenance eligible",
);
assert!(
!state.should_idle_exit_at(300, 100_000),
"an in-flight Internal guard must defer idle-exit",
);
drop(conn);
assert!(
!state.should_idle_exit_at(300, 1000),
"recent Internal activity (clock at 970) must defer idle-exit at now=1000",
);
assert!(
state.should_idle_exit_at(300, 970 + 300),
"once housekeeping completes and the window elapses, idle-exit fires",
);
}
#[tokio::test]
async fn defer_reason_io_pressure_defers_and_clears() {
let state = test_state().await;
state.set_last_activity_secs_for_test(0);
let probe = TestProbe::new(true);
assert_eq!(
state.maintenance_defer_reason_at(&probe, 1000),
Some(DeferReason::IoPressure),
);
probe.set(false);
assert_eq!(state.maintenance_defer_reason_at(&probe, 1000), None);
}
#[tokio::test]
async fn defer_reason_active_takes_priority_over_io_pressure() {
let state = test_state().await;
state.set_last_activity_secs_for_test(0);
let _conn = state.enter_connection(WorkClass::External);
let probe = TestProbe::new(true);
assert_eq!(
state.maintenance_defer_reason_at(&probe, 1000),
Some(DeferReason::Active),
);
}
#[tokio::test]
async fn defer_reason_idle_boundary_59s_is_still_active() {
let state = test_state().await;
state.set_last_activity_secs_for_test(0);
let probe = TestProbe::new(false);
assert_eq!(
state.maintenance_defer_reason_at(&probe, 59),
Some(DeferReason::Active),
"59s since last activity must not count as idle",
);
}
#[tokio::test]
async fn defer_reason_idle_boundary_60s_is_idle() {
let state = test_state().await;
state.set_last_activity_secs_for_test(0);
let probe = TestProbe::new(false);
assert_eq!(
state.maintenance_defer_reason_at(&probe, 60),
None,
"60s since last activity must count as idle (inclusive boundary)",
);
}
#[tokio::test]
async fn defer_reason_idle_boundary_61s_is_idle() {
let state = test_state().await;
state.set_last_activity_secs_for_test(0);
let probe = TestProbe::new(false);
assert_eq!(state.maintenance_defer_reason_at(&probe, 61), None);
}
#[tokio::test(start_paused = true)]
async fn maintenance_loop_warns_after_eleven_consecutive_defers() {
let _coord = crate::debt::OBSERVED_HARD_COORD.read().await;
let state = test_state().await;
let _conn = state.enter_connection(WorkClass::External);
let capture = EventCapture::default();
let subscriber = Registry::default().with(capture.clone());
let _guard = tracing::subscriber::set_default(subscriber);
let cancel = CancellationToken::new();
let task = tokio::spawn(maintenance_loop(
state.clone(),
cancel.clone(),
Duration::from_secs(100),
Arc::new(TestProbe::new(false)),
));
tokio::task::yield_now().await;
tokio::time::advance(Duration::from_secs(111)).await;
tokio::task::yield_now().await;
let mut saw_eleven = false;
for _ in 0..30 {
tokio::time::advance(Duration::from_secs(60)).await;
tokio::task::yield_now().await;
let warned = {
let events = capture.0.lock().expect("capture lock");
events
.iter()
.any(|e| e.numbers.get("consecutive_defers").copied() == Some(11))
};
if warned {
saw_eleven = true;
break;
}
}
assert!(
saw_eleven,
"expected a warn event with consecutive_defers == 11 within the advance budget"
);
cancel.cancel();
task.await.expect("maintenance_loop task");
}
#[tokio::test(start_paused = true)]
async fn maintenance_loop_shuts_down_promptly_during_interval_sleep() {
let state = test_state().await;
let cancel = CancellationToken::new();
let task = tokio::spawn(maintenance_loop(
state.clone(),
cancel.clone(),
Duration::from_secs(100),
Arc::new(TestProbe::new(false)),
));
cancel.cancel();
task.await.expect("maintenance_loop task");
}
#[tokio::test(start_paused = true)]
async fn maintenance_loop_shuts_down_promptly_during_recheck_sleep() {
let state = test_state().await;
let _conn = state.enter_connection(WorkClass::External);
let cancel = CancellationToken::new();
let task = tokio::spawn(maintenance_loop(
state.clone(),
cancel.clone(),
Duration::from_secs(100),
Arc::new(TestProbe::new(false)),
));
tokio::task::yield_now().await;
tokio::time::advance(Duration::from_secs(111)).await;
tokio::task::yield_now().await;
cancel.cancel();
task.await.expect("maintenance_loop task");
}
#[tokio::test]
async fn maintenance_interval_zero_disables_the_background_task() {
let tmp = tempfile::tempdir().expect("tempdir");
let mut cfg = Config::default();
cfg.embeddings.enabled = false;
cfg.storage.ground_dir = tmp.path().to_string_lossy().into_owned();
cfg.daemon.maintenance_interval_secs = 0;
let capture = EventCapture::default();
let subscriber = Registry::default().with(capture.clone());
let _guard = tracing::subscriber::set_default(subscriber);
let state = DaemonState::open(cfg, None).await.expect("open");
assert!(state.take_maintenance_task().await.is_none());
let events = capture.0.lock().expect("capture lock");
let disabled = events.iter().any(|e| {
e.strings
.get("message")
.map(|m| m.contains("disabled"))
.unwrap_or(false)
});
assert!(
disabled,
"expected a log message mentioning maintenance is disabled"
);
}
#[test]
fn jittered_sleep_secs_stays_within_interval_plus_ten_percent() {
let mut seen = std::collections::HashSet::new();
for _ in 0..1000 {
let jittered = jittered_sleep_secs(100);
assert!(
(100..=110).contains(&jittered),
"jittered_sleep_secs(100) = {jittered} out of [100, 110]",
);
seen.insert(jittered);
}
assert!(
seen.contains(&100) && seen.contains(&110),
"expected both closed-interval endpoints reachable across 1000 iterations, got {seen:?}",
);
}
#[test]
fn jittered_sleep_secs_zero_interval_returns_zero() {
assert_eq!(jittered_sleep_secs(0), 0);
}
#[test]
fn jittered_sleep_secs_below_ten_adds_no_jitter() {
assert_eq!(jittered_sleep_secs(5), 5);
}
}