use super::concurrency::{DreamCycleGauge, acquire_dream_permit};
use super::config::{DreamConfig, DreamStats};
use super::cycle::{
DedupOutcome, compact_pass, content_prune_pass, dedup_pass, prune_pass, refresh_closets,
};
use super::fading::detect_fading;
use super::guard::CompactionGuard;
use super::kg_compact::kg_compact_pass;
use super::recall_benchmark::run_benchmark;
use super::semantic::{SemanticPassOutcome, semantic_consolidation_pass};
use super::settled;
use crate::memory_core::embed::Embedder;
use crate::memory_core::palace::PalaceId;
use crate::memory_core::registry::PalaceRegistry;
use crate::memory_core::retrieval::PalaceHandle;
use crate::memory_core::semantic_consolidation::SemanticConsolidator;
use anyhow::{Context, Result};
use std::sync::Arc;
use std::sync::atomic::{AtomicBool, AtomicU64, Ordering};
use std::time::Duration;
use super::helpers::now_secs;
pub type AfterCycle = Arc<dyn Fn(&PalaceId, Option<&DreamStats>) + Send + Sync>;
pub struct Dreamer {
pub config: DreamConfig,
pub(super) last_activity: Arc<AtomicU64>,
pub(super) consolidator: Option<Arc<SemanticConsolidator>>,
pub(super) semantic_consolidation_disabled: AtomicBool,
pub(super) after_cycle: Option<AfterCycle>,
pub(super) embedder: Option<Arc<dyn Embedder + Send + Sync>>,
}
impl Dreamer {
pub fn new(config: DreamConfig) -> Self {
Self {
config,
last_activity: Arc::new(AtomicU64::new(now_secs())),
consolidator: None,
semantic_consolidation_disabled: AtomicBool::new(false),
after_cycle: None,
embedder: None,
}
}
pub fn with_consolidator(config: DreamConfig, consolidator: Arc<SemanticConsolidator>) -> Self {
Self {
config,
last_activity: Arc::new(AtomicU64::new(now_secs())),
consolidator: Some(consolidator),
semantic_consolidation_disabled: AtomicBool::new(false),
after_cycle: None,
embedder: None,
}
}
pub fn with_after_cycle(mut self, hook: AfterCycle) -> Self {
self.after_cycle = Some(hook);
self
}
#[cfg(test)]
pub(super) fn with_embedder(mut self, embedder: Arc<dyn Embedder + Send + Sync>) -> Self {
self.embedder = Some(embedder);
self
}
pub fn touch(&self) {
self.last_activity.store(now_secs(), Ordering::Relaxed);
}
pub fn is_idle(&self) -> bool {
let last = self.last_activity.load(Ordering::Relaxed);
now_secs().saturating_sub(last) >= self.config.idle_secs
}
pub fn is_semantic_consolidation_disabled(&self) -> bool {
self.semantic_consolidation_disabled.load(Ordering::Relaxed)
}
pub fn start_with_shutdown(
self: Arc<Self>,
registry: PalaceRegistry,
palace_id: PalaceId,
first_tick_stagger: Duration,
mut shutdown: tokio::sync::watch::Receiver<bool>,
) -> tokio::task::JoinHandle<()> {
tokio::spawn(async move {
let interval = Duration::from_secs(self.config.idle_secs.max(1));
let mut wait = interval + first_tick_stagger;
loop {
tokio::select! {
_ = tokio::time::sleep(wait) => { wait = interval; }
res = shutdown.changed() => {
if res.is_err() || *shutdown.borrow() {
tracing::info!(palace = %palace_id, "dreamer shutting down");
return;
}
}
}
if *shutdown.borrow() {
tracing::info!(palace = %palace_id, "dreamer shutting down");
return;
}
if !self.is_idle() || !registry.may_run_maintenance() {
continue;
}
let Some(handle) = registry.peek(&palace_id) else {
continue;
};
log_cycle_outcome(&palace_id, self.dream_cycle(&handle).await);
}
})
}
pub async fn dream_cycle(&self, handle: &Arc<PalaceHandle>) -> Result<DreamStats> {
let _permit = acquire_dream_permit().await;
let Some(_compaction_guard) = CompactionGuard::try_claim(handle.is_compacting.clone())
else {
tracing::info!(palace = %handle.id, "dream cycle skipped: one is already running");
return Ok(DreamStats::default());
};
let _in_flight = DreamCycleGauge::enter();
let outcome = self.run_claimed_cycle(handle).await;
if let Some(hook) = &self.after_cycle {
hook(&handle.id, outcome.as_ref().ok());
}
outcome
}
async fn run_claimed_cycle(&self, handle: &Arc<PalaceHandle>) -> Result<DreamStats> {
let started = std::time::Instant::now();
let budget = Duration::from_millis(self.config.max_cycle_ms);
let drawers_before = handle.drawers.read().len() as u64;
let fingerprint = settled::corpus_fingerprint(handle, &self.config);
let unchanged = handle
.data_dir
.as_deref()
.is_some_and(|dir| settled::is_settled(dir, &fingerprint));
if unchanged {
tracing::debug!(
palace = %handle.id,
"dream cycle: corpus unchanged since its last settled cycle; \
skipping dedup, recall benchmark and semantic consolidation"
);
}
let benchmark = self.config.recall_benchmark_enabled && !unchanged;
let recall_score_before = if benchmark {
run_benchmark(handle, self.embedder.clone()).await
} else {
None
};
let content_pruned = if self.config.content_prune_enabled {
content_prune_pass(handle, started, budget, self.config.content_prune_min_words)
.await
.context("dream content prune pass")?
} else {
0
};
let dedup = if unchanged {
DedupOutcome::default()
} else {
let embedder = self.embedder.clone();
dedup_pass(
handle,
started,
budget,
self.config.dedup_threshold,
embedder,
)
.await
.context("dream dedup pass")?
};
let pruned = prune_pass(handle, started, budget, self.config.prune_importance)
.await
.context("dream prune pass")?;
let compacted = compact_pass(handle, started, budget)
.await
.context("dream compact pass")?;
let closets_updated = refresh_closets(handle);
let semantic = if unchanged {
SemanticPassOutcome::SETTLED_NOOP
} else {
semantic_consolidation_pass(
handle,
&self.config,
self.consolidator.clone(),
&self.semantic_consolidation_disabled,
)
.await
};
let passes_complete = started.elapsed() < budget && dedup.failed == 0 && semantic.settles;
if let Err(e) = handle.flush() {
tracing::warn!("dream flush failed: {e:#}");
}
let drawers_after = handle.drawers.read().len() as u64;
let fading = detect_fading(handle, &self.config.fading);
let recall_score_after = if benchmark {
run_benchmark(handle, self.embedder.clone()).await
} else {
None
};
let mut stats = DreamStats {
merged: dedup.merged,
pruned,
closets_updated,
compacted,
content_pruned,
semantically_consolidated: semantic.consolidated,
semantic_llm_calls: semantic.llm_calls,
semantic_cache_hits: semantic.cache_hits,
duration_ms: started.elapsed().as_millis() as u64,
drawers_before,
drawers_after,
compression_ratio: 0.0, recall_score_before,
recall_score_after,
kg_bytes_reclaimed: 0,
kg_bytes_after: 0,
kg_history_rows_pruned: 0,
fading,
};
stats.update_compression_ratio();
crate::memory_core::maintenance_log::warn_removed(
&handle.id,
"dream cycle",
stats.merged + stats.pruned + stats.content_pruned,
);
if self.config.compact {
match kg_compact_pass(handle, &self.config, false).await {
Ok(report) => {
stats.kg_bytes_reclaimed = report.bytes_reclaimed();
stats.kg_bytes_after = report.bytes_after;
stats.kg_history_rows_pruned = report.history_rows_pruned;
if report.ran() {
tracing::info!(palace = %handle.id, "kg.redb: {}", report.summary());
} else {
tracing::debug!(palace = %handle.id, "kg.redb: {}", report.summary());
}
}
Err(e) => tracing::warn!(
palace = %handle.id,
"kg.redb compaction failed (non-fatal; the live file is unchanged): {e:#}"
),
}
}
if let Some(data_dir) = handle.data_dir.as_ref() {
use super::config::PersistedDreamStats;
let persisted = PersistedDreamStats {
last_run_at: chrono::Utc::now(),
stats: stats.clone(),
};
if let Err(e) = persisted.save(data_dir) {
tracing::warn!(palace = %handle.id, "persist dream_stats.json failed: {e:#}");
}
if !unchanged
&& passes_complete
&& changed_nothing(&stats)
&& let Err(e) = settled::record_settled(data_dir, &fingerprint)
{
tracing::warn!(palace = %handle.id, "persist dream_settled.json failed: {e:#}");
}
}
Ok(stats)
}
}
fn changed_nothing(stats: &DreamStats) -> bool {
[
stats.merged,
stats.pruned,
stats.content_pruned,
stats.compacted,
stats.semantically_consolidated,
]
.iter()
.all(|n| *n == 0)
}
fn log_cycle_outcome(palace_id: &PalaceId, outcome: Result<DreamStats>) {
match outcome {
Ok(stats) => tracing::info!(
palace = %palace_id,
merged = stats.merged,
pruned = stats.pruned,
content_pruned = stats.content_pruned,
compacted = stats.compacted,
closets_updated = stats.closets_updated,
semantically_consolidated = stats.semantically_consolidated,
semantic_llm_calls = stats.semantic_llm_calls,
duration_ms = stats.duration_ms,
kg_bytes_reclaimed = stats.kg_bytes_reclaimed,
kg_history_rows_pruned = stats.kg_history_rows_pruned,
drawers_before = stats.drawers_before,
drawers_after = stats.drawers_after,
compression_ratio = stats.compression_ratio,
fading = stats.fading.len(),
"dream cycle complete"
),
Err(e) => tracing::warn!(palace = %palace_id, "dream cycle failed: {e:#}"),
}
}