use super::concurrency::acquire_dream_permit_within;
use super::config::DreamConfig;
use super::helpers::char_safe_prefix;
use crate::memory_core::palace::{Drawer, RoomType};
use crate::memory_core::retrieval::PalaceHandle;
use crate::memory_core::semantic_consolidation::{
SemanticConsolidator, resolve_consolidation_provider, validate_ollama_model,
};
use crate::memory_core::timeouts;
use anyhow::Result;
use std::collections::HashSet;
use std::sync::Arc;
use std::sync::atomic::{AtomicBool, Ordering};
use uuid::Uuid;
#[derive(Debug, Clone, Copy, Default, PartialEq, Eq)]
pub struct RoomConsolidationStats {
pub summary_facts_created: usize,
pub facts_evicted: usize,
}
pub(super) fn build_consolidator_from_config(
config: &DreamConfig,
) -> Result<Option<Arc<SemanticConsolidator>>> {
if !config.semantic.enabled {
return Ok(None);
}
use crate::memory_core::semantic_consolidation::{
ConsolidationProvider, OllamaInference, OpenRouterInference,
};
let backend: Arc<dyn crate::memory_core::semantic_consolidation::Inference> =
match resolve_consolidation_provider(
&config.semantic.model,
&config.openrouter_api_key,
config.local_model_enabled,
) {
ConsolidationProvider::Local { base_url, model } => {
validate_ollama_model(&model)?;
Arc::new(OllamaInference::new(base_url, model))
}
ConsolidationProvider::OpenRouter { api_key, model } => {
Arc::new(OpenRouterInference::new(api_key, model))
}
ConsolidationProvider::Unavailable { reason, .. } => anyhow::bail!(reason),
};
Ok(Some(Arc::new(SemanticConsolidator::new(
backend,
config.semantic.clone(),
))))
}
async fn apply_consolidation_result(
handle: &Arc<PalaceHandle>,
result: &crate::memory_core::semantic_consolidation::ConsolidationResult,
) -> (usize, Vec<Uuid>) {
let mut canonical_count = 0usize;
let mut superseded_ids: Vec<Uuid> = Vec::new();
for canonical in &result.canonical_drawers {
match handle
.remember(
canonical.content.clone(),
RoomType::General,
canonical.tags.clone(),
canonical.importance,
)
.await
{
Ok(canonical_id) => {
canonical_count += 1;
record_provenance_and_collect_superseded(
handle,
canonical_id,
&canonical.canonical_for,
&mut superseded_ids,
)
.await;
}
Err(e) => {
tracing::warn!(
content = char_safe_prefix(&canonical.content, 80),
"dream semantic: failed to add canonical drawer: {e:#}"
);
}
}
}
apply_alias_and_flag_passes(handle, result).await;
(canonical_count, superseded_ids)
}
pub(super) async fn record_provenance_and_collect_superseded(
handle: &Arc<PalaceHandle>,
canonical_id: Uuid,
originals: &[Uuid],
superseded_ids: &mut Vec<Uuid>,
) {
for &orig_id in originals {
match crate::memory_core::share::assert_superseded_by(
&handle.kg,
orig_id,
canonical_id,
"dream:semantic_consolidation",
)
.await
{
Ok(()) => superseded_ids.push(orig_id),
Err(e) => {
tracing::warn!(
orig = %orig_id,
canonical = %canonical_id,
"failed to write superseded_by triple: {e:#} — original \
retained (no eviction without recorded provenance)"
);
}
}
}
}
async fn apply_alias_and_flag_passes(
handle: &Arc<PalaceHandle>,
result: &crate::memory_core::semantic_consolidation::ConsolidationResult,
) {
for (from, to) in &result.aliases {
let triple = crate::memory_core::store::kg::Triple {
subject: from.clone(),
predicate: "alias_of".to_string(),
object: to.clone(),
valid_from: chrono::Utc::now(),
valid_to: None,
confidence: 1.0,
provenance: Some("dream:semantic_consolidation".to_string()),
};
if let Err(e) = handle.kg.assert(triple).await {
tracing::warn!(
from,
to,
"dream semantic: failed to write alias triple: {e:#}"
);
}
}
for (id, reason) in &result.flagged_ids {
tracing::info!(
palace = %handle.id,
drawer_id = %id,
reason,
"dream semantic: flagged drawer for human review (contradiction)"
);
}
}
pub(super) async fn semantic_consolidation_pass(
handle: &Arc<PalaceHandle>,
config: &DreamConfig,
injected: Option<Arc<SemanticConsolidator>>,
disabled: &AtomicBool,
) -> (usize, usize, usize) {
if !config.semantic.enabled {
tracing::debug!(
palace = %handle.id,
"skipping semantic consolidation: disabled in config"
);
return (0, 0, 0);
}
let consolidator: Arc<SemanticConsolidator> = match injected {
Some(c) => c,
None => {
if disabled.load(Ordering::Relaxed) {
return (0, 0, 0);
}
match build_consolidator_from_config(config) {
Ok(Some(c)) => c,
Ok(None) => {
tracing::debug!(
palace = %handle.id,
"skipping semantic consolidation: disabled in config"
);
return (0, 0, 0);
}
Err(e) => {
tracing::warn!(
palace = %handle.id,
model = %config.semantic.model,
"semantic consolidation disabled for this palace: {e:#}"
);
disabled.store(true, Ordering::Relaxed);
return (0, 0, 0);
}
}
}
};
let snapshot: Vec<Drawer> = handle
.drawers
.read()
.iter()
.filter(|d| !d.drawer_type.is_protected())
.cloned()
.collect();
if snapshot.is_empty() {
return (0, 0, 0);
}
let result = consolidator.consolidate(&snapshot).await;
let (canonical_count, _superseded) = apply_consolidation_result(handle, &result).await;
tracing::debug!(
palace = %handle.id,
canonical_added = canonical_count,
aliases = result.aliases.len(),
flagged = result.flagged_ids.len(),
llm_calls = result.llm_calls,
cache_hits = result.cache_hits,
"semantic consolidation phase complete"
);
(canonical_count, result.llm_calls, result.cache_hits)
}
pub async fn consolidate_scoped(
handle: &Arc<PalaceHandle>,
config: &DreamConfig,
room: Option<RoomType>,
max_age_days: i64,
injected: Option<Arc<SemanticConsolidator>>,
) -> Result<RoomConsolidationStats> {
consolidate_scoped_within(
handle,
config,
room,
max_age_days,
injected,
timeouts::dream_permit_wait_timeout(),
)
.await
}
pub(super) async fn consolidate_scoped_within(
handle: &Arc<PalaceHandle>,
config: &DreamConfig,
room: Option<RoomType>,
max_age_days: i64,
injected: Option<Arc<SemanticConsolidator>>,
permit_wait: std::time::Duration,
) -> Result<RoomConsolidationStats> {
if max_age_days <= 0 {
tracing::debug!(
palace = %handle.id,
max_age_days,
"dream_consolidate_room: non-positive age window; no-op"
);
return Ok(RoomConsolidationStats::default());
}
let _permit = acquire_dream_permit_within(permit_wait)
.await
.map_err(|busy| {
tracing::warn!(palace = %handle.id, "dream_consolidate_room: {busy}");
busy
})?;
let consolidator: Arc<SemanticConsolidator> = match injected {
Some(c) => c,
None => match build_consolidator_from_config(config) {
Ok(Some(c)) => c,
Ok(None) => {
tracing::debug!(
palace = %handle.id,
"dream_consolidate_room: inference unavailable; no-op"
);
return Ok(RoomConsolidationStats::default());
}
Err(e) => {
tracing::error!(
palace = %handle.id,
model = %config.semantic.model,
"dream_consolidate_room: semantic consolidation misconfigured: {e:#}"
);
return Err(e);
}
},
};
let cutoff = chrono::Utc::now() - chrono::Duration::days(max_age_days);
let snapshot: Vec<Drawer> = handle
.list_drawers(room, None, usize::MAX)
.into_iter()
.filter(|d| !d.drawer_type.is_protected())
.filter(|d| d.created_at <= cutoff)
.collect();
if snapshot.is_empty() {
return Ok(RoomConsolidationStats::default());
}
let result = consolidator.consolidate(&snapshot).await;
let (summary_facts_created, superseded_ids) = apply_consolidation_result(handle, &result).await;
let mut evicted = 0usize;
let mut seen: HashSet<Uuid> = HashSet::new();
for id in superseded_ids {
if !seen.insert(id) {
continue;
}
match handle.forget(id).await {
Ok(outcome) if outcome.is_deleted() => evicted += 1,
Ok(_) => tracing::warn!(?id, "dream_consolidate_room: superseded id not in palace"),
Err(e) => tracing::warn!(?id, "dream_consolidate_room: evict failed: {e:#}"),
}
}
if let Err(e) = handle.flush() {
tracing::warn!(palace = %handle.id, "dream_consolidate_room flush failed: {e:#}");
}
Ok(RoomConsolidationStats {
summary_facts_created,
facts_evicted: evicted,
})
}