use super::dream::{DreamConfig, Dreamer};
use super::maintenance_log::{
DeletionReason, MAINTENANCE_LOG_ROTATED_FILENAME, MaintenanceDeletion, RecordOutcome, append,
journal_path, read_journal, record,
};
use super::palace::{Drawer, Palace, PalaceId, RoomType};
use super::retrieval::{ForgetOutcome, PalaceHandle, seed_shared_embedder_with_mock};
use super::store::kg::KnowledgeGraph;
use chrono::{Duration, Utc};
use std::io::Write;
use std::path::{Path, PathBuf};
use std::sync::{Arc, Mutex};
use tempfile::{TempDir, tempdir};
use tracing_subscriber::fmt::MakeWriter;
use uuid::Uuid;
#[derive(Clone, Default)]
struct Capture(Arc<Mutex<Vec<u8>>>);
impl Capture {
fn text(&self) -> String {
String::from_utf8_lossy(&self.0.lock().unwrap()).into_owned()
}
}
impl Write for Capture {
fn write(&mut self, buf: &[u8]) -> std::io::Result<usize> {
self.0.lock().unwrap().extend_from_slice(buf);
Ok(buf.len())
}
fn flush(&mut self) -> std::io::Result<()> {
Ok(())
}
}
impl<'a> MakeWriter<'a> for Capture {
type Writer = Capture;
fn make_writer(&'a self) -> Self::Writer {
self.clone()
}
}
fn capture_warn() -> (Capture, tracing::subscriber::DefaultGuard) {
let cap = Capture::default();
let subscriber = tracing_subscriber::fmt()
.with_writer(cap.clone())
.with_ansi(false)
.with_max_level(tracing::Level::WARN)
.finish();
(cap, tracing::subscriber::set_default(subscriber))
}
fn palace_at(root: &Path, name: &str) -> Palace {
let data_dir = root.join(name);
std::fs::create_dir_all(&data_dir).unwrap();
Palace {
id: PalaceId::new(name),
name: name.into(),
description: None,
created_at: Utc::now(),
data_dir,
}
}
fn open_palace(name: &str) -> (TempDir, PathBuf, Arc<PalaceHandle>) {
seed_shared_embedder_with_mock();
let dir = tempdir().unwrap();
let palace = palace_at(dir.path(), name);
let handle = PalaceHandle::open(&palace).unwrap();
(dir, palace.data_dir, handle)
}
fn dedup_only_config() -> DreamConfig {
DreamConfig {
semantic: super::semantic_consolidation::SemanticConsolidationConfig {
enabled: false,
..Default::default()
},
..DreamConfig::default()
}
}
#[tokio::test]
async fn dream_dedup_records_the_removed_and_surviving_drawer() {
let (_dir, data_dir, handle) = open_palace("dedup-trail");
let content = "Rust uses HNSW for vector search";
let lose = handle
.remember(content.into(), RoomType::Backend, vec![], 0.7)
.await
.unwrap();
let keep = handle
.remember(content.into(), RoomType::Backend, vec![], 0.6)
.await
.unwrap();
let (log, _guard) = capture_warn();
let stats = Dreamer::new(dedup_only_config())
.dream_cycle(&handle)
.await
.unwrap();
assert_eq!(stats.merged, 1);
let journal = read_journal(&data_dir).unwrap();
assert_eq!(journal.records.len(), 1, "{:?}", journal.records);
let rec = &journal.records[0];
assert_eq!(rec.palace, "dedup-trail");
assert_eq!(rec.drawer_id, lose);
assert_eq!(rec.reason, DeletionReason::DreamDedup);
assert_eq!(rec.survivor_id, Some(keep));
let score = rec.score.expect("dedup records its score");
assert!(score >= dedup_only_config().dedup_threshold, "{score}");
assert_eq!(rec.pid, std::process::id());
let text = log.text();
assert!(
text.contains("WARN") && text.contains("dream cycle removed 1 drawer"),
"{text}"
);
}
#[tokio::test]
async fn every_maintenance_removal_logs_its_id_and_reason() {
let (_dir, _data_dir, handle) = open_palace("dedup-log");
let content = "Dream dedup merges near-duplicate drawers above the threshold";
let lose = handle
.remember(content.into(), RoomType::Backend, vec![], 0.7)
.await
.unwrap();
handle
.remember(content.into(), RoomType::Backend, vec![], 0.6)
.await
.unwrap();
let (log, _guard) = capture_warn();
let stats = Dreamer::new(dedup_only_config())
.dream_cycle(&handle)
.await
.unwrap();
assert_eq!(stats.merged, 1);
let text = log.text();
let line = text
.lines()
.find(|l| l.contains(&lose.to_string()))
.unwrap_or_else(|| panic!("no log line names removed drawer {lose}: {text}"));
for needle in ["WARN", "dedup-log", "dream_dedup"] {
assert!(line.contains(needle), "missing {needle} in {line}");
}
}
#[test]
fn the_open_time_purge_records_each_expired_drawer() {
let dir = tempdir().unwrap();
let palace = palace_at(dir.path(), "open-purge-trail");
let mut expired = Drawer::new(Uuid::new_v4(), "an expired session event");
expired.expires_at = Some(Utc::now() - Duration::days(1));
let expired_id = expired.id;
{
let kg = KnowledgeGraph::open(&palace.data_dir.join("kg.db")).unwrap();
kg.upsert_drawer_sync(&expired).unwrap();
}
let handle = PalaceHandle::open(&palace).unwrap();
assert!(handle.drawers.read().is_empty());
let journal = read_journal(&palace.data_dir).unwrap();
assert_eq!(journal.records.len(), 1, "{:?}", journal.records);
let rec = &journal.records[0];
assert_eq!(rec.palace, "open-purge-trail");
assert_eq!(rec.drawer_id, expired_id);
assert_eq!(rec.reason, DeletionReason::ExpiredPurgeAtOpen);
assert_eq!(rec.survivor_id, None);
}
#[tokio::test]
async fn purge_expired_records_each_drawer() {
let (_dir, data_dir, handle) = open_palace("ttl-purge-trail");
let mut expired = Drawer::new(Uuid::new_v4(), "expired");
expired.expires_at = Some(Utc::now() - Duration::days(1));
let expired_id = expired.id;
handle.add_drawer(expired);
assert_eq!(handle.purge_expired().await.unwrap(), 1);
let records = read_journal(&data_dir).unwrap().records;
assert_eq!(records.len(), 1, "{records:?}");
assert_eq!(records[0].drawer_id, expired_id);
assert_eq!(records[0].reason, DeletionReason::ExpiredPurge);
}
#[tokio::test]
async fn a_user_forget_writes_a_user_forget_journal_record() {
let (_dir, data_dir, handle) = open_palace("user-forget");
let text = "a fact the user chose to forget";
let id = handle
.remember(text.into(), RoomType::General, vec![], 0.5)
.await
.unwrap();
assert_eq!(handle.forget(id).await.unwrap(), ForgetOutcome::Deleted);
let records = read_journal(&data_dir).unwrap().records;
assert_eq!(records.len(), 1, "{records:?}");
let rec = &records[0];
assert_eq!(rec.palace, "user-forget");
assert_eq!(rec.drawer_id, id);
assert_eq!(rec.reason, DeletionReason::UserForget);
assert!(rec.drawer.is_none(), "a user forget must keep no copy");
assert_eq!(
rec.content_hash,
Some(super::content_hash::memory_content_hash(text))
);
}
#[tokio::test]
async fn forgotten_content_cannot_be_recovered_from_the_journal() {
let (_dir, data_dir, handle) = open_palace("forget-secret");
let secret = "marigold-quokka-anniversary-9283";
let id = handle
.remember(
format!("the surprise party codeword is {secret}"),
RoomType::General,
vec![],
0.5,
)
.await
.unwrap();
assert_eq!(handle.forget(id).await.unwrap(), ForgetOutcome::Deleted);
let journal = std::fs::read_to_string(journal_path(&data_dir)).unwrap();
assert!(
journal.contains(&id.to_string()),
"the record names the drawer"
);
assert!(
!journal.contains(secret) && !journal.contains("surprise party"),
"forgotten content leaked into the journal: {journal}"
);
}
#[tokio::test]
async fn a_failed_record_write_logs_the_record_and_still_deletes() {
let (_dir, data_dir, handle) = open_palace("broken-journal");
std::fs::create_dir_all(journal_path(&data_dir)).unwrap();
let survivor = Uuid::new_v4();
let mut doomed = Drawer::new(Uuid::new_v4(), "a drawer dedup will remove");
doomed.importance = 0.4;
let doomed_id = doomed.id;
handle.add_drawer(doomed);
let (log, _guard) = capture_warn();
let outcome = handle
.forget_for_maintenance(
doomed_id,
DeletionReason::DreamDedup,
Some((survivor, Some(0.97))),
)
.await
.unwrap();
assert_eq!(outcome, ForgetOutcome::Deleted, "the deletion proceeds");
assert!(handle.drawers.read().iter().all(|d| d.id != doomed_id));
let text = log.text();
for needle in [
"ERROR",
"broken-journal",
&doomed_id.to_string(),
&survivor.to_string(),
"dream_dedup",
"0.97",
"a drawer dedup will remove",
] {
assert!(text.contains(needle), "missing {needle} in {text}");
}
let rec = MaintenanceDeletion::new(&handle.id, doomed_id, DeletionReason::DreamPrune);
assert_eq!(record(Some(&data_dir), &rec), RecordOutcome::LoggedOnly);
}
#[tokio::test]
async fn a_failed_snapshot_save_after_the_delete_still_journals_the_copy() {
let (_dir, data_dir, handle) = open_palace("broken-snapshot");
let content = "a drawer whose snapshot save will fail";
let id = handle
.remember(content.into(), RoomType::General, vec![], 0.5)
.await
.unwrap();
let snapshot = data_dir.join("l1_cache.json");
let _ = std::fs::remove_file(&snapshot);
std::fs::create_dir_all(snapshot.join("occupied")).unwrap();
let err = handle
.forget_for_maintenance(id, DeletionReason::DreamPrune, None)
.await
.expect_err("the failed snapshot save is reported");
assert!(format!("{err:#}").contains("L1 snapshot"), "{err:#}");
assert!(handle.drawers.read().iter().all(|d| d.id != id));
let journal = read_journal(&data_dir).unwrap();
let rec = journal
.records
.iter()
.find(|r| r.drawer_id == id)
.unwrap_or_else(|| panic!("no record for the removed drawer: {:?}", journal.records));
assert_eq!(
rec.drawer.as_ref().map(|d| d.content.as_str()),
Some(content)
);
}
#[test]
fn the_journal_rotates_and_reads_back_oldest_first() {
let dir = tempdir().unwrap();
let palace = PalaceId::new("rotate");
let first = MaintenanceDeletion::new(&palace, Uuid::new_v4(), DeletionReason::DreamPrune);
let second = MaintenanceDeletion::new(&palace, Uuid::new_v4(), DeletionReason::ExpiredPurge)
.with_survivor(Uuid::new_v4(), None);
append(dir.path(), &first, u64::MAX).unwrap();
append(dir.path(), &second, 1).unwrap();
assert!(dir.path().join(MAINTENANCE_LOG_ROTATED_FILENAME).exists());
let mut live = std::fs::OpenOptions::new()
.append(true)
.open(journal_path(dir.path()))
.unwrap();
live.write_all(b"{\"torn\":\n").unwrap();
let journal = read_journal(dir.path()).unwrap();
assert_eq!(journal.records, vec![first, second]);
assert_eq!(journal.malformed, 1);
}
#[test]
fn concurrent_appends_across_the_rotation_boundary_lose_no_record() {
const WRITERS: usize = 16;
let palace = PalaceId::new("rotate-race");
for round in 0..20 {
let dir = tempdir().unwrap();
let prefill = 40;
for _ in 0..prefill {
let rec = MaintenanceDeletion::new(&palace, Uuid::new_v4(), DeletionReason::DreamPrune);
append(dir.path(), &rec, u64::MAX).unwrap();
}
let rotate_at = std::fs::metadata(journal_path(dir.path())).unwrap().len();
let barrier = Arc::new(std::sync::Barrier::new(WRITERS));
let handles: Vec<_> = (0..WRITERS)
.map(|_| {
let (barrier, dir, palace) = (
Arc::clone(&barrier),
dir.path().to_path_buf(),
palace.clone(),
);
std::thread::spawn(move || {
let rec = MaintenanceDeletion::new(
&palace,
Uuid::new_v4(),
DeletionReason::DreamDedup,
);
barrier.wait();
append(&dir, &rec, rotate_at).map_err(|e| format!("{e:#}"))
})
})
.collect();
let results: Vec<_> = handles.into_iter().map(|h| h.join().unwrap()).collect();
let journal = read_journal(dir.path()).unwrap();
let errors: Vec<_> = results.iter().filter_map(|r| r.as_ref().err()).collect();
assert_eq!(
journal.records.len(),
prefill + WRITERS,
"round {round}: a record was lost; append errors: {errors:?}"
);
assert!(errors.is_empty(), "round {round}: {errors:?}");
}
}
#[test]
fn a_lock_or_rotation_failure_still_appends_the_record() {
let dir = tempdir().unwrap();
let palace = PalaceId::new("lock-fail");
let sidecar = crate::file_lock::lock_path(&journal_path(dir.path()));
std::fs::create_dir_all(&sidecar).unwrap();
let first = MaintenanceDeletion::new(&palace, Uuid::new_v4(), DeletionReason::DreamPrune);
append(dir.path(), &first, u64::MAX).unwrap();
std::fs::remove_dir(&sidecar).unwrap();
let rotated = dir.path().join(MAINTENANCE_LOG_ROTATED_FILENAME);
std::fs::create_dir_all(rotated.join("occupied")).unwrap();
let second = MaintenanceDeletion::new(&palace, Uuid::new_v4(), DeletionReason::ExpiredPurge);
append(dir.path(), &second, 1).unwrap();
let live = std::fs::read_to_string(journal_path(dir.path())).unwrap();
assert!(live.contains(&first.drawer_id.to_string()), "{live}");
assert!(live.contains(&second.drawer_id.to_string()), "{live}");
}
#[tokio::test]
async fn every_dream_removal_journals_a_recoverable_copy() {
use super::retrieval::RememberOptions;
let (_dir, data_dir, handle) = open_palace("recoverable");
let dup = "Rust uses HNSW for vector search";
let short = "hello world";
let stale = "very stale fact nobody cares about";
for content in [dup, dup] {
handle
.remember(content.into(), RoomType::Backend, vec!["rust".into()], 0.6)
.await
.unwrap();
}
handle
.remember_with_options(
short.into(),
RoomType::General,
vec![],
0.5,
RememberOptions::forced(),
)
.await
.unwrap();
let stale_id = handle
.remember(stale.into(), RoomType::General, vec![], 0.01)
.await
.unwrap();
for d in handle.drawers.write().iter_mut() {
if d.id == stale_id {
d.created_at = Utc::now() - Duration::days(60);
}
}
let originals: std::collections::HashMap<String, (String, Vec<String>)> = handle
.drawers
.read()
.iter()
.map(|d| (d.id.to_string(), (d.content().to_string(), d.tags.clone())))
.collect();
let stats = Dreamer::new(dedup_only_config())
.dream_cycle(&handle)
.await
.unwrap();
assert_eq!(
(stats.merged, stats.content_pruned, stats.pruned),
(1, 1, 1),
"{stats:?}"
);
let raw = std::fs::read_to_string(journal_path(&data_dir)).unwrap();
let lines: Vec<serde_json::Value> = raw
.lines()
.map(|l| serde_json::from_str(l).unwrap())
.collect();
assert_eq!(lines.len(), 3, "{raw}");
for line in &lines {
let id = line["drawer_id"].as_str().unwrap();
let (content, tags) = &originals[id];
let copy = &line["drawer"];
assert_eq!(copy["content"], content.as_str(), "{line}");
assert_eq!(copy["tags"], serde_json::json!(tags), "{line}");
assert!(copy["importance"].is_number(), "{line}");
}
}