use pulsedb::{CollectiveId, Config, ExperienceUpdate, NewExperience, PulseDB, WatchEventType};
use tempfile::tempdir;
const DIM: usize = 384;
fn dummy_embedding() -> Vec<f32> {
vec![0.1; DIM]
}
fn open_db() -> (PulseDB, tempfile::TempDir) {
let dir = tempdir().unwrap();
let path = dir.path().join("test.db");
let db = PulseDB::open(&path, Config::default()).unwrap();
(db, dir)
}
fn open_db_with_collective() -> (PulseDB, CollectiveId, tempfile::TempDir) {
let (db, dir) = open_db();
let cid = db.create_collective("test-collective").unwrap();
(db, cid, dir)
}
fn minimal_experience(collective_id: CollectiveId) -> NewExperience {
NewExperience {
collective_id,
content: "Cross-process watch test experience".to_string(),
embedding: Some(dummy_embedding()),
..Default::default()
}
}
#[test]
fn test_get_current_sequence_starts_at_zero() {
let (db, _dir) = open_db();
assert_eq!(db.get_current_sequence().unwrap(), 0);
}
#[test]
fn test_get_current_sequence_returns_latest() {
let (db, cid, _dir) = open_db_with_collective();
db.record_experience(minimal_experience(cid)).unwrap();
assert_eq!(db.get_current_sequence().unwrap(), 2);
db.record_experience(minimal_experience(cid)).unwrap();
assert_eq!(db.get_current_sequence().unwrap(), 3);
db.record_experience(minimal_experience(cid)).unwrap();
assert_eq!(db.get_current_sequence().unwrap(), 4);
}
#[test]
fn test_poll_changes_empty_database() {
let (db, _dir) = open_db();
let (events, seq) = db.poll_changes(0).unwrap();
assert!(events.is_empty());
assert_eq!(seq, 0);
}
#[test]
fn test_poll_changes_after_record() {
let (db, cid, _dir) = open_db_with_collective();
let exp_id = db.record_experience(minimal_experience(cid)).unwrap();
let (events, seq) = db.poll_changes(0).unwrap();
assert_eq!(events.len(), 1); assert_eq!(seq, 2);
let event = &events[0];
assert_eq!(event.experience_id, exp_id);
assert_eq!(event.collective_id, cid);
assert_eq!(event.event_type, WatchEventType::Created);
}
#[test]
fn test_poll_changes_incremental() {
let (db, cid, _dir) = open_db_with_collective();
db.record_experience(minimal_experience(cid)).unwrap();
db.record_experience(minimal_experience(cid)).unwrap();
db.record_experience(minimal_experience(cid)).unwrap();
let (events, seq) = db.poll_changes(0).unwrap();
assert_eq!(events.len(), 3);
assert_eq!(seq, 4);
db.record_experience(minimal_experience(cid)).unwrap();
db.record_experience(minimal_experience(cid)).unwrap();
let (events, seq) = db.poll_changes(4).unwrap();
assert_eq!(events.len(), 2);
assert_eq!(seq, 6);
let (events, seq) = db.poll_changes(6).unwrap();
assert_eq!(events.len(), 0);
assert_eq!(seq, 6);
}
#[test]
fn test_poll_changes_mixed_operations() {
let (db, cid, _dir) = open_db_with_collective();
let exp_id = db.record_experience(minimal_experience(cid)).unwrap();
db.update_experience(
exp_id,
ExperienceUpdate {
importance: Some(0.99),
..Default::default()
},
)
.unwrap();
db.reinforce_experience(exp_id).unwrap();
db.update_experience(
exp_id,
ExperienceUpdate {
archived: Some(true),
..Default::default()
},
)
.unwrap();
db.update_experience(
exp_id,
ExperienceUpdate {
archived: Some(false),
..Default::default()
},
)
.unwrap();
db.delete_experience(exp_id).unwrap();
let (events, seq) = db.poll_changes(0).unwrap();
assert_eq!(events.len(), 6);
assert_eq!(seq, 7);
assert_eq!(events[0].event_type, WatchEventType::Created);
assert_eq!(events[1].event_type, WatchEventType::Updated); assert_eq!(events[2].event_type, WatchEventType::Updated); assert_eq!(events[3].event_type, WatchEventType::Archived); assert_eq!(events[4].event_type, WatchEventType::Updated); assert_eq!(events[5].event_type, WatchEventType::Deleted);
}
#[test]
fn test_poll_changes_batch_limit() {
let (db, cid, _dir) = open_db_with_collective();
for _ in 0..10 {
db.record_experience(minimal_experience(cid)).unwrap();
}
let (events, seq) = db.poll_changes_batch(0, 3).unwrap();
assert_eq!(events.len(), 2);
assert_eq!(seq, 3);
let (events, seq) = db.poll_changes_batch(3, 3).unwrap();
assert_eq!(events.len(), 3);
assert_eq!(seq, 6);
let (events, seq) = db.poll_changes_batch(6, 3).unwrap();
assert_eq!(events.len(), 3);
assert_eq!(seq, 9);
let (events, seq) = db.poll_changes_batch(9, 100).unwrap();
assert_eq!(events.len(), 2);
assert_eq!(seq, 11);
}
#[test]
fn test_sequence_survives_close_reopen() {
let dir = tempdir().unwrap();
let path = dir.path().join("test.db");
let cid;
{
let db = PulseDB::open(&path, Config::default()).unwrap();
cid = db.create_collective("test").unwrap();
db.record_experience(minimal_experience(cid)).unwrap();
db.record_experience(minimal_experience(cid)).unwrap();
db.record_experience(minimal_experience(cid)).unwrap();
assert_eq!(db.get_current_sequence().unwrap(), 4);
db.close().unwrap();
}
{
let db = PulseDB::open(&path, Config::default()).unwrap();
assert_eq!(db.get_current_sequence().unwrap(), 4);
let (events, seq) = db.poll_changes(0).unwrap();
assert_eq!(events.len(), 3);
assert_eq!(seq, 4);
assert!(events
.iter()
.all(|e| e.event_type == WatchEventType::Created));
db.record_experience(minimal_experience(cid)).unwrap();
assert_eq!(db.get_current_sequence().unwrap(), 5);
db.close().unwrap();
}
}
#[test]
fn test_poll_changes_collective_isolation() {
let (db, _dir) = open_db();
let cid_a = db.create_collective("alpha").unwrap();
let cid_b = db.create_collective("beta").unwrap();
db.record_experience(minimal_experience(cid_a)).unwrap();
db.record_experience(minimal_experience(cid_b)).unwrap();
db.record_experience(minimal_experience(cid_a)).unwrap();
let (events, seq) = db.poll_changes(0).unwrap();
assert_eq!(events.len(), 3);
assert_eq!(seq, 5);
assert_eq!(events[0].collective_id, cid_a);
assert_eq!(events[1].collective_id, cid_b);
assert_eq!(events[2].collective_id, cid_a);
}