use crate::memory_core::palace::Drawer;
use crate::memory_core::retrieval::{PalaceHandle, shared_embedder};
use crate::memory_core::store::vector::VectorStore;
use crate::memory_core::timeouts;
use anyhow::{Context, Result};
use std::collections::HashMap;
use std::sync::Arc;
use std::time::{Duration, SystemTime, UNIX_EPOCH};
use uuid::Uuid;
pub(crate) const STOP_WORDS: &[&str] = &[
"the", "a", "an", "is", "are", "was", "were", "be", "been", "being", "of", "in", "on", "at",
"to", "for", "with", "and", "or", "but", "not", "no", "yes", "i", "you", "he", "she", "it",
"we", "they", "this", "that", "these", "those", "as", "by", "from", "into", "over", "under",
"if", "then", "than", "so", "do", "does", "did", "have", "has", "had", "will", "would",
"shall", "should", "can", "could", "may", "might", "must", "about", "any", "all", "some",
"more", "most", "such",
];
pub fn extract_keywords(content: &str) -> Vec<String> {
let mut seen: std::collections::HashSet<String> = std::collections::HashSet::new();
let mut out: Vec<String> = Vec::new();
for raw in content.split_whitespace() {
let token: String = raw
.chars()
.filter(|c| c.is_alphanumeric())
.flat_map(|c| c.to_lowercase())
.collect();
if token.len() < 3 {
continue;
}
if STOP_WORDS.iter().any(|s| *s == token) {
continue;
}
if seen.insert(token.clone()) {
out.push(token);
}
}
out
}
pub(crate) fn is_low_quality_content(content: &str, min_words: usize) -> bool {
if crate::memory_core::filter::blocklist_match(content).is_some() {
return true;
}
let word_count = content.split_whitespace().count();
word_count < min_words
}
pub(crate) fn char_safe_prefix(s: &str, max_bytes: usize) -> &str {
&s[..s.floor_char_boundary(max_bytes)]
}
pub(crate) fn now_secs() -> u64 {
SystemTime::now()
.duration_since(UNIX_EPOCH)
.map(|d| d.as_secs())
.unwrap_or(0)
}
pub(crate) const RULING_TAG: &str = "ruling";
pub(crate) const MERGE_MAX_BYTES: usize = 4 * 1024;
const MERGE_SEPARATOR: &str = "\n\nAlso: ";
pub(crate) fn merged_drawer(survivor: &Drawer, loser: &Drawer) -> Option<Drawer> {
let mut merged = survivor.clone();
if !survivor.content().contains(loser.content()) {
let len = survivor.content().len() + MERGE_SEPARATOR.len() + loser.content().len();
if len > MERGE_MAX_BYTES {
return None;
}
merged.set_content(format!(
"{}{MERGE_SEPARATOR}{}",
survivor.content(),
loser.content()
));
}
merged.importance = merged.importance.max(loser.importance);
for tag in &loser.tags {
if !merged.tags.contains(tag) {
merged.tags.push(tag.clone());
}
}
Some(merged)
}
pub(crate) fn pick_survivor<'a>(a: &'a Drawer, b: &'a Drawer) -> Option<(&'a Drawer, &'a Drawer)> {
if a.is_tier_c() && b.is_tier_c() {
return None;
}
let rank = |d: &Drawer| (d.is_tier_c(), d.tags.iter().any(|t| t == RULING_TAG));
let order = rank(a)
.cmp(&rank(b))
.then(a.created_at.cmp(&b.created_at))
.then(a.importance.total_cmp(&b.importance));
Some(if order.is_lt() { (b, a) } else { (a, b) })
}
pub(crate) async fn persist_merge(
handle: &Arc<PalaceHandle>,
a: Uuid,
b: Uuid,
) -> Result<Option<(Uuid, Uuid)>> {
let _write_guard = timeouts::lock_with_timeout(
&handle.write_mutex,
timeouts::write_lock_timeout(),
handle.id.as_str(),
)
.await?;
let _order = handle.commit_mutex.lock().await;
let (a, b) = {
let drawers = handle.drawers.read();
let find = |id: Uuid| drawers.iter().find(|d| d.id == id).cloned();
(find(a), find(b))
};
let (Some(a), Some(b)) = (a, b) else {
return Ok(None);
};
let Some((survivor, loser)) = pick_survivor(&a, &b) else {
return Ok(None);
};
let Some(merged) = merged_drawer(survivor, loser) else {
tracing::info!(
palace = %handle.id, survivor = %survivor.id, loser = %loser.id,
"dream dedup: merge would pass {MERGE_MAX_BYTES} bytes; both drawers kept"
);
return Ok(None);
};
handle
.kg
.upsert_drawer(&merged)
.await
.with_context(|| format!("persist merged dedup survivor {}", merged.id))?;
let pair = (survivor.id, loser.id);
if let Some(row) = handle.drawers.write().iter_mut().find(|d| d.id == pair.0) {
*row = merged;
}
Ok(Some(pair))
}
pub(crate) async fn rebuild_index_from_drawers(
handle: &Arc<PalaceHandle>,
started: std::time::Instant,
budget: Duration,
) -> Result<usize> {
let _write_guard = timeouts::lock_with_timeout(
&handle.write_mutex,
timeouts::write_lock_timeout(),
handle.id.as_str(),
)
.await?;
let snapshot: Vec<Drawer> = handle.drawers.read().clone();
handle
.vector_store
.reset()
.context("reset vector index for rebuild")?;
if snapshot.is_empty() {
return Ok(0);
}
let embedder = shared_embedder()
.await
.context("acquire shared embedder for dream rebuild")?;
let mut rebuilt: usize = 0;
for drawer in snapshot.iter() {
if started.elapsed() >= budget {
break;
}
let body = drawer.content().to_string();
let vecs = embedder
.embed_batch(std::slice::from_ref(&body))
.await
.with_context(|| format!("re-embed drawer {}", drawer.id))?;
if let Some(v) = vecs.into_iter().next() {
handle
.vector_store
.upsert(drawer.id, v)
.await
.with_context(|| format!("re-upsert drawer {}", drawer.id))?;
rebuilt += 1;
}
}
Ok(rebuilt)
}
pub(crate) fn build_closet_index(drawers: &[Drawer]) -> HashMap<String, Vec<Uuid>> {
let mut new_index: HashMap<String, Vec<Uuid>> = HashMap::new();
for drawer in drawers.iter() {
for kw in extract_keywords(drawer.content()) {
new_index.entry(kw).or_default().push(drawer.id);
}
}
new_index
}