use rusqlite::params;
use crate::hlc::HLCTimestamp;
use crate::types::*;
use super::YantrikDB;
fn vec_seed(seed: f32, dim: usize) -> Vec<f32> {
let raw: Vec<f32> = (0..dim).map(|i| (seed + i as f32) * 0.1).collect();
let norm: f32 = raw.iter().map(|x| x * x).sum::<f32>().sqrt();
raw.iter().map(|x| x / norm).collect()
}
fn empty_meta() -> serde_json::Value {
serde_json::json!({})
}
#[test]
fn test_new_and_stats() {
let db = YantrikDB::new(":memory:", 8).unwrap();
let s = db.stats(None).unwrap();
assert_eq!(s.active_memories, 0);
assert_eq!(s.edges, 0);
}
#[test]
fn test_actor_id_auto_generated() {
let db = YantrikDB::new(":memory:", 8).unwrap();
assert_eq!(db.actor_id().len(), 36); }
#[test]
fn test_actor_id_explicit() {
let db = YantrikDB::new_with_actor(":memory:", 8, "device-A").unwrap();
assert_eq!(db.actor_id(), "device-A");
}
#[test]
fn test_record_auto_extracts_entities() {
let db = YantrikDB::new(":memory:", 8).unwrap();
let rid = db
.record(
"Alice Chen is the CEO of Acme Corp",
"semantic",
0.8,
0.0,
604800.0,
&empty_meta(),
&vec_seed(1.0, 8),
"default",
0.8,
"people",
"user",
None,
)
.unwrap();
db.apply_pending_ops_once(100).unwrap();
let entities: Vec<String> = {
let conn = db.conn();
let mut stmt = conn
.prepare("SELECT entity_name FROM memory_entities WHERE memory_rid = ?1")
.unwrap();
stmt.query_map(params![rid], |r| r.get::<_, String>(0))
.unwrap()
.collect::<std::result::Result<Vec<_>, _>>()
.unwrap()
};
assert!(
entities.contains(&"Alice Chen".to_string()),
"got: {:?}",
entities
);
assert!(
entities.contains(&"Acme Corp".to_string()),
"got: {:?}",
entities
);
let entity_count: i64 = db
.conn()
.query_row("SELECT COUNT(*) FROM entities", [], |r| r.get(0))
.unwrap();
assert!(
entity_count >= 2,
"expected >= 2 entities, got {}",
entity_count
);
}
#[test]
fn test_record_batch_auto_extracts_entities() {
let db = YantrikDB::new(":memory:", 8).unwrap();
let inputs = vec![
RecordInput {
text: "Alice Chen is the CEO of Acme Corp".to_string(),
memory_type: "semantic".to_string(),
importance: 0.8,
valence: 0.0,
half_life: 604800.0,
metadata: empty_meta(),
embedding: vec_seed(1.0, 8),
namespace: "default".to_string(),
certainty: 0.8,
domain: "people".to_string(),
source: "user".to_string(),
emotional_state: None,
},
RecordInput {
text: "Sarah Kim is the CTO of Acme Corp".to_string(),
memory_type: "semantic".to_string(),
importance: 0.8,
valence: 0.0,
half_life: 604800.0,
metadata: empty_meta(),
embedding: vec_seed(1.05, 8),
namespace: "default".to_string(),
certainty: 0.8,
domain: "people".to_string(),
source: "user".to_string(),
emotional_state: None,
},
];
let rids = db.record_batch(&inputs).unwrap();
assert_eq!(rids.len(), 2);
let total_links: i64 = db
.conn()
.query_row("SELECT COUNT(*) FROM memory_entities", [], |r| r.get(0))
.unwrap();
assert!(
total_links >= 3,
"expected batch to link both memories to entities, got {} links",
total_links
);
let load_entities = |rid: &str| -> Vec<String> {
let conn = db.conn();
let mut stmt = conn
.prepare("SELECT entity_name FROM memory_entities WHERE memory_rid = ?1")
.unwrap();
stmt.query_map(params![rid], |r| r.get::<_, String>(0))
.unwrap()
.collect::<std::result::Result<Vec<_>, _>>()
.unwrap()
};
let m1_entities = load_entities(&rids[0]);
let m2_entities = load_entities(&rids[1]);
assert!(m1_entities.contains(&"Alice Chen".to_string()));
assert!(m2_entities.contains(&"Sarah Kim".to_string()));
assert!(!m1_entities.contains(&"Sarah Kim".to_string()));
assert!(!m2_entities.contains(&"Alice Chen".to_string()));
}
#[test]
fn test_record_and_get() {
let db = YantrikDB::new(":memory:", 8).unwrap();
let emb = vec_seed(1.0, 8);
let rid = db
.record(
"hello world",
"episodic",
0.8,
0.0,
604800.0,
&empty_meta(),
&emb,
"default",
0.8,
"general",
"user",
None,
)
.unwrap();
assert_eq!(rid.len(), 36);
let mem = db.get(&rid).unwrap().unwrap();
assert_eq!(mem.text, "hello world");
assert_eq!(mem.memory_type, "episodic");
assert_eq!(mem.importance, 0.8);
assert_eq!(mem.consolidation_status, "active");
}
#[test]
fn test_record_updates_stats() {
let db = YantrikDB::new(":memory:", 8).unwrap();
db.record(
"one",
"episodic",
0.5,
0.0,
604800.0,
&empty_meta(),
&vec_seed(1.0, 8),
"default",
0.8,
"general",
"user",
None,
)
.unwrap();
db.record(
"two",
"episodic",
0.5,
0.0,
604800.0,
&empty_meta(),
&vec_seed(2.0, 8),
"default",
0.8,
"general",
"user",
None,
)
.unwrap();
assert_eq!(db.stats(None).unwrap().active_memories, 2);
}
#[test]
fn test_recall_basic() {
let db = YantrikDB::new(":memory:", 8).unwrap();
db.record(
"the cat sat on the mat",
"episodic",
0.5,
0.0,
604800.0,
&empty_meta(),
&vec_seed(1.0, 8),
"default",
0.8,
"general",
"user",
None,
)
.unwrap();
db.record(
"dogs are loyal friends",
"episodic",
0.5,
0.0,
604800.0,
&empty_meta(),
&vec_seed(5.0, 8),
"default",
0.8,
"general",
"user",
None,
)
.unwrap();
db.record(
"cats love warm places",
"episodic",
0.5,
0.0,
604800.0,
&empty_meta(),
&vec_seed(1.1, 8),
"default",
0.8,
"general",
"user",
None,
)
.unwrap();
let results = db
.recall(
&vec_seed(1.0, 8),
2,
None,
None,
false,
false,
None,
false,
None,
None,
None,
None,
None,
)
.unwrap();
assert_eq!(results.len(), 2);
}
#[test]
fn test_recall_empty() {
let db = YantrikDB::new(":memory:", 8).unwrap();
let results = db
.recall(
&vec_seed(1.0, 8),
5,
None,
None,
false,
false,
None,
false,
None,
None,
None,
None,
None,
)
.unwrap();
assert!(results.is_empty());
}
#[test]
fn test_relate_and_get_edges() {
let db = YantrikDB::new(":memory:", 8).unwrap();
let eid = db.relate("Alice", "Bob", "knows", 1.0).unwrap();
assert_eq!(eid.len(), 36);
let edges = db.get_edges("Alice").unwrap();
assert_eq!(edges.len(), 1);
assert_eq!(edges[0].src, "Alice");
assert_eq!(edges[0].dst, "Bob");
}
#[test]
fn test_forget() {
let db = YantrikDB::new(":memory:", 8).unwrap();
let rid = db
.record(
"forget me",
"episodic",
0.5,
0.0,
604800.0,
&empty_meta(),
&vec_seed(1.0, 8),
"default",
0.8,
"general",
"user",
None,
)
.unwrap();
assert!(db.forget(&rid).unwrap());
let mem = db.get(&rid).unwrap().unwrap();
assert_eq!(mem.consolidation_status, "tombstoned");
}
#[test]
fn test_forget_nonexistent() {
let db = YantrikDB::new(":memory:", 8).unwrap();
assert!(!db.forget("nonexistent").unwrap());
}
#[test]
fn test_decay_fresh() {
let db = YantrikDB::new(":memory:", 8).unwrap();
db.record(
"fresh",
"episodic",
0.9,
0.0,
604800.0,
&empty_meta(),
&vec_seed(1.0, 8),
"default",
0.8,
"general",
"user",
None,
)
.unwrap();
let decayed = db.decay(0.01).unwrap();
assert!(decayed.is_empty());
}
#[test]
fn test_oplog_has_hlc() {
let db = YantrikDB::new(":memory:", 8).unwrap();
db.record(
"test",
"episodic",
0.5,
0.0,
604800.0,
&empty_meta(),
&vec_seed(1.0, 8),
"default",
0.8,
"general",
"user",
None,
)
.unwrap();
let hlc_bytes: Vec<u8> = db
.conn()
.query_row(
"SELECT hlc FROM oplog ORDER BY rowid DESC LIMIT 1",
[],
|row| row.get(0),
)
.unwrap();
assert_eq!(hlc_bytes.len(), 16);
let ts = HLCTimestamp::from_bytes(&hlc_bytes).unwrap();
assert!(ts.millis > 0);
}
#[test]
fn test_oplog_has_embedding_hash() {
let db = YantrikDB::new(":memory:", 8).unwrap();
db.record(
"test",
"episodic",
0.5,
0.0,
604800.0,
&empty_meta(),
&vec_seed(1.0, 8),
"default",
0.8,
"general",
"user",
None,
)
.unwrap();
let hash: Vec<u8> = db
.conn()
.query_row(
"SELECT embedding_hash FROM oplog WHERE op_type = 'record' LIMIT 1",
[],
|row| row.get(0),
)
.unwrap();
assert_eq!(hash.len(), 32); }
#[test]
fn test_oplog_enriched_payload() {
let db = YantrikDB::new(":memory:", 8).unwrap();
db.record(
"test payload",
"semantic",
0.7,
0.3,
1000.0,
&serde_json::json!({"key": "val"}),
&vec_seed(1.0, 8),
"default",
0.8,
"general",
"user",
None,
)
.unwrap();
let payload_str: String = db
.conn()
.query_row(
"SELECT payload FROM oplog WHERE op_type = 'record' LIMIT 1",
[],
|row| row.get(0),
)
.unwrap();
let payload: serde_json::Value = serde_json::from_str(&payload_str).unwrap();
assert_eq!(payload["type"], "semantic");
assert_eq!(payload["text"], "test payload");
assert_eq!(payload["importance"], 0.7);
assert_eq!(payload["valence"], 0.3);
assert_eq!(payload["half_life"], 1000.0);
assert!(payload["rid"].is_string());
assert!(payload["created_at"].is_number());
assert!(payload["metadata"]["key"] == "val");
}
#[test]
fn test_schema_v3_has_conflicts_table() {
let db = YantrikDB::new(":memory:", 8).unwrap();
let count: i64 = db
.conn()
.query_row(
"SELECT COUNT(*) FROM sqlite_master WHERE type='table' AND name='conflicts'",
[],
|row| row.get(0),
)
.unwrap();
assert_eq!(count, 1);
}
#[test]
fn test_resolve_keep_a() {
let db = YantrikDB::new(":memory:", 8).unwrap();
let rid_a = db
.record(
"birthday March 5",
"episodic",
0.5,
0.0,
604800.0,
&empty_meta(),
&vec_seed(1.0, 8),
"default",
0.8,
"general",
"user",
None,
)
.unwrap();
let rid_b = db
.record(
"birthday March 15",
"episodic",
0.5,
0.0,
604800.0,
&empty_meta(),
&vec_seed(2.0, 8),
"default",
0.8,
"general",
"user",
None,
)
.unwrap();
let conflict = crate::conflict::create_conflict(
&db,
&crate::types::ConflictType::IdentityFact,
&rid_a,
&rid_b,
Some("User"),
Some("birthday"),
"conflicting birthdays",
)
.unwrap();
let result = db
.resolve_conflict(
&conflict.conflict_id,
"keep_a",
Some(&rid_a),
None,
Some("User confirmed March 5"),
)
.unwrap();
assert!(result.loser_tombstoned);
let mem_b = db.get(&rid_b).unwrap().unwrap();
assert_eq!(mem_b.consolidation_status, "tombstoned");
let resolved = db.get_conflict(&conflict.conflict_id).unwrap().unwrap();
assert_eq!(resolved.status, "resolved");
assert_eq!(resolved.strategy.as_deref(), Some("keep_a"));
}
#[test]
fn test_resolve_keep_both() {
let db = YantrikDB::new(":memory:", 8).unwrap();
let rid_a = db
.record(
"a",
"episodic",
0.5,
0.0,
604800.0,
&empty_meta(),
&vec_seed(1.0, 8),
"default",
0.8,
"general",
"user",
None,
)
.unwrap();
let rid_b = db
.record(
"b",
"episodic",
0.5,
0.0,
604800.0,
&empty_meta(),
&vec_seed(2.0, 8),
"default",
0.8,
"general",
"user",
None,
)
.unwrap();
let conflict = crate::conflict::create_conflict(
&db,
&crate::types::ConflictType::Minor,
&rid_a,
&rid_b,
None,
None,
"test",
)
.unwrap();
let result = db
.resolve_conflict(&conflict.conflict_id, "keep_both", None, None, None)
.unwrap();
assert!(!result.loser_tombstoned);
let mem_a = db.get(&rid_a).unwrap().unwrap();
let mem_b = db.get(&rid_b).unwrap().unwrap();
assert_eq!(mem_a.consolidation_status, "active");
assert_eq!(mem_b.consolidation_status, "active");
}
#[test]
fn test_correct_memory() {
let db = YantrikDB::new(":memory:", 8).unwrap();
let rid = db
.record(
"favorite color is green",
"episodic",
0.7,
0.0,
604800.0,
&empty_meta(),
&vec_seed(1.0, 8),
"default",
0.8,
"general",
"user",
None,
)
.unwrap();
let result = db
.correct(
&rid,
Some("favorite color is blue"),
None, Some(0.9), None, "User corrected their favorite color", )
.unwrap();
assert_eq!(result.corrected_rid, rid);
assert_eq!(result.original_rid, rid);
assert!(!result.original_tombstoned);
assert_eq!(result.revision_num, 1);
let updated = db.get(&rid).unwrap().unwrap();
assert_ne!(updated.consolidation_status, "tombstoned");
assert_eq!(updated.text, "favorite color is blue");
assert!((updated.importance - 0.9).abs() < 1e-9);
}
#[test]
fn test_get_conflicts_filtered() {
let db = YantrikDB::new(":memory:", 8).unwrap();
let rid_a = db
.record(
"a",
"episodic",
0.5,
0.0,
604800.0,
&empty_meta(),
&vec_seed(1.0, 8),
"default",
0.8,
"general",
"user",
None,
)
.unwrap();
let rid_b = db
.record(
"b",
"episodic",
0.5,
0.0,
604800.0,
&empty_meta(),
&vec_seed(2.0, 8),
"default",
0.8,
"general",
"user",
None,
)
.unwrap();
let rid_c = db
.record(
"c",
"episodic",
0.5,
0.0,
604800.0,
&empty_meta(),
&vec_seed(3.0, 8),
"default",
0.8,
"general",
"user",
None,
)
.unwrap();
crate::conflict::create_conflict(
&db,
&crate::types::ConflictType::IdentityFact,
&rid_a,
&rid_b,
Some("User"),
Some("birthday"),
"test 1",
)
.unwrap();
crate::conflict::create_conflict(
&db,
&crate::types::ConflictType::Preference,
&rid_b,
&rid_c,
Some("User"),
Some("prefers"),
"test 2",
)
.unwrap();
let all = db.get_conflicts(None, None, None, None, None, 50).unwrap();
assert_eq!(all.len(), 2);
let identity_only = db
.get_conflicts(None, Some("identity_fact"), None, None, None, 50)
.unwrap();
assert_eq!(identity_only.len(), 1);
let critical = db
.get_conflicts(None, None, None, Some("critical"), None, 50)
.unwrap();
assert_eq!(critical.len(), 1);
}
#[test]
fn test_dismiss_conflict() {
let db = YantrikDB::new(":memory:", 8).unwrap();
let rid_a = db
.record(
"a",
"episodic",
0.5,
0.0,
604800.0,
&empty_meta(),
&vec_seed(1.0, 8),
"default",
0.8,
"general",
"user",
None,
)
.unwrap();
let rid_b = db
.record(
"b",
"episodic",
0.5,
0.0,
604800.0,
&empty_meta(),
&vec_seed(2.0, 8),
"default",
0.8,
"general",
"user",
None,
)
.unwrap();
let conflict = crate::conflict::create_conflict(
&db,
&crate::types::ConflictType::Minor,
&rid_a,
&rid_b,
None,
None,
"test",
)
.unwrap();
db.dismiss_conflict(&conflict.conflict_id, Some("Not really a conflict"))
.unwrap();
let c = db.get_conflict(&conflict.conflict_id).unwrap().unwrap();
assert_eq!(c.status, "dismissed");
}
#[test]
fn test_stats_include_conflicts() {
let db = YantrikDB::new(":memory:", 8).unwrap();
let s = db.stats(None).unwrap();
assert_eq!(s.open_conflicts, 0);
assert_eq!(s.resolved_conflicts, 0);
let rid_a = db
.record(
"a",
"episodic",
0.5,
0.0,
604800.0,
&empty_meta(),
&vec_seed(1.0, 8),
"default",
0.8,
"general",
"user",
None,
)
.unwrap();
let rid_b = db
.record(
"b",
"episodic",
0.5,
0.0,
604800.0,
&empty_meta(),
&vec_seed(2.0, 8),
"default",
0.8,
"general",
"user",
None,
)
.unwrap();
crate::conflict::create_conflict(
&db,
&crate::types::ConflictType::Minor,
&rid_a,
&rid_b,
None,
None,
"test",
)
.unwrap();
let s = db.stats(None).unwrap();
assert_eq!(s.open_conflicts, 1);
assert_eq!(s.resolved_conflicts, 0);
}
#[test]
fn test_schema_v4_has_trigger_log_and_patterns() {
let db = YantrikDB::new(":memory:", 8).unwrap();
let count: i64 = db.conn().query_row(
"SELECT COUNT(*) FROM sqlite_master WHERE type='table' AND name IN ('trigger_log', 'patterns')",
[], |row| row.get(0),
).unwrap();
assert_eq!(count, 2);
}
#[test]
fn test_schema_v19_has_rfc007_tables() {
let db = YantrikDB::new(":memory:", 8).unwrap();
let count: i64 = db
.conn()
.query_row(
"SELECT COUNT(*) FROM sqlite_master WHERE type='table' AND name IN \
('propositions', 'variables', 'state_assertions', 'rule_edges', 'scenario_specs')",
[],
|row| row.get(0),
)
.unwrap();
assert_eq!(
count, 5,
"RFC 007 Phase 0 should create all five new tables"
);
}
#[test]
fn test_schema_v19_claims_has_proposition_id() {
let db = YantrikDB::new(":memory:", 8).unwrap();
let conn = db.conn();
let mut stmt = conn.prepare("PRAGMA table_info(claims)").unwrap();
let cols: Vec<String> = stmt
.query_map([], |row| row.get::<_, String>(1))
.unwrap()
.filter_map(|r| r.ok())
.collect();
assert!(
cols.contains(&"proposition_id".to_string()),
"claims table should have a proposition_id column after V19. Got: {:?}",
cols
);
}
#[test]
fn test_schema_v19_rule_edge_whitelist_enforced() {
let db = YantrikDB::new(":memory:", 8).unwrap();
db.conn()
.execute(
"INSERT INTO variables (variable_id, name, namespace, value_space, scope, created_at) \
VALUES ('v1', 'var_a', 'default', '{}', 'generic', 0.0)",
[],
)
.unwrap();
db.conn()
.execute(
"INSERT INTO variables (variable_id, name, namespace, value_space, scope, created_at) \
VALUES ('v2', 'var_b', 'default', '{}', 'generic', 0.0)",
[],
)
.unwrap();
let result = db.conn().execute(
"INSERT INTO rule_edges (rule_id, parent_variable_id, child_variable_id, edge_type, \
direction_confidence, persistence, scope, source, namespace, created_at) \
VALUES ('r1', 'v1', 'v2', 'implies', 'high', 'instantaneous', 'generic', 'test', 'default', 0.0)",
[],
);
assert!(
result.is_err(),
"rule_edges should reject edge_type='implies' — only whitelist (causal_promotes, causal_inhibits, requires) allowed"
);
db.conn().execute(
"INSERT INTO rule_edges (rule_id, parent_variable_id, child_variable_id, edge_type, \
direction_confidence, persistence, scope, source, namespace, created_at) \
VALUES ('r2', 'v1', 'v2', 'causal_promotes', 'high', 'instantaneous', 'generic', 'test', 'default', 0.0)",
[],
).unwrap();
}
#[test]
fn test_schema_v20_has_rfc008_phase1_tables() {
let db = YantrikDB::new(":memory:", 8).unwrap();
let count: i64 = db
.conn()
.query_row(
"SELECT COUNT(*) FROM sqlite_master WHERE type='table' AND name IN \
('mobility_state', 'actor_profile', 'compression_artifact')",
[],
|row| row.get(0),
)
.unwrap();
assert_eq!(
count, 3,
"RFC 008 Phase 1 should create mobility_state, actor_profile, compression_artifact"
);
}
#[test]
fn test_schema_v20_claims_has_mobility_signals() {
let db = YantrikDB::new(":memory:", 8).unwrap();
let conn = db.conn();
let mut stmt = conn.prepare("PRAGMA table_info(claims)").unwrap();
let cols: Vec<String> = stmt
.query_map([], |row| row.get::<_, String>(1))
.unwrap()
.filter_map(|r| r.ok())
.collect();
for expected in &[
"regime_tag",
"self_generated",
"source_lineage",
"modality_signal",
] {
assert!(
cols.contains(&expected.to_string()),
"claims table should have column {} after V20. Got: {:?}",
expected,
cols
);
}
}
#[test]
fn test_schema_v20_actor_profile_whitelist_enforced() {
let db = YantrikDB::new(":memory:", 8).unwrap();
db.conn()
.execute(
"INSERT INTO actor_profile (actor_id, actor_type, regime, last_updated) \
VALUES ('ext_medical', 'extractor', 'medical', 0.0)",
[],
)
.unwrap();
let bad = db.conn().execute(
"INSERT INTO actor_profile (actor_id, actor_type, regime, last_updated) \
VALUES ('weird', 'hallucinator', 'default', 0.0)",
[],
);
assert!(
bad.is_err(),
"actor_profile should reject actor_type not in the whitelist"
);
}
#[test]
fn test_schema_v20_compression_artifact_status_whitelist() {
let db = YantrikDB::new(":memory:", 8).unwrap();
db.conn().execute(
"INSERT INTO compression_artifact (artifact_id, source_span_json, abstraction_operator, \
reversibility_pointer, namespace, created_at, status) \
VALUES ('a1', '[]', 'consolidate_v1', 'raw:1-100', 'default', 0.0, 'active')",
[],
).unwrap();
let bad = db.conn().execute(
"INSERT INTO compression_artifact (artifact_id, source_span_json, abstraction_operator, \
reversibility_pointer, namespace, created_at, status) \
VALUES ('a2', '[]', 'x', 'y', 'default', 0.0, 'freshly_minted')",
[],
);
assert!(
bad.is_err(),
"compression_artifact.status whitelist should reject unknown values"
);
}
#[test]
fn test_schema_v20_mobility_state_roundtrip() {
let db = YantrikDB::new(":memory:", 8).unwrap();
db.conn()
.execute(
"INSERT INTO propositions (proposition_id, src, rel_type, dst, namespace, created_at) \
VALUES ('p1', 'Alice', 'works_at', 'Acme', 'default', 0.0)",
[],
)
.unwrap();
db.conn()
.execute(
"INSERT INTO mobility_state (proposition_id, regime, snapshot_ts, \
support_mass, attack_mass, self_gen_local, modality_consilience, \
tier_write_components) \
VALUES ('p1', 'default', 100.0, 2.0, 0.5, 0.0, 1.0, \
'[\"support_mass\",\"attack_mass\",\"self_gen_local\",\"modality_consilience\"]')",
[],
)
.unwrap();
let (s, a, psi_l, chi): (f64, f64, f64, f64) = db
.conn()
.query_row(
"SELECT support_mass, attack_mass, self_gen_local, modality_consilience \
FROM mobility_state WHERE proposition_id='p1'",
[],
|row| Ok((row.get(0)?, row.get(1)?, row.get(2)?, row.get(3)?)),
)
.unwrap();
assert_eq!(s, 2.0);
assert_eq!(a, 0.5);
assert_eq!(psi_l, 0.0);
assert_eq!(chi, 1.0);
let ancestral: Option<f64> = db
.conn()
.query_row(
"SELECT self_gen_ancestral FROM mobility_state WHERE proposition_id='p1'",
[],
|row| row.get(0),
)
.unwrap();
assert!(
ancestral.is_none(),
"untouched background components should remain NULL"
);
}
fn seed_mobility_claim(
db: &YantrikDB,
proposition_id: &str,
extractor: &str,
polarity: i32,
source_lineage_json: &str,
self_gen: i32,
modality: &str,
regime: &str,
) {
let claim_id = format!("c_{}", uuid_like(extractor, source_lineage_json, polarity));
db.conn()
.execute(
"INSERT INTO claims (claim_id, src, dst, rel_type, created_at, \
extractor, polarity, namespace, proposition_id, regime_tag, \
self_generated, source_lineage, modality_signal, weight) \
VALUES (?1, 'X', 'Y', 'rel', 0.0, ?2, ?3, 'default', ?4, ?5, ?6, ?7, ?8, 1.0)",
rusqlite::params![
claim_id,
extractor,
polarity,
proposition_id,
regime,
self_gen,
source_lineage_json,
modality,
],
)
.unwrap();
}
fn uuid_like(a: &str, b: &str, pol: i32) -> String {
use std::collections::hash_map::DefaultHasher;
use std::hash::{Hash, Hasher};
let mut h = DefaultHasher::new();
(a, b, pol).hash(&mut h);
format!("c_{:x}", h.finish())
}
fn seed_proposition(db: &YantrikDB, pid: &str) {
let src = format!("src_{}", pid);
let rel = format!("rel_{}", pid);
let dst = format!("dst_{}", pid);
db.conn()
.execute(
"INSERT OR IGNORE INTO propositions (proposition_id, src, rel_type, dst, namespace, created_at) \
VALUES (?1, ?2, ?3, ?4, 'default', 0.0)",
rusqlite::params![pid, src, rel, dst],
)
.unwrap();
}
#[test]
fn test_mobility_single_support_claim() {
let db = YantrikDB::new(":memory:", 8).unwrap();
seed_proposition(&db, "p_single");
seed_mobility_claim(
&db,
"p_single",
"ext_a",
1,
"[\"src_1\"]",
0,
"text",
"default",
);
let state = db
.compute_write_tier_mobility("p_single", "default")
.unwrap();
assert_eq!(state.support_mass, Some(1.0));
assert_eq!(state.attack_mass, Some(0.0));
assert_eq!(state.self_gen_local, Some(0.0));
assert!((state.modality_consilience.unwrap() - 1.0 / 6.0).abs() < 1e-9);
}
#[test]
fn test_mobility_three_independent_sources() {
let db = YantrikDB::new(":memory:", 8).unwrap();
seed_proposition(&db, "p_ind");
seed_mobility_claim(
&db,
"p_ind",
"ext_a",
1,
"[\"src_a\"]",
0,
"text",
"default",
);
seed_mobility_claim(
&db,
"p_ind",
"ext_b",
1,
"[\"src_b\"]",
0,
"image",
"default",
);
seed_mobility_claim(
&db,
"p_ind",
"ext_c",
1,
"[\"src_c\"]",
0,
"numeric",
"default",
);
let state = db.compute_write_tier_mobility("p_ind", "default").unwrap();
assert!(
(state.support_mass.unwrap() - 3.0).abs() < 1e-9,
"expected 3.0 for independent claims, got {:?}",
state.support_mass
);
assert!((state.modality_consilience.unwrap() - 0.5).abs() < 1e-9);
}
#[test]
fn test_mobility_duplicate_sources_discounted() {
let db = YantrikDB::new(":memory:", 8).unwrap();
seed_proposition(&db, "p_dup");
seed_mobility_claim(
&db,
"p_dup",
"ext_a",
1,
"[\"src_shared\"]",
0,
"text",
"default",
);
seed_mobility_claim(
&db,
"p_dup",
"ext_b",
1,
"[\"src_shared\"]",
0,
"text",
"default",
);
seed_mobility_claim(
&db,
"p_dup",
"ext_c",
1,
"[\"src_shared\"]",
0,
"text",
"default",
);
let state = db.compute_write_tier_mobility("p_dup", "default").unwrap();
let s = state.support_mass.unwrap();
assert!(
s < 2.5,
"shared-source claims should discount below 2.5, got {}",
s
);
assert!(s > 1.8, "discount shouldn't be excessive, got {}", s);
}
#[test]
fn test_mobility_support_and_attack_separate() {
let db = YantrikDB::new(":memory:", 8).unwrap();
seed_proposition(&db, "p_mixed");
seed_mobility_claim(
&db,
"p_mixed",
"ext_a",
1,
"[\"src_a\"]",
0,
"text",
"default",
);
seed_mobility_claim(
&db,
"p_mixed",
"ext_b",
1,
"[\"src_b\"]",
0,
"image",
"default",
);
seed_mobility_claim(
&db,
"p_mixed",
"ext_c",
-1,
"[\"src_c\"]",
0,
"text",
"default",
);
let state = db
.compute_write_tier_mobility("p_mixed", "default")
.unwrap();
assert!((state.support_mass.unwrap() - 2.0).abs() < 1e-9);
assert!((state.attack_mass.unwrap() - 1.0).abs() < 1e-9);
assert_eq!(state.self_gen_local, Some(0.0));
}
#[test]
fn test_mobility_upsert_and_read_back() {
let db = YantrikDB::new(":memory:", 8).unwrap();
seed_proposition(&db, "p_rw");
seed_mobility_claim(&db, "p_rw", "ext_a", 1, "[\"src_1\"]", 0, "text", "default");
let computed = db.compute_write_tier_mobility("p_rw", "default").unwrap();
db.upsert_mobility_state(&computed).unwrap();
let read_back = db.get_mobility_state("p_rw", "default").unwrap().unwrap();
assert_eq!(read_back.support_mass, computed.support_mass);
assert_eq!(read_back.attack_mass, computed.attack_mass);
assert_eq!(read_back.tier_write_components.len(), 4);
assert!(read_back.self_gen_ancestral.is_none());
assert!(read_back.novelty_isolation.is_none());
}
#[test]
fn test_mobility_missing_returns_none() {
let db = YantrikDB::new(":memory:", 8).unwrap();
let result = db.get_mobility_state("nonexistent", "default").unwrap();
assert!(result.is_none());
}
#[test]
fn test_mobility_self_generated_discount() {
let db = YantrikDB::new(":memory:", 8).unwrap();
seed_proposition(&db, "p_self");
seed_mobility_claim(
&db,
"p_self",
"self_mode_a",
1,
"[\"self_1\"]",
1,
"text",
"default",
);
seed_mobility_claim(
&db,
"p_self",
"self_mode_b",
1,
"[\"self_1\"]",
1,
"text",
"default",
);
let state = db.compute_write_tier_mobility("p_self", "default").unwrap();
let s = state.support_mass.unwrap();
assert!(
s < 1.2,
"self-gen shared-lineage claims should collapse, got {}",
s
);
assert_eq!(state.self_gen_local, Some(1.0));
}
#[test]
fn test_m3_schema_v21_columns_present() {
let db = YantrikDB::new(":memory:", 8).unwrap();
let conn = db.conn();
let mut stmt = conn.prepare("PRAGMA table_info(mobility_state)").unwrap();
let rows = stmt
.query_map([], |r| r.get::<_, String>(1))
.unwrap()
.collect::<std::result::Result<Vec<_>, _>>()
.unwrap();
for expected in [
"formula_version",
"content_hash",
"live_claim_count",
"state_status",
"computed_at",
] {
assert!(
rows.iter().any(|c| c == expected),
"V21 column {} missing from mobility_state",
expected
);
}
}
#[test]
fn test_m3_state_populated_with_hash_and_status() {
let db = YantrikDB::new(":memory:", 8).unwrap();
seed_proposition(&db, "p_hash");
seed_mobility_claim(
&db,
"p_hash",
"ext_a",
1,
"[\"src_1\"]",
0,
"text",
"default",
);
let state = db.compute_write_tier_mobility("p_hash", "default").unwrap();
assert_eq!(
state.formula_version,
crate::engine::warrant::FORMULA_VERSION
);
assert!(
!state.content_hash.is_empty(),
"content_hash should be populated"
);
assert_eq!(state.state_status, "fresh");
assert_eq!(state.live_claim_count, 1);
assert!(state.computed_at > 0, "computed_at should be set");
let read = db.get_mobility_state("p_hash", "default").unwrap().unwrap();
assert_eq!(read.content_hash, state.content_hash);
assert_eq!(read.state_status, "fresh");
assert_eq!(read.live_claim_count, 1);
}
#[test]
fn test_m3_idempotent_recompute_on_unchanged_live_set() {
let db = YantrikDB::new(":memory:", 8).unwrap();
seed_proposition(&db, "p_idem");
seed_mobility_claim(
&db,
"p_idem",
"ext_a",
1,
"[\"src_1\"]",
0,
"text",
"default",
);
let first = db.compute_write_tier_mobility("p_idem", "default").unwrap();
let second = db.compute_write_tier_mobility("p_idem", "default").unwrap();
assert_eq!(
first.content_hash, second.content_hash,
"hash must be stable on unchanged set"
);
assert_eq!(
first.snapshot_ts, second.snapshot_ts,
"idempotent call should return the same row"
);
}
#[test]
fn test_m3_hash_discriminates_on_claim_change() {
let db = YantrikDB::new(":memory:", 8).unwrap();
seed_proposition(&db, "p_disc");
seed_mobility_claim(
&db,
"p_disc",
"ext_a",
1,
"[\"src_a\"]",
0,
"text",
"default",
);
let h1 = db
.compute_write_tier_mobility("p_disc", "default")
.unwrap()
.content_hash;
seed_mobility_claim(
&db,
"p_disc",
"ext_b",
1,
"[\"src_b\"]",
0,
"image",
"default",
);
let h2 = db
.compute_write_tier_mobility("p_disc", "default")
.unwrap()
.content_hash;
assert_ne!(h1, h2, "content_hash must change when live set changes");
}
#[test]
fn test_m3_order_invariant_through_db() {
fn setup_and_compute(order: &[(&str, &str, &str)]) -> (f64, String) {
let db = YantrikDB::new(":memory:", 8).unwrap();
seed_proposition(&db, "p_ord");
for (ext, src, modality) in order {
let lineage = format!("[\"{}\"]", src);
seed_mobility_claim(&db, "p_ord", ext, 1, &lineage, 0, modality, "default");
}
let state = db.compute_write_tier_mobility("p_ord", "default").unwrap();
(state.support_mass.unwrap(), state.content_hash)
}
let (mass_abc, hash_abc) = setup_and_compute(&[
("ext_a", "src_a", "text"),
("ext_b", "src_b", "image"),
("ext_c", "src_c", "numeric"),
]);
let (mass_cba, hash_cba) = setup_and_compute(&[
("ext_c", "src_c", "numeric"),
("ext_b", "src_b", "image"),
("ext_a", "src_a", "text"),
]);
assert!(
(mass_abc - mass_cba).abs() < 1e-9,
"order invariance broken: {} vs {}",
mass_abc,
mass_cba
);
assert_eq!(hash_abc, hash_cba, "content_hash must be order-invariant");
}
#[test]
fn test_m3_ingest_claim_triggers_mobility() {
let db = YantrikDB::new(":memory:", 8).unwrap();
let claim_id = db
.ingest_claim(
"Alice", "works_at", "Acme", "default", 1, "asserted", None, None, "manual", None,
"medium", None, None, None, 1.0,
)
.unwrap();
assert!(!claim_id.is_empty());
let prop_id: String = db
.conn()
.query_row(
"SELECT proposition_id FROM propositions \
WHERE src = 'Alice' AND rel_type = 'works_at' AND dst = 'Acme' AND namespace = 'default'",
[],
|row| row.get(0),
)
.unwrap();
assert!(!prop_id.is_empty());
let claim_prop_id: String = db
.conn()
.query_row(
"SELECT proposition_id FROM claims WHERE claim_id = ?1",
rusqlite::params![claim_id],
|row| row.get(0),
)
.unwrap();
assert_eq!(claim_prop_id, prop_id);
let state = db.get_mobility_state(&prop_id, "default").unwrap().unwrap();
assert_eq!(state.state_status, "fresh");
assert_eq!(state.live_claim_count, 1);
assert!(state.content_hash.len() >= 32);
assert!((state.support_mass.unwrap() - 1.0).abs() < 1e-9);
}
fn seed_contest_claim(
db: &YantrikDB,
proposition_id: &str,
claim_id: &str,
extractor: &str,
polarity: i32,
source_lineage_json: &str,
source_memory_rid: Option<&str>,
namespace: &str,
valid_from: Option<f64>,
valid_to: Option<f64>,
) {
let src = format!("src_{}", proposition_id);
let dst = format!("dst_{}", proposition_id);
let rel = format!("rel_{}", proposition_id);
db.conn()
.execute(
"INSERT INTO claims (claim_id, src, dst, rel_type, created_at, \
extractor, polarity, namespace, proposition_id, regime_tag, \
self_generated, source_lineage, modality_signal, weight, \
source_memory_rid, valid_from, valid_to) \
VALUES (?1, ?2, ?3, ?4, 0.0, ?5, ?6, ?7, ?8, 'default', 0, \
?9, 'text', 1.0, ?10, ?11, ?12)",
rusqlite::params![
claim_id,
src,
dst,
rel,
extractor,
polarity,
namespace,
proposition_id,
source_lineage_json,
source_memory_rid,
valid_from,
valid_to,
],
)
.unwrap();
}
#[test]
fn test_m4_schema_v22_contest_state_table() {
let db = YantrikDB::new(":memory:", 8).unwrap();
let conn = db.conn();
let mut stmt = conn.prepare("PRAGMA table_info(contest_state)").unwrap();
let rows: Vec<String> = stmt
.query_map([], |r| r.get::<_, String>(1))
.unwrap()
.collect::<std::result::Result<Vec<_>, _>>()
.unwrap();
for expected in [
"proposition_id",
"regime",
"support_mass",
"attack_mass",
"support_effective_independence",
"attack_effective_independence",
"support_distinct_source_count",
"attack_distinct_source_count",
"same_source_opposite_polarity_count",
"same_artifact_extractor_polarity_conflict_count",
"temporal_overlap_conflict_count",
"temporal_separable_opposition_count",
"referent_schema_heterogeneity_count",
"heuristic_flags",
"derivation_version",
"content_hash",
"live_claim_count",
"state_status",
"computed_at",
] {
assert!(
rows.iter().any(|c| c == expected),
"contest_state column {} missing",
expected
);
}
}
#[test]
fn test_m4_contest_state_missing_returns_none() {
let db = YantrikDB::new(":memory:", 8).unwrap();
assert!(db
.get_contest_state("nonexistent", "default")
.unwrap()
.is_none());
}
#[test]
fn test_m4_contest_single_support_basic_fields() {
let db = YantrikDB::new(":memory:", 8).unwrap();
seed_proposition(&db, "p_c1");
seed_contest_claim(
&db,
"p_c1",
"c1",
"ext_a",
1,
"[\"src_1\"]",
None,
"default",
None,
None,
);
let c = db.compute_contest_state("p_c1", "default").unwrap();
assert_eq!(c.live_claim_count, 1);
assert_eq!(c.state_status, "fresh");
assert!(!c.content_hash.is_empty());
assert!((c.support_mass - 1.0).abs() < 1e-9);
assert_eq!(c.attack_mass, 0.0);
assert_eq!(c.support_distinct_source_count, 1);
assert_eq!(c.attack_distinct_source_count, 0);
assert_eq!(
c.heuristic_flags, 0,
"no flags should fire for a lone claim"
);
}
#[test]
fn test_m4_contest_same_source_opposite_polarity() {
let db = YantrikDB::new(":memory:", 8).unwrap();
seed_proposition(&db, "p_same_src");
seed_contest_claim(
&db,
"p_same_src",
"c_sup",
"ext_a",
1,
"[\"src_x\"]",
None,
"default",
None,
None,
);
seed_contest_claim(
&db,
"p_same_src",
"c_att",
"ext_b",
-1,
"[\"src_x\"]",
None,
"default",
None,
None,
);
let c = db.compute_contest_state("p_same_src", "default").unwrap();
assert_eq!(c.same_source_opposite_polarity_count, 1);
assert!(c.heuristic_flags & crate::engine::warrant::contest_flags::SAME_SOURCE_CONFLICT != 0);
}
#[test]
fn test_m4_contest_same_artifact_extractor_conflict() {
let db = YantrikDB::new(":memory:", 8).unwrap();
seed_proposition(&db, "p_artifact");
seed_contest_claim(
&db,
"p_artifact",
"c1",
"extractor_a",
1,
"[]",
Some("doc_42"),
"default",
None,
None,
);
seed_contest_claim(
&db,
"p_artifact",
"c2",
"extractor_b",
-1,
"[]",
Some("doc_42"),
"default",
None,
None,
);
let c = db.compute_contest_state("p_artifact", "default").unwrap();
assert_eq!(c.same_artifact_extractor_polarity_conflict_count, 1);
assert!(
c.heuristic_flags & crate::engine::warrant::contest_flags::SAME_ARTIFACT_EXTRACTOR_CONFLICT
!= 0
);
}
#[test]
fn test_m4_contest_same_artifact_same_extractor_is_not_conflict() {
let db = YantrikDB::new(":memory:", 8).unwrap();
seed_proposition(&db, "p_same_ext");
seed_contest_claim(
&db,
"p_same_ext",
"c1",
"extractor_a",
1,
"[]",
Some("doc_42"),
"default",
None,
None,
);
seed_contest_claim(
&db,
"p_same_ext",
"c2",
"extractor_a",
-1,
"[]",
Some("doc_42"),
"default",
None,
None,
);
let c = db.compute_contest_state("p_same_ext", "default").unwrap();
assert_eq!(
c.same_artifact_extractor_polarity_conflict_count, 0,
"same-extractor same-artifact conflict is not an EXTRACTOR conflict"
);
}
#[test]
fn test_m4_contest_temporal_separable_vs_overlap() {
let db = YantrikDB::new(":memory:", 8).unwrap();
seed_proposition(&db, "p_temporal");
seed_contest_claim(
&db,
"p_temporal",
"c_sup",
"ext_a",
1,
"[\"s1\"]",
None,
"default",
Some(0.0),
Some(10.0),
);
seed_contest_claim(
&db,
"p_temporal",
"c_att",
"ext_b",
-1,
"[\"s2\"]",
None,
"default",
Some(20.0),
Some(30.0),
);
let c = db.compute_contest_state("p_temporal", "default").unwrap();
assert_eq!(c.temporal_separable_opposition_count, 1);
assert_eq!(c.temporal_overlap_conflict_count, 0);
assert!(
c.heuristic_flags & crate::engine::warrant::contest_flags::PRESENT_TENSE_CONFLICT == 0,
"disjoint intervals should NOT set PRESENT_TENSE_CONFLICT"
);
}
#[test]
fn test_m4_contest_temporal_overlap_is_present_tense_conflict() {
let db = YantrikDB::new(":memory:", 8).unwrap();
seed_proposition(&db, "p_overlap");
seed_contest_claim(
&db,
"p_overlap",
"c_sup",
"ext_a",
1,
"[\"s1\"]",
None,
"default",
Some(0.0),
Some(20.0),
);
seed_contest_claim(
&db,
"p_overlap",
"c_att",
"ext_b",
-1,
"[\"s2\"]",
None,
"default",
Some(10.0),
Some(30.0),
);
let c = db.compute_contest_state("p_overlap", "default").unwrap();
assert_eq!(c.temporal_overlap_conflict_count, 1);
assert_eq!(c.temporal_separable_opposition_count, 0);
assert!(c.heuristic_flags & crate::engine::warrant::contest_flags::PRESENT_TENSE_CONFLICT != 0);
}
#[test]
fn test_m4_contest_referent_heterogeneity() {
let db = YantrikDB::new(":memory:", 8).unwrap();
db.conn()
.execute(
"INSERT INTO propositions (proposition_id, src, rel_type, dst, namespace, created_at) \
VALUES ('p_het', 'X', 'rel', 'Y', 'default', 0.0)",
[],
)
.unwrap();
seed_contest_claim(
&db, "p_het", "c1", "ext_a", 1, "[\"s1\"]", None, "ns_a", None, None,
);
seed_contest_claim(
&db, "p_het", "c2", "ext_b", 1, "[\"s2\"]", None, "ns_b", None, None,
);
let c = db.compute_contest_state("p_het", "default").unwrap();
assert!(
c.referent_schema_heterogeneity_count > 0,
"two namespaces should register heterogeneity"
);
assert!(
c.heuristic_flags & crate::engine::warrant::contest_flags::REFERENT_HETEROGENEITY_PRESENT
!= 0
);
}
#[test]
fn test_m4_contest_duplication_risk_flag() {
let db = YantrikDB::new(":memory:", 8).unwrap();
seed_proposition(&db, "p_dup_risk");
for i in 0..4 {
let cid = format!("c_dup_{}", i);
let ext = format!("ext_{}", i);
seed_contest_claim(
&db,
"p_dup_risk",
&cid,
&ext,
1,
"[\"shared\"]",
None,
"default",
None,
None,
);
}
let c = db.compute_contest_state("p_dup_risk", "default").unwrap();
assert!(
c.support_distinct_source_count == 1,
"all share the same lineage element"
);
assert!(
c.support_mass > 2.0,
"four-way shared lineage should still give σ > 2"
);
let _ = c.heuristic_flags;
}
#[test]
fn test_m4_contest_order_invariant_through_db() {
fn setup(order: &[(&str, i32, &str)]) -> (f64, i64, String) {
let db = YantrikDB::new(":memory:", 8).unwrap();
seed_proposition(&db, "p_ord");
for (ext, pol, lineage) in order {
let cid = format!("c_{}_{}", ext, pol);
seed_contest_claim(
&db, "p_ord", &cid, ext, *pol, lineage, None, "default", None, None,
);
}
let c = db.compute_contest_state("p_ord", "default").unwrap();
(
c.support_mass,
c.same_source_opposite_polarity_count,
c.content_hash,
)
}
let (m1, n1, h1) = setup(&[
("ext_a", 1, "[\"x\"]"),
("ext_b", -1, "[\"x\"]"),
("ext_c", 1, "[\"y\"]"),
]);
let (m2, n2, h2) = setup(&[
("ext_c", 1, "[\"y\"]"),
("ext_b", -1, "[\"x\"]"),
("ext_a", 1, "[\"x\"]"),
]);
assert!(
(m1 - m2).abs() < 1e-9,
"support_mass must be order-invariant"
);
assert_eq!(n1, n2, "counters must be order-invariant");
assert_eq!(h1, h2, "content_hash must be order-invariant");
}
#[test]
fn test_m4_contest_idempotent_recompute() {
let db = YantrikDB::new(":memory:", 8).unwrap();
seed_proposition(&db, "p_idem");
seed_contest_claim(
&db, "p_idem", "c1", "ext_a", 1, "[\"s1\"]", None, "default", None, None,
);
let first = db.compute_contest_state("p_idem", "default").unwrap();
let second = db.compute_contest_state("p_idem", "default").unwrap();
assert_eq!(first.content_hash, second.content_hash);
assert_eq!(
first.computed_at, second.computed_at,
"idempotent recompute should not re-stamp"
);
}
#[test]
fn test_m4_ingest_claim_triggers_contest_state() {
let db = YantrikDB::new(":memory:", 8).unwrap();
db.ingest_claim(
"Alice", "works_at", "Acme", "default", 1, "asserted", None, None, "manual", None,
"medium", None, None, None, 1.0,
)
.unwrap();
let prop_id: String = db
.conn()
.query_row(
"SELECT proposition_id FROM propositions \
WHERE src = 'Alice' AND rel_type = 'works_at' AND dst = 'Acme' AND namespace = 'default'",
[],
|row| row.get(0),
)
.unwrap();
let c = db.get_contest_state(&prop_id, "default").unwrap().unwrap();
assert_eq!(c.state_status, "fresh");
assert_eq!(c.live_claim_count, 1);
assert!(!c.content_hash.is_empty());
assert_eq!(
c.derivation_version,
crate::engine::warrant::CONTEST_DERIVATION_VERSION
);
}
#[test]
fn test_m4_contest_independence_matches_omega_sum() {
let db = YantrikDB::new(":memory:", 8).unwrap();
seed_proposition(&db, "p_ind");
for i in 0..3 {
let cid = format!("c_ind_{}", i);
let ext = format!("ext_{}", i);
seed_contest_claim(
&db,
"p_ind",
&cid,
&ext,
1,
"[\"shared\"]",
None,
"default",
None,
None,
);
}
let c = db.compute_contest_state("p_ind", "default").unwrap();
assert!(
(c.support_effective_independence - 2.0).abs() < 0.01,
"expected ≈ 2.0, got {}",
c.support_effective_independence
);
}
#[test]
fn test_m45_list_flagged_propositions_empty_when_nothing_flagged() {
let db = YantrikDB::new(":memory:", 8).unwrap();
seed_proposition(&db, "p_clean");
seed_contest_claim(
&db,
"p_clean",
"c1",
"ext_a",
1,
"[\"src_1\"]",
None,
"default",
None,
None,
);
db.compute_contest_state("p_clean", "default").unwrap();
let flagged = db
.list_flagged_propositions(
crate::engine::warrant::contest_flags::SAME_SOURCE_CONFLICT,
10,
)
.unwrap();
assert!(
flagged.is_empty(),
"clean proposition should not be flagged"
);
}
#[test]
fn test_m45_list_flagged_propositions_returns_matching() {
let db = YantrikDB::new(":memory:", 8).unwrap();
seed_proposition(&db, "p_flag");
seed_contest_claim(
&db,
"p_flag",
"c_s",
"ext_a",
1,
"[\"src_x\"]",
None,
"default",
None,
None,
);
seed_contest_claim(
&db,
"p_flag",
"c_a",
"ext_b",
-1,
"[\"src_x\"]",
None,
"default",
None,
None,
);
db.compute_contest_state("p_flag", "default").unwrap();
seed_proposition(&db, "p_clean");
seed_contest_claim(
&db,
"p_clean",
"c1",
"ext_a",
1,
"[\"src_y\"]",
None,
"default",
None,
None,
);
db.compute_contest_state("p_clean", "default").unwrap();
let flagged = db
.list_flagged_propositions(
crate::engine::warrant::contest_flags::SAME_SOURCE_CONFLICT,
10,
)
.unwrap();
assert_eq!(flagged.len(), 1);
assert_eq!(flagged[0].proposition_id, "p_flag");
}
#[test]
fn test_m45_list_flagged_propositions_combined_mask() {
let db = YantrikDB::new(":memory:", 8).unwrap();
seed_proposition(&db, "p_same_src");
seed_contest_claim(
&db,
"p_same_src",
"c_src_s",
"ext_a",
1,
"[\"src_x\"]",
None,
"default",
None,
None,
);
seed_contest_claim(
&db,
"p_same_src",
"c_src_a",
"ext_b",
-1,
"[\"src_x\"]",
None,
"default",
None,
None,
);
db.compute_contest_state("p_same_src", "default").unwrap();
seed_proposition(&db, "p_artifact");
seed_contest_claim(
&db,
"p_artifact",
"c_art_s",
"extractor_a",
1,
"[]",
Some("doc_1"),
"default",
None,
None,
);
seed_contest_claim(
&db,
"p_artifact",
"c_art_a",
"extractor_b",
-1,
"[]",
Some("doc_1"),
"default",
None,
None,
);
db.compute_contest_state("p_artifact", "default").unwrap();
let mask = crate::engine::warrant::contest_flags::SAME_SOURCE_CONFLICT
| crate::engine::warrant::contest_flags::SAME_ARTIFACT_EXTRACTOR_CONFLICT;
let flagged = db.list_flagged_propositions(mask, 10).unwrap();
assert_eq!(
flagged.len(),
2,
"combined mask should match both propositions"
);
let prop_ids: Vec<&str> = flagged.iter().map(|c| c.proposition_id.as_str()).collect();
assert!(prop_ids.contains(&"p_same_src"));
assert!(prop_ids.contains(&"p_artifact"));
}
#[test]
fn test_m45_list_flagged_propositions_zero_mask_returns_empty() {
let db = YantrikDB::new(":memory:", 8).unwrap();
seed_proposition(&db, "p_any");
seed_contest_claim(
&db,
"p_any",
"c_s",
"ext_a",
1,
"[\"src_x\"]",
None,
"default",
None,
None,
);
seed_contest_claim(
&db,
"p_any",
"c_a",
"ext_b",
-1,
"[\"src_x\"]",
None,
"default",
None,
None,
);
db.compute_contest_state("p_any", "default").unwrap();
let flagged = db.list_flagged_propositions(0, 10).unwrap();
assert!(flagged.is_empty(), "mask = 0 should match nothing");
}
#[test]
fn test_m45_inspect_contest_conflicts_missing_returns_none() {
let db = YantrikDB::new(":memory:", 8).unwrap();
let report = db
.inspect_contest_conflicts("nonexistent", "default")
.unwrap();
assert!(report.is_none());
}
#[test]
fn test_m45_inspect_contest_returns_exemplar_pairs() {
let db = YantrikDB::new(":memory:", 8).unwrap();
seed_proposition(&db, "p_insp");
seed_contest_claim(
&db,
"p_insp",
"c_support",
"ext_a",
1,
"[\"src_shared\"]",
None,
"default",
None,
None,
);
seed_contest_claim(
&db,
"p_insp",
"c_attack",
"ext_b",
-1,
"[\"src_shared\"]",
None,
"default",
None,
None,
);
db.compute_contest_state("p_insp", "default").unwrap();
let report = db
.inspect_contest_conflicts("p_insp", "default")
.unwrap()
.unwrap();
assert!(
report.heuristic_flags & crate::engine::warrant::contest_flags::SAME_SOURCE_CONFLICT != 0
);
assert_eq!(report.same_source_opposite_polarity_pairs.len(), 1);
let pair = &report.same_source_opposite_polarity_pairs[0];
assert!(
(pair.0 == "c_support" && pair.1 == "c_attack")
|| (pair.0 == "c_attack" && pair.1 == "c_support"),
"exemplar pair should contain the actual conflicting claim_ids, got ({}, {})",
pair.0,
pair.1
);
}
#[test]
fn test_m45_inspect_temporal_overlap_pairs() {
let db = YantrikDB::new(":memory:", 8).unwrap();
seed_proposition(&db, "p_temp");
seed_contest_claim(
&db,
"p_temp",
"c_s",
"ext_a",
1,
"[\"s1\"]",
None,
"default",
Some(0.0),
Some(20.0),
);
seed_contest_claim(
&db,
"p_temp",
"c_a",
"ext_b",
-1,
"[\"s2\"]",
None,
"default",
Some(10.0),
Some(30.0),
);
db.compute_contest_state("p_temp", "default").unwrap();
let report = db
.inspect_contest_conflicts("p_temp", "default")
.unwrap()
.unwrap();
assert_eq!(report.temporal_overlap_conflict_pairs.len(), 1);
let db2 = YantrikDB::new(":memory:", 8).unwrap();
seed_proposition(&db2, "p_sep");
seed_contest_claim(
&db2,
"p_sep",
"c_s",
"ext_a",
1,
"[\"s1\"]",
None,
"default",
Some(0.0),
Some(10.0),
);
seed_contest_claim(
&db2,
"p_sep",
"c_a",
"ext_b",
-1,
"[\"s2\"]",
None,
"default",
Some(20.0),
Some(30.0),
);
db2.compute_contest_state("p_sep", "default").unwrap();
let report2 = db2
.inspect_contest_conflicts("p_sep", "default")
.unwrap()
.unwrap();
assert_eq!(
report2.temporal_overlap_conflict_pairs.len(),
0,
"disjoint intervals should not appear as overlap conflicts"
);
}
#[test]
fn test_m45_end_to_end_ingest_then_flagged_query() {
let db = YantrikDB::new(":memory:", 8).unwrap();
db.ingest_claim(
"Alice", "works_at", "Acme", "default", 1, "asserted", None, None, "manual", None,
"medium", None, None, None, 1.0,
)
.unwrap();
db.ingest_claim(
"Alice",
"works_at",
"Acme",
"default",
-1,
"denied",
None,
None,
"alt_manual",
None,
"medium",
None,
None,
None,
1.0,
)
.unwrap();
let prop_id: String = db
.conn()
.query_row(
"SELECT proposition_id FROM propositions \
WHERE src = 'Alice' AND rel_type = 'works_at' AND dst = 'Acme' AND namespace = 'default'",
[],
|row| row.get(0),
)
.unwrap();
let state = db.get_contest_state(&prop_id, "default").unwrap().unwrap();
assert_eq!(state.live_claim_count, 2);
assert!(
state.heuristic_flags & crate::engine::warrant::contest_flags::PRESENT_TENSE_CONFLICT != 0,
"opposite-polarity claims without temporal bounds should flag as present-tense conflict"
);
let flagged = db
.list_flagged_propositions(
crate::engine::warrant::contest_flags::PRESENT_TENSE_CONFLICT,
10,
)
.unwrap();
assert_eq!(flagged.len(), 1);
assert_eq!(flagged[0].proposition_id, prop_id);
}
#[test]
fn test_m3_second_ingest_updates_mobility_deterministically() {
let db = YantrikDB::new(":memory:", 8).unwrap();
db.ingest_claim(
"Alice", "works_at", "Acme", "default", 1, "asserted", None, None, "source_a", None,
"medium", None, None, None, 1.0,
)
.unwrap();
let prop_id: String = db
.conn()
.query_row(
"SELECT proposition_id FROM propositions \
WHERE src = 'Alice' AND rel_type = 'works_at' AND dst = 'Acme' AND namespace = 'default'",
[],
|row| row.get(0),
)
.unwrap();
let h1 = db
.get_mobility_state(&prop_id, "default")
.unwrap()
.unwrap()
.content_hash;
db.ingest_claim(
"Alice", "works_at", "Acme", "default", 1, "asserted", None, None, "source_b", None,
"medium", None, None, None, 1.0,
)
.unwrap();
let state2 = db.get_mobility_state(&prop_id, "default").unwrap().unwrap();
assert_eq!(state2.live_claim_count, 2);
assert_ne!(h1, state2.content_hash, "hash must change after new claim");
assert!((state2.support_mass.unwrap() - 2.0).abs() < 1e-9);
}
#[test]
fn test_schema_v20_migration_from_v19() {
use rusqlite::Connection;
let conn = Connection::open_in_memory().unwrap();
conn.execute_batch(
"
CREATE TABLE propositions (
proposition_id TEXT PRIMARY KEY,
src TEXT NOT NULL, rel_type TEXT NOT NULL, dst TEXT NOT NULL,
namespace TEXT NOT NULL, created_at REAL NOT NULL,
UNIQUE(src, rel_type, dst, namespace)
);
CREATE TABLE claims (
claim_id TEXT PRIMARY KEY,
src TEXT NOT NULL, dst TEXT NOT NULL, rel_type TEXT NOT NULL,
weight REAL NOT NULL DEFAULT 1.0,
created_at REAL NOT NULL,
tombstoned INTEGER NOT NULL DEFAULT 0,
polarity INTEGER NOT NULL DEFAULT 1,
modality TEXT NOT NULL DEFAULT 'asserted',
valid_from REAL, valid_to REAL,
extractor TEXT NOT NULL DEFAULT 'manual',
extractor_version TEXT,
confidence_band TEXT NOT NULL DEFAULT 'medium',
source_memory_rid TEXT,
span_start INTEGER, span_end INTEGER,
namespace TEXT NOT NULL DEFAULT 'default',
proposition_id TEXT
);
",
)
.unwrap();
conn.execute(
"INSERT INTO propositions (proposition_id, src, rel_type, dst, namespace, created_at) \
VALUES ('p1', 'Alice', 'works_at', 'Acme', 'default', 0.0)",
[],
)
.unwrap();
conn.execute("INSERT INTO claims (claim_id, src, dst, rel_type, created_at, extractor, namespace, proposition_id) \
VALUES ('c1', 'Alice', 'Acme', 'works_at', 0.0, 'manual', 'default', 'p1')", []).unwrap();
conn.execute_batch(crate::schema::MIGRATE_V19_TO_V20)
.unwrap();
let count: i64 = conn
.query_row(
"SELECT COUNT(*) FROM sqlite_master WHERE type='table' AND name IN \
('mobility_state', 'actor_profile', 'compression_artifact')",
[],
|r| r.get(0),
)
.unwrap();
assert_eq!(count, 3);
let (regime, self_gen, lineage, modality): (String, i64, String, String) = conn
.query_row(
"SELECT regime_tag, self_generated, source_lineage, modality_signal \
FROM claims WHERE claim_id='c1'",
[],
|r| Ok((r.get(0)?, r.get(1)?, r.get(2)?, r.get(3)?)),
)
.unwrap();
assert_eq!(regime, "default");
assert_eq!(self_gen, 0);
assert_eq!(lineage, "[]");
assert_eq!(modality, "text");
}
#[test]
fn test_schema_v19_backfill_from_v18() {
use rusqlite::Connection;
let conn = Connection::open_in_memory().unwrap();
conn.execute_batch(
"
CREATE TABLE claims (
claim_id TEXT PRIMARY KEY,
src TEXT NOT NULL,
dst TEXT NOT NULL,
rel_type TEXT NOT NULL,
weight REAL NOT NULL DEFAULT 1.0,
created_at REAL NOT NULL,
tombstoned INTEGER NOT NULL DEFAULT 0,
polarity INTEGER NOT NULL DEFAULT 1,
modality TEXT NOT NULL DEFAULT 'asserted',
valid_from REAL, valid_to REAL,
extractor TEXT NOT NULL DEFAULT 'manual',
extractor_version TEXT,
confidence_band TEXT NOT NULL DEFAULT 'medium',
source_memory_rid TEXT,
span_start INTEGER, span_end INTEGER,
namespace TEXT NOT NULL DEFAULT 'default'
);
",
)
.unwrap();
conn.execute(
"INSERT INTO claims (claim_id, src, dst, rel_type, created_at, extractor, namespace) \
VALUES ('c1', 'Alice', 'Acme', 'works_at', 0.0, 'source_a', 'ns1')",
[],
)
.unwrap();
conn.execute(
"INSERT INTO claims (claim_id, src, dst, rel_type, created_at, extractor, namespace) \
VALUES ('c2', 'Alice', 'Acme', 'works_at', 0.0, 'source_b', 'ns1')",
[],
)
.unwrap();
conn.execute(
"INSERT INTO claims (claim_id, src, dst, rel_type, created_at, extractor, namespace) \
VALUES ('c3', 'Bob', 'Beta', 'works_at', 0.0, 'source_a', 'ns1')",
[],
)
.unwrap();
conn.execute("INSERT INTO claims (claim_id, src, dst, rel_type, created_at, extractor, namespace, tombstoned) \
VALUES ('c4', 'Carol', 'Gamma', 'works_at', 0.0, 'source_a', 'ns1', 1)", []).unwrap();
conn.execute_batch(crate::schema::MIGRATE_V18_TO_V19)
.unwrap();
let prop_count: i64 = conn
.query_row("SELECT COUNT(*) FROM propositions", [], |r| r.get(0))
.unwrap();
assert_eq!(
prop_count, 2,
"backfill should create one proposition per unique non-tombstoned tuple"
);
let alice_props: Vec<String> = conn
.prepare("SELECT proposition_id FROM claims WHERE src='Alice'")
.unwrap()
.query_map([], |r| r.get::<_, String>(0))
.unwrap()
.filter_map(|r| r.ok())
.collect();
assert_eq!(alice_props.len(), 2);
assert_eq!(
alice_props[0], alice_props[1],
"claims on the same tuple should share a proposition_id"
);
let bob_prop: String = conn
.query_row(
"SELECT proposition_id FROM claims WHERE src='Bob'",
[],
|r| r.get(0),
)
.unwrap();
assert_ne!(alice_props[0], bob_prop);
let carol_prop: Option<String> = conn
.query_row(
"SELECT proposition_id FROM claims WHERE src='Carol'",
[],
|r| r.get(0),
)
.unwrap();
assert!(
carol_prop.is_none(),
"tombstoned claims should not be backfilled"
);
}
#[test]
fn test_think_empty_db() {
let db = YantrikDB::new(":memory:", 8).unwrap();
let config = ThinkConfig {
run_consolidation: false,
run_conflict_scan: false,
run_pattern_mining: false,
..Default::default()
};
let result = db.think(&config).unwrap();
assert!(result.triggers.is_empty());
assert_eq!(result.consolidation_count, 0);
assert_eq!(result.conflicts_found, 0);
assert!(result.duration_ms < 5000);
}
#[test]
fn test_think_with_decayed_memories() {
let db = YantrikDB::new(":memory:", 8).unwrap();
let rid = db
.record(
"important deadline",
"episodic",
0.9,
0.0,
100.0,
&empty_meta(),
&vec_seed(1.0, 8),
"default",
0.8,
"general",
"user",
None,
)
.unwrap();
let ts = std::time::SystemTime::now()
.duration_since(std::time::UNIX_EPOCH)
.unwrap()
.as_secs_f64();
db.conn()
.execute(
"UPDATE memories SET last_access = ?1 WHERE rid = ?2",
rusqlite::params![ts - 10000.0, rid],
)
.unwrap();
let config = ThinkConfig {
run_consolidation: false,
run_conflict_scan: false,
run_pattern_mining: false,
..Default::default()
};
let result = db.think(&config).unwrap();
assert!(!result.triggers.is_empty());
assert_eq!(result.triggers[0].trigger_type, "decay_review");
}
#[test]
fn test_think_records_last_think_at() {
let db = YantrikDB::new(":memory:", 8).unwrap();
let config = ThinkConfig {
run_consolidation: false,
run_conflict_scan: false,
run_pattern_mining: false,
..Default::default()
};
db.think(&config).unwrap();
let val: String = db
.conn()
.query_row(
"SELECT value FROM meta WHERE key = 'last_think_at'",
[],
|row| row.get(0),
)
.unwrap();
let ts: f64 = val.parse().unwrap();
assert!(ts > 0.0);
}
#[test]
fn test_trigger_lifecycle() {
let db = YantrikDB::new(":memory:", 8).unwrap();
let trigger = crate::types::Trigger {
trigger_type: "decay_review".to_string(),
reason: "test".to_string(),
urgency: 0.8,
source_rids: vec!["rid-1".to_string()],
suggested_action: "test".to_string(),
context: std::collections::HashMap::new(),
};
let ts = std::time::SystemTime::now()
.duration_since(std::time::UNIX_EPOCH)
.unwrap()
.as_secs_f64();
let tid = crate::triggers::persist_trigger(&db, &trigger, ts)
.unwrap()
.unwrap();
let pending = db.get_pending_triggers(10).unwrap();
assert_eq!(pending.len(), 1);
assert_eq!(pending[0].status, "pending");
assert!(db.deliver_trigger(&tid).unwrap());
let history = db.get_trigger_history(None, 10).unwrap();
assert_eq!(history[0].status, "delivered");
assert!(db.acknowledge_trigger(&tid).unwrap());
assert!(db.act_on_trigger(&tid).unwrap());
let history = db.get_trigger_history(None, 10).unwrap();
assert_eq!(history[0].status, "acted");
}
#[test]
fn test_stats_include_triggers_and_patterns() {
let db = YantrikDB::new(":memory:", 8).unwrap();
let s = db.stats(None).unwrap();
assert_eq!(s.pending_triggers, 0);
assert_eq!(s.active_patterns, 0);
}
#[test]
fn test_entity_type_stored_on_relate() {
let db = YantrikDB::new(":memory:", 8).unwrap();
db.relate("Sarah", "Mike", "knows", 1.0).unwrap();
db.relate("Sarah", "Flipkart", "works_at", 1.0).unwrap();
db.relate("Sarah", "Bangalore", "lives_in", 1.0).unwrap();
db.relate("FAISS", "recommendation engine", "used_in", 1.0)
.unwrap();
let etype: String = db
.conn()
.query_row(
"SELECT entity_type FROM entities WHERE name = 'Sarah'",
[],
|r| r.get(0),
)
.unwrap();
assert_eq!(etype, "person");
let etype: String = db
.conn()
.query_row(
"SELECT entity_type FROM entities WHERE name = 'Mike'",
[],
|r| r.get(0),
)
.unwrap();
assert_eq!(etype, "person");
let etype: String = db
.conn()
.query_row(
"SELECT entity_type FROM entities WHERE name = 'Flipkart'",
[],
|r| r.get(0),
)
.unwrap();
assert_eq!(etype, "organization");
let etype: String = db
.conn()
.query_row(
"SELECT entity_type FROM entities WHERE name = 'Bangalore'",
[],
|r| r.get(0),
)
.unwrap();
assert_eq!(etype, "place");
let etype: String = db
.conn()
.query_row(
"SELECT entity_type FROM entities WHERE name = 'FAISS'",
[],
|r| r.get(0),
)
.unwrap();
assert_eq!(etype, "tech");
let etype: String = db
.conn()
.query_row(
"SELECT entity_type FROM entities WHERE name = 'recommendation engine'",
[],
|r| r.get(0),
)
.unwrap();
assert_eq!(etype, "unknown");
}
#[test]
fn test_recall_deterministic_with_skip_reinforce() {
let db = YantrikDB::new(":memory:", 8).unwrap();
for i in 0..10 {
db.record(
&format!("memory {i}"),
"episodic",
0.5,
0.0,
604800.0,
&empty_meta(),
&vec_seed(i as f32, 8),
"default",
0.8,
"general",
"user",
None,
)
.unwrap();
}
let query = vec_seed(3.0, 8);
let r1 = db
.recall(
&query, 5, None, None, false, false, None, true, None, None, None, None, None,
)
.unwrap();
let r2 = db
.recall(
&query, 5, None, None, false, false, None, true, None, None, None, None, None,
)
.unwrap();
let r3 = db
.recall(
&query, 5, None, None, false, false, None, true, None, None, None, None, None,
)
.unwrap();
let rids1: Vec<&str> = r1.iter().map(|r| r.rid.as_str()).collect();
let rids2: Vec<&str> = r2.iter().map(|r| r.rid.as_str()).collect();
let rids3: Vec<&str> = r3.iter().map(|r| r.rid.as_str()).collect();
assert_eq!(rids1, rids2);
assert_eq!(rids2, rids3);
for i in 0..5 {
assert!(
(r1[i].score - r2[i].score).abs() < 1e-4,
"score drift too large between calls: {} vs {}",
r1[i].score,
r2[i].score
);
}
}
#[test]
fn test_reinforce_mutates_but_skip_does_not() {
let db = YantrikDB::new(":memory:", 8).unwrap();
let rid = db
.record(
"test",
"episodic",
0.5,
0.0,
1000.0,
&empty_meta(),
&vec_seed(1.0, 8),
"default",
0.8,
"general",
"user",
None,
)
.unwrap();
let original_hl = db.get(&rid).unwrap().unwrap().half_life;
db.recall(
&vec_seed(1.0, 8),
1,
None,
None,
false,
false,
None,
true,
None,
None,
None,
None,
None,
)
.unwrap();
let after_skip = db.get(&rid).unwrap().unwrap().half_life;
assert!((original_hl - after_skip).abs() < 1e-10);
db.recall(
&vec_seed(1.0, 8),
1,
None,
None,
false,
false,
None,
false,
None,
None,
None,
None,
None,
)
.unwrap();
let after_reinforce = db.get(&rid).unwrap().unwrap().half_life;
assert!(after_reinforce > original_hl);
}
#[test]
fn test_graph_expansion_off_no_graph_results() {
let db = YantrikDB::new(":memory:", 8).unwrap();
let r1 = db
.record(
"Alice discussed plan",
"episodic",
0.5,
0.0,
604800.0,
&empty_meta(),
&vec_seed(1.0, 8),
"default",
0.8,
"general",
"user",
None,
)
.unwrap();
let r2 = db
.record(
"Bob reviewed code",
"episodic",
0.5,
0.0,
604800.0,
&empty_meta(),
&vec_seed(5.0, 8),
"default",
0.8,
"general",
"user",
None,
)
.unwrap();
db.relate("Alice", "Bob", "knows", 1.0).unwrap();
db.link_memory_entity(&r1, "Alice").unwrap();
db.link_memory_entity(&r2, "Bob").unwrap();
let results = db
.recall(
&vec_seed(1.0, 8),
10,
None,
None,
false,
false,
Some("Alice"),
false,
None,
None,
None,
None,
None,
)
.unwrap();
for r in &results {
assert!(
(r.scores.graph_proximity - 0.0).abs() < 1e-10,
"graph_proximity should be 0.0 when expansion is disabled"
);
}
}
#[test]
fn test_graph_expansion_on_boosts_connected_memory() {
let db = YantrikDB::new(":memory:", 8).unwrap();
let r1 = db
.record(
"Alice discussed the project plan",
"episodic",
0.5,
0.0,
604800.0,
&empty_meta(),
&vec_seed(1.0, 8),
"default",
0.8,
"general",
"user",
None,
)
.unwrap();
let r2 = db
.record(
"Bob reviewed the code",
"episodic",
0.5,
0.0,
604800.0,
&empty_meta(),
&vec_seed(5.0, 8),
"default",
0.8,
"general",
"user",
None,
)
.unwrap();
db.relate("Alice", "Bob", "knows", 1.0).unwrap();
db.link_memory_entity(&r1, "Alice").unwrap();
db.link_memory_entity(&r2, "Bob").unwrap();
let results = db
.recall(
&vec_seed(1.0, 8),
10,
None,
None,
false,
true,
Some("What is Alice working on?"),
true,
None,
None,
None,
None,
None,
)
.unwrap();
let alice_result = results.iter().find(|r| r.rid == r1).unwrap();
assert!(
alice_result.scores.graph_proximity > 0.0,
"Alice memory should have graph proximity when expansion is on"
);
}
#[test]
fn test_backfill_uses_word_boundaries() {
let db = YantrikDB::new(":memory:", 8).unwrap();
db.relate("data", "pipeline", "part_of", 1.0).unwrap();
let r1 = db
.record(
"the data is clean",
"episodic",
0.5,
0.0,
604800.0,
&empty_meta(),
&vec_seed(1.0, 8),
"default",
0.8,
"general",
"user",
None,
)
.unwrap();
let r2 = db
.record(
"the database is fast",
"episodic",
0.5,
0.0,
604800.0,
&empty_meta(),
&vec_seed(2.0, 8),
"default",
0.8,
"general",
"user",
None,
)
.unwrap();
let _count = db.backfill_memory_entities().unwrap();
let linked_to_data: Vec<String> = db
.conn()
.prepare("SELECT memory_rid FROM memory_entities WHERE entity_name = 'data'")
.unwrap()
.query_map([], |row| row.get(0))
.unwrap()
.collect::<std::result::Result<Vec<_>, _>>()
.unwrap();
assert!(
linked_to_data.contains(&r1),
"memory with 'data' as word should be linked"
);
assert!(
!linked_to_data.contains(&r2),
"memory with 'database' should NOT be linked (word boundary)"
);
}
#[test]
fn test_recall_scores_bounded() {
let db = YantrikDB::new(":memory:", 8).unwrap();
for i in 0..10 {
db.record(
&format!("memory {i}"),
"episodic",
(i as f64) * 0.1, ((i as f64) - 5.0) * 0.2, 604800.0,
&empty_meta(),
&vec_seed(i as f32, 8),
"default",
0.8,
"general",
"user",
None,
)
.unwrap();
}
let results = db
.recall(
&vec_seed(5.0, 8),
10,
None,
None,
false,
false,
None,
true,
None,
None,
None,
None,
None,
)
.unwrap();
for r in &results {
assert!(
r.score >= 0.0,
"score should be non-negative, got {}",
r.score
);
assert!(
r.score < 5.0,
"score should be reasonably bounded, got {}",
r.score
);
assert!(r.scores.similarity >= -1.0 && r.scores.similarity <= 1.0);
assert!(r.scores.decay >= 0.0 && r.scores.decay <= 1.0);
assert!(r.scores.recency >= 0.0 && r.scores.recency <= 1.0);
}
}
#[test]
fn test_link_memory_entity_idempotent() {
let db = YantrikDB::new(":memory:", 8).unwrap();
let rid = db
.record(
"test",
"episodic",
0.5,
0.0,
604800.0,
&empty_meta(),
&vec_seed(1.0, 8),
"default",
0.8,
"general",
"user",
None,
)
.unwrap();
db.relate("Alice", "Bob", "knows", 1.0).unwrap();
db.link_memory_entity(&rid, "Alice").unwrap();
db.link_memory_entity(&rid, "Alice").unwrap();
let count: i64 = db
.conn()
.query_row(
"SELECT COUNT(*) FROM memory_entities WHERE memory_rid = ?1 AND entity_name = 'Alice'",
params![rid],
|r| r.get(0),
)
.unwrap();
assert_eq!(count, 1);
}
#[test]
fn test_schema_v5_has_memory_entities() {
let db = YantrikDB::new(":memory:", 8).unwrap();
let count: i64 = db
.conn()
.query_row(
"SELECT COUNT(*) FROM sqlite_master WHERE type='table' AND name='memory_entities'",
[],
|r| r.get(0),
)
.unwrap();
assert_eq!(count, 1);
}
#[test]
fn test_recall_top_k_respected_with_graph_expansion() {
let db = YantrikDB::new(":memory:", 8).unwrap();
for i in 0..20 {
let rid = db
.record(
&format!("memory about topic {i}"),
"episodic",
0.5,
0.0,
604800.0,
&empty_meta(),
&vec_seed(i as f32, 8),
"default",
0.8,
"general",
"user",
None,
)
.unwrap();
let entity = format!("Entity{i}");
db.relate(
&entity,
&format!("Entity{}", (i + 1) % 20),
"related_to",
1.0,
)
.unwrap();
db.link_memory_entity(&rid, &entity).unwrap();
}
let results = db
.recall(
&vec_seed(0.0, 8),
5,
None,
None,
false,
true,
Some("Entity0 topic"),
true,
None,
None,
None,
None,
None,
)
.unwrap();
assert!(
results.len() <= 5,
"results should not exceed top_k=5, got {}",
results.len()
);
}
#[test]
fn test_schema_v6_has_storage_tier() {
let db = YantrikDB::new(":memory:", 8).unwrap();
let rid = db
.record(
"tier test",
"episodic",
0.5,
0.0,
604800.0,
&empty_meta(),
&vec_seed(1.0, 8),
"default",
0.8,
"general",
"user",
None,
)
.unwrap();
let mem = db.get(&rid).unwrap().unwrap();
assert_eq!(mem.storage_tier, "hot");
}
#[test]
fn test_schema_v7_has_fts5_and_join_tables() {
let db = YantrikDB::new(":memory:", 8).unwrap();
let _rid = db
.record(
"The quick brown fox jumps over the lazy dog",
"episodic",
0.5,
0.0,
604800.0,
&empty_meta(),
&vec_seed(1.0, 8),
"default",
0.8,
"general",
"user",
None,
)
.unwrap();
let conn = db.conn();
let count: i64 = conn
.query_row(
"SELECT COUNT(*) FROM memories_fts WHERE memories_fts MATCH 'quick brown'",
[],
|row| row.get(0),
)
.unwrap();
assert_eq!(count, 1, "FTS5 should index inserted memory");
let _: i64 = conn
.query_row("SELECT COUNT(*) FROM trigger_source_rids", [], |row| {
row.get(0)
})
.unwrap();
let _: i64 = conn
.query_row("SELECT COUNT(*) FROM pattern_evidence", [], |row| {
row.get(0)
})
.unwrap();
let _: i64 = conn
.query_row("SELECT COUNT(*) FROM pattern_entities", [], |row| {
row.get(0)
})
.unwrap();
let ver: String = conn
.query_row(
"SELECT value FROM meta WHERE key = 'schema_version'",
[],
|row| row.get(0),
)
.unwrap();
assert_eq!(ver, crate::schema::SCHEMA_VERSION.to_string());
}
#[test]
fn test_fts5_search_multiple_memories() {
let db = YantrikDB::new(":memory:", 8).unwrap();
db.record(
"Alice loves Rust programming",
"semantic",
0.5,
0.0,
604800.0,
&empty_meta(),
&vec_seed(1.0, 8),
"default",
0.8,
"general",
"user",
None,
)
.unwrap();
db.record(
"Bob prefers Python scripting",
"semantic",
0.5,
0.0,
604800.0,
&empty_meta(),
&vec_seed(0.5, 8),
"default",
0.8,
"general",
"user",
None,
)
.unwrap();
db.record(
"Alice and Bob work on Rust projects",
"episodic",
0.5,
0.0,
604800.0,
&empty_meta(),
&vec_seed(0.3, 8),
"default",
0.8,
"general",
"user",
None,
)
.unwrap();
let conn = db.conn();
let count: i64 = conn
.query_row(
"SELECT COUNT(*) FROM memories_fts WHERE memories_fts MATCH 'rust'",
[],
|row| row.get(0),
)
.unwrap();
assert_eq!(count, 2, "FTS5 should find 2 memories containing 'rust'");
let count: i64 = conn
.query_row(
"SELECT COUNT(*) FROM memories_fts WHERE memories_fts MATCH 'alice'",
[],
|row| row.get(0),
)
.unwrap();
assert_eq!(count, 2, "FTS5 should find 2 memories containing 'alice'");
let count: i64 = conn
.query_row(
"SELECT COUNT(*) FROM memories_fts WHERE memories_fts MATCH 'python'",
[],
|row| row.get(0),
)
.unwrap();
assert_eq!(count, 1);
}
#[test]
fn test_archive_memory() {
let db = YantrikDB::new(":memory:", 8).unwrap();
let rid = db
.record(
"to archive",
"episodic",
0.5,
0.0,
604800.0,
&empty_meta(),
&vec_seed(1.0, 8),
"default",
0.8,
"general",
"user",
None,
)
.unwrap();
assert!(db.archive(&rid).unwrap());
let mem = db.get(&rid).unwrap().unwrap();
assert_eq!(mem.storage_tier, "cold");
let results = db
.recall(
&vec_seed(1.0, 8),
10,
None,
None,
false,
false,
None,
true,
None,
None,
None,
None,
None,
)
.unwrap();
assert!(
results.iter().all(|r| r.rid != rid),
"archived memory should not appear in recall"
);
assert_eq!(db.stats(None).unwrap().archived_memories, 1);
}
#[test]
fn test_hydrate_memory() {
let db = YantrikDB::new(":memory:", 8).unwrap();
let emb = vec_seed(2.0, 8);
let rid = db
.record(
"to hydrate",
"episodic",
0.5,
0.0,
604800.0,
&empty_meta(),
&emb,
"default",
0.8,
"general",
"user",
None,
)
.unwrap();
db.archive(&rid).unwrap();
assert!(db.hydrate(&rid).unwrap());
let mem = db.get(&rid).unwrap().unwrap();
assert_eq!(mem.storage_tier, "hot");
let results = db
.recall(
&emb, 10, None, None, false, false, None, true, None, None, None, None, None,
)
.unwrap();
assert!(
results.iter().any(|r| r.rid == rid),
"hydrated memory should appear in recall"
);
assert_eq!(db.stats(None).unwrap().archived_memories, 0);
}
#[test]
fn test_archive_idempotent() {
let db = YantrikDB::new(":memory:", 8).unwrap();
let rid = db
.record(
"idempotent",
"episodic",
0.5,
0.0,
604800.0,
&empty_meta(),
&vec_seed(1.0, 8),
"default",
0.8,
"general",
"user",
None,
)
.unwrap();
assert!(db.archive(&rid).unwrap());
assert!(!db.archive(&rid).unwrap()); }
#[test]
fn test_record_batch() {
let db = YantrikDB::new(":memory:", 8).unwrap();
let inputs: Vec<RecordInput> = (0..10)
.map(|i| RecordInput {
text: format!("batch memory {i}"),
memory_type: "episodic".to_string(),
importance: 0.5,
valence: 0.0,
half_life: 604800.0,
metadata: serde_json::json!({}),
embedding: vec_seed(i as f32, 8),
namespace: "default".to_string(),
certainty: 0.8,
domain: "general".to_string(),
source: "user".to_string(),
emotional_state: None,
})
.collect();
let rids = db.record_batch(&inputs).unwrap();
assert_eq!(rids.len(), 10);
for rid in &rids {
assert!(db.get(rid).unwrap().is_some());
}
assert_eq!(db.stats(None).unwrap().active_memories, 10);
}
#[test]
fn test_record_batch_empty() {
let db = YantrikDB::new(":memory:", 8).unwrap();
let rids = db.record_batch(&[]).unwrap();
assert!(rids.is_empty());
}
#[test]
fn test_evict() {
let db = YantrikDB::new(":memory:", 8).unwrap();
for i in 0..20 {
db.record(
&format!("evict mem {i}"),
"episodic",
0.5,
0.0,
604800.0,
&empty_meta(),
&vec_seed(i as f32, 8),
"default",
0.8,
"general",
"user",
None,
)
.unwrap();
}
assert_eq!(db.stats(None).unwrap().active_memories, 20);
let archived = db.evict(10).unwrap();
assert_eq!(archived.len(), 10);
let stats = db.stats(None).unwrap();
assert_eq!(stats.archived_memories, 10);
let results = db
.recall(
&vec_seed(0.0, 8),
20,
None,
None,
false,
false,
None,
true,
None,
None,
None,
None,
None,
)
.unwrap();
for r in &results {
assert!(
!archived.contains(&r.rid),
"evicted memory should not be in recall"
);
}
}
#[test]
fn test_evict_no_action_when_under_limit() {
let db = YantrikDB::new(":memory:", 8).unwrap();
for i in 0..5 {
db.record(
&format!("small db {i}"),
"episodic",
0.5,
0.0,
604800.0,
&empty_meta(),
&vec_seed(i as f32, 8),
"default",
0.8,
"general",
"user",
None,
)
.unwrap();
}
let archived = db.evict(10).unwrap();
assert!(archived.is_empty());
}
#[test]
fn test_query_builder_basic() {
let db = YantrikDB::new(":memory:", 8).unwrap();
for i in 0..10 {
db.record(
&format!("memory {i}"),
"episodic",
0.5,
0.0,
604800.0,
&empty_meta(),
&vec_seed(i as f32, 8),
"default",
0.8,
"general",
"user",
None,
)
.unwrap();
}
let results = db
.query(RecallQuery::new(vec_seed(0.0, 8)).top_k(3).skip_reinforce())
.unwrap();
assert_eq!(results.len(), 3);
}
#[test]
fn test_query_builder_with_filters() {
let db = YantrikDB::new(":memory:", 8).unwrap();
db.record(
"episodic one",
"episodic",
0.5,
0.0,
604800.0,
&empty_meta(),
&vec_seed(1.0, 8),
"work",
0.8,
"general",
"user",
None,
)
.unwrap();
db.record(
"semantic one",
"semantic",
0.5,
0.0,
604800.0,
&empty_meta(),
&vec_seed(2.0, 8),
"work",
0.8,
"general",
"user",
None,
)
.unwrap();
db.record(
"episodic two",
"episodic",
0.5,
0.0,
604800.0,
&empty_meta(),
&vec_seed(3.0, 8),
"personal",
0.8,
"general",
"user",
None,
)
.unwrap();
let results = db
.query(
RecallQuery::new(vec_seed(1.0, 8))
.top_k(10)
.memory_type("episodic")
.namespace("work")
.skip_reinforce(),
)
.unwrap();
assert_eq!(results.len(), 1);
assert_eq!(results[0].memory_type, "episodic");
assert_eq!(results[0].namespace, "work");
}
#[test]
fn test_query_builder_contributions_present() {
let db = YantrikDB::new(":memory:", 8).unwrap();
db.record(
"test mem",
"episodic",
0.8,
0.5,
604800.0,
&empty_meta(),
&vec_seed(1.0, 8),
"default",
0.8,
"general",
"user",
None,
)
.unwrap();
let results = db
.query(RecallQuery::new(vec_seed(1.0, 8)).top_k(1).skip_reinforce())
.unwrap();
assert_eq!(results.len(), 1);
let r = &results[0];
assert!(r.scores.valence_multiplier >= 1.0);
assert!(r.scores.contributions.similarity >= 0.0);
assert!(r.scores.contributions.decay >= 0.0);
assert!(r.scores.contributions.recency >= 0.0);
assert!(r.scores.contributions.importance >= 0.0);
}
fn test_key() -> [u8; 32] {
let mut key = [0u8; 32];
for (i, b) in key.iter_mut().enumerate() {
*b = (i as u8).wrapping_mul(7).wrapping_add(42);
}
key
}
#[test]
fn test_encrypted_record_and_get() {
let key = test_key();
let db = YantrikDB::new_encrypted(":memory:", 8, &key).unwrap();
assert!(db.is_encrypted());
let meta = serde_json::json!({"source": "test", "topic": "encryption"});
let emb = vec_seed(1.0, 8);
let rid = db
.record(
"secret memory",
"episodic",
0.8,
0.3,
604800.0,
&meta,
&emb,
"default",
0.8,
"general",
"user",
None,
)
.unwrap();
let mem = db.get(&rid).unwrap().unwrap();
assert_eq!(mem.text, "secret memory");
assert_eq!(mem.memory_type, "episodic");
assert_eq!(mem.importance, 0.8);
assert_eq!(mem.metadata["source"], "test");
assert_eq!(mem.metadata["topic"], "encryption");
}
#[test]
fn test_encrypted_data_not_plaintext_in_db() {
let key = test_key();
let db = YantrikDB::new_encrypted(":memory:", 8, &key).unwrap();
let rid = db
.record(
"secret memory",
"episodic",
0.5,
0.0,
604800.0,
&empty_meta(),
&vec_seed(1.0, 8),
"default",
0.8,
"general",
"user",
None,
)
.unwrap();
let stored_text: String = db
.conn()
.query_row(
"SELECT text FROM memories WHERE rid = ?1",
params![rid],
|r| r.get(0),
)
.unwrap();
assert_ne!(
stored_text, "secret memory",
"text should be encrypted in DB"
);
let stored_meta: String = db
.conn()
.query_row(
"SELECT metadata FROM memories WHERE rid = ?1",
params![rid],
|r| r.get(0),
)
.unwrap();
assert_ne!(stored_meta, "{}", "metadata should be encrypted in DB");
}
#[test]
fn test_encrypted_recall_roundtrip() {
let key = test_key();
let db = YantrikDB::new_encrypted(":memory:", 8, &key).unwrap();
db.record(
"cat sat on mat",
"episodic",
0.5,
0.0,
604800.0,
&empty_meta(),
&vec_seed(1.0, 8),
"default",
0.8,
"general",
"user",
None,
)
.unwrap();
db.record(
"dog ran in park",
"episodic",
0.5,
0.0,
604800.0,
&empty_meta(),
&vec_seed(5.0, 8),
"default",
0.8,
"general",
"user",
None,
)
.unwrap();
db.record(
"cats love warmth",
"episodic",
0.5,
0.0,
604800.0,
&empty_meta(),
&vec_seed(1.1, 8),
"default",
0.8,
"general",
"user",
None,
)
.unwrap();
let results = db
.recall(
&vec_seed(1.0, 8),
2,
None,
None,
false,
false,
None,
true,
None,
None,
None,
None,
None,
)
.unwrap();
assert_eq!(results.len(), 2);
assert!(results.iter().any(|r| r.text.contains("cat")));
}
#[test]
fn test_encrypted_record_batch() {
let key = test_key();
let db = YantrikDB::new_encrypted(":memory:", 8, &key).unwrap();
let inputs: Vec<RecordInput> = (0..5)
.map(|i| RecordInput {
text: format!("encrypted batch {i}"),
memory_type: "episodic".to_string(),
importance: 0.5,
valence: 0.0,
half_life: 604800.0,
metadata: serde_json::json!({"idx": i}),
embedding: vec_seed(i as f32, 8),
namespace: "default".to_string(),
certainty: 0.8,
domain: "general".to_string(),
source: "user".to_string(),
emotional_state: None,
})
.collect();
let rids = db.record_batch(&inputs).unwrap();
assert_eq!(rids.len(), 5);
for (i, rid) in rids.iter().enumerate() {
let mem = db.get(rid).unwrap().unwrap();
assert_eq!(mem.text, format!("encrypted batch {i}"));
assert_eq!(mem.metadata["idx"], i);
}
}
#[test]
fn test_encrypted_archive_hydrate() {
let key = test_key();
let db = YantrikDB::new_encrypted(":memory:", 8, &key).unwrap();
let emb = vec_seed(2.0, 8);
let rid = db
.record(
"to archive encrypted",
"episodic",
0.5,
0.0,
604800.0,
&empty_meta(),
&emb,
"default",
0.8,
"general",
"user",
None,
)
.unwrap();
assert!(db.archive(&rid).unwrap());
let mem = db.get(&rid).unwrap().unwrap();
assert_eq!(mem.storage_tier, "cold");
assert_eq!(mem.text, "to archive encrypted");
assert!(db.hydrate(&rid).unwrap());
let mem = db.get(&rid).unwrap().unwrap();
assert_eq!(mem.storage_tier, "hot");
let results = db
.recall(
&emb, 10, None, None, false, false, None, true, None, None, None, None, None,
)
.unwrap();
assert!(results.iter().any(|r| r.rid == rid));
}
#[test]
fn test_encrypted_correct_memory() {
let key = test_key();
let db = YantrikDB::new_encrypted(":memory:", 8, &key).unwrap();
let rid = db
.record(
"color is green",
"semantic",
0.7,
0.0,
604800.0,
&empty_meta(),
&vec_seed(1.0, 8),
"default",
0.8,
"general",
"user",
None,
)
.unwrap();
let result = db
.correct(
&rid,
Some("color is blue"),
None,
Some(0.9),
None,
"fixed", )
.unwrap();
assert_eq!(result.corrected_rid, rid);
assert!(!result.original_tombstoned);
let updated = db.get(&rid).unwrap().unwrap();
assert_eq!(updated.text, "color is blue");
assert!((updated.importance - 0.9).abs() < 1e-9);
}
#[test]
fn test_unencrypted_db_unaffected() {
let db = YantrikDB::new(":memory:", 8).unwrap();
assert!(!db.is_encrypted());
let rid = db
.record(
"plaintext memory",
"episodic",
0.5,
0.0,
604800.0,
&empty_meta(),
&vec_seed(1.0, 8),
"default",
0.8,
"general",
"user",
None,
)
.unwrap();
let mem = db.get(&rid).unwrap().unwrap();
assert_eq!(mem.text, "plaintext memory");
let stored_text: String = db
.conn()
.query_row(
"SELECT text FROM memories WHERE rid = ?1",
params![rid],
|r| r.get(0),
)
.unwrap();
assert_eq!(stored_text, "plaintext memory");
}
#[test]
fn test_encrypted_db_wrong_key_fails() {
use tempfile::NamedTempFile;
let tmp = NamedTempFile::new().unwrap();
let path = tmp.path().to_str().unwrap();
let key_a = test_key();
{
let db = YantrikDB::new_encrypted(path, 8, &key_a).unwrap();
db.record(
"secret",
"episodic",
0.5,
0.0,
604800.0,
&empty_meta(),
&vec_seed(1.0, 8),
"default",
0.8,
"general",
"user",
None,
)
.unwrap();
db.close().unwrap();
}
let mut key_b = [0u8; 32];
key_b[0] = 99;
let result = YantrikDB::new_encrypted(path, 8, &key_b);
assert!(
result.is_err(),
"Opening encrypted DB with wrong key should fail"
);
}
#[test]
fn test_encrypted_db_reopen_same_key() {
use tempfile::NamedTempFile;
let tmp = NamedTempFile::new().unwrap();
let path = tmp.path().to_str().unwrap();
let key = test_key();
let rid;
{
let db = YantrikDB::new_encrypted(path, 8, &key).unwrap();
rid = db
.record(
"persistent secret",
"episodic",
0.5,
0.0,
604800.0,
&empty_meta(),
&vec_seed(1.0, 8),
"default",
0.8,
"general",
"user",
None,
)
.unwrap();
db.close().unwrap();
}
{
let db = YantrikDB::new_encrypted(path, 8, &key).unwrap();
let mem = db.get(&rid).unwrap().unwrap();
assert_eq!(mem.text, "persistent secret");
}
}
#[test]
fn test_open_encrypted_db_without_key_fails() {
use tempfile::NamedTempFile;
let tmp = NamedTempFile::new().unwrap();
let path = tmp.path().to_str().unwrap();
let key = test_key();
{
let db = YantrikDB::new_encrypted(path, 8, &key).unwrap();
db.record(
"data",
"episodic",
0.5,
0.0,
604800.0,
&empty_meta(),
&vec_seed(1.0, 8),
"default",
0.8,
"general",
"user",
None,
)
.unwrap();
db.close().unwrap();
}
let result = YantrikDB::new(path, 8);
assert!(
result.is_err(),
"Opening encrypted DB without key should fail"
);
}
#[test]
fn test_record_with_dimensions() {
let db = YantrikDB::new(":memory:", 8).unwrap();
let emb = vec_seed(1.0, 8);
let rid = db
.record(
"meeting notes for Q1 planning",
"episodic",
0.7,
0.2,
604800.0,
&empty_meta(),
&emb,
"default",
0.9,
"work",
"document",
Some("joy"),
)
.unwrap();
let mem = db.get(&rid).unwrap().unwrap();
assert_eq!(mem.text, "meeting notes for Q1 planning");
assert!(
(mem.certainty - 0.9).abs() < 1e-6,
"certainty should be 0.9, got {}",
mem.certainty
);
assert_eq!(mem.domain, "work");
assert_eq!(mem.source, "document");
assert_eq!(mem.emotional_state, Some("joy".to_string()));
}
#[test]
fn test_domain_filter() {
let db = YantrikDB::new(":memory:", 8).unwrap();
db.record(
"work task A",
"episodic",
0.5,
0.0,
604800.0,
&empty_meta(),
&vec_seed(1.0, 8),
"default",
0.8,
"work",
"user",
None,
)
.unwrap();
db.record(
"health checkup",
"episodic",
0.5,
0.0,
604800.0,
&empty_meta(),
&vec_seed(2.0, 8),
"default",
0.8,
"health",
"user",
None,
)
.unwrap();
db.record(
"work task B",
"episodic",
0.5,
0.0,
604800.0,
&empty_meta(),
&vec_seed(3.0, 8),
"default",
0.8,
"work",
"user",
None,
)
.unwrap();
let results = db
.recall(
&vec_seed(1.0, 8),
10,
None,
None,
false,
false,
None,
true,
None,
Some("work"),
None,
None,
None,
)
.unwrap();
assert_eq!(
results.len(),
2,
"Expected 2 work-domain memories, got {}",
results.len()
);
for r in &results {
assert_eq!(r.domain, "work");
}
}
#[test]
fn test_source_filter() {
let db = YantrikDB::new(":memory:", 8).unwrap();
db.record(
"user input A",
"episodic",
0.5,
0.0,
604800.0,
&empty_meta(),
&vec_seed(1.0, 8),
"default",
0.8,
"general",
"user",
None,
)
.unwrap();
db.record(
"system log",
"episodic",
0.5,
0.0,
604800.0,
&empty_meta(),
&vec_seed(2.0, 8),
"default",
0.8,
"general",
"system",
None,
)
.unwrap();
db.record(
"user input B",
"episodic",
0.5,
0.0,
604800.0,
&empty_meta(),
&vec_seed(3.0, 8),
"default",
0.8,
"general",
"user",
None,
)
.unwrap();
let results = db
.recall(
&vec_seed(1.0, 8),
10,
None,
None,
false,
false,
None,
true,
None,
None,
Some("user"),
None,
None,
)
.unwrap();
assert_eq!(
results.len(),
2,
"Expected 2 user-source memories, got {}",
results.len()
);
for r in &results {
assert_eq!(r.source, "user");
}
}
#[test]
fn test_domain_and_source_combined_filter() {
let db = YantrikDB::new(":memory:", 8).unwrap();
db.record(
"work from user",
"episodic",
0.5,
0.0,
604800.0,
&empty_meta(),
&vec_seed(1.0, 8),
"default",
0.8,
"work",
"user",
None,
)
.unwrap();
db.record(
"work from system",
"episodic",
0.5,
0.0,
604800.0,
&empty_meta(),
&vec_seed(2.0, 8),
"default",
0.8,
"work",
"system",
None,
)
.unwrap();
db.record(
"health from user",
"episodic",
0.5,
0.0,
604800.0,
&empty_meta(),
&vec_seed(3.0, 8),
"default",
0.8,
"health",
"user",
None,
)
.unwrap();
db.record(
"health from system",
"episodic",
0.5,
0.0,
604800.0,
&empty_meta(),
&vec_seed(4.0, 8),
"default",
0.8,
"health",
"system",
None,
)
.unwrap();
let results = db
.recall(
&vec_seed(1.0, 8),
10,
None,
None,
false,
false,
None,
true,
None,
Some("work"),
Some("user"),
None,
None,
)
.unwrap();
assert_eq!(
results.len(),
1,
"Expected 1 work+user memory, got {}",
results.len()
);
assert_eq!(results[0].domain, "work");
assert_eq!(results[0].source, "user");
let results = db
.recall(
&vec_seed(4.0, 8),
10,
None,
None,
false,
false,
None,
true,
None,
Some("health"),
Some("system"),
None,
None,
)
.unwrap();
assert_eq!(
results.len(),
1,
"Expected 1 health+system memory, got {}",
results.len()
);
assert_eq!(results[0].domain, "health");
assert_eq!(results[0].source, "system");
}
#[test]
fn test_dimensions_preserved_on_correct() {
let db = YantrikDB::new(":memory:", 8).unwrap();
let rid = db
.record(
"the sky is green",
"semantic",
0.7,
0.0,
604800.0,
&empty_meta(),
&vec_seed(1.0, 8),
"default",
0.6,
"work",
"document",
Some("surprise"),
)
.unwrap();
let result = db
.correct(
&rid,
Some("the sky is blue"),
None,
Some(0.9),
None,
"color fix", )
.unwrap();
assert!(!result.original_tombstoned);
assert_eq!(result.corrected_rid, rid);
let updated = db.get(&rid).unwrap().unwrap();
assert_eq!(updated.text, "the sky is blue");
assert_eq!(
updated.domain, "work",
"domain should be preserved after correction"
);
assert_eq!(
updated.source, "document",
"source should be preserved after correction"
);
}
#[test]
fn test_batch_record_with_dimensions() {
let db = YantrikDB::new(":memory:", 8).unwrap();
let inputs: Vec<RecordInput> = vec![
RecordInput {
text: "batch work meeting".to_string(),
memory_type: "episodic".to_string(),
importance: 0.6,
valence: 0.1,
half_life: 604800.0,
metadata: serde_json::json!({"batch": true}),
embedding: vec_seed(1.0, 8),
namespace: "default".to_string(),
certainty: 0.95,
domain: "work".to_string(),
source: "document".to_string(),
emotional_state: Some("focus".to_string()),
},
RecordInput {
text: "batch health jog".to_string(),
memory_type: "episodic".to_string(),
importance: 0.4,
valence: 0.3,
half_life: 604800.0,
metadata: serde_json::json!({"batch": true}),
embedding: vec_seed(2.0, 8),
namespace: "default".to_string(),
certainty: 0.7,
domain: "health".to_string(),
source: "user".to_string(),
emotional_state: None,
},
RecordInput {
text: "batch personal diary".to_string(),
memory_type: "semantic".to_string(),
importance: 0.3,
valence: -0.1,
half_life: 604800.0,
metadata: serde_json::json!({"batch": true}),
embedding: vec_seed(3.0, 8),
namespace: "default".to_string(),
certainty: 0.5,
domain: "personal".to_string(),
source: "system".to_string(),
emotional_state: Some("calm".to_string()),
},
];
let rids = db.record_batch(&inputs).unwrap();
assert_eq!(rids.len(), 3);
let m0 = db.get(&rids[0]).unwrap().unwrap();
assert_eq!(m0.text, "batch work meeting");
assert!((m0.certainty - 0.95).abs() < 1e-6);
assert_eq!(m0.domain, "work");
assert_eq!(m0.source, "document");
assert_eq!(m0.emotional_state, Some("focus".to_string()));
let m1 = db.get(&rids[1]).unwrap().unwrap();
assert_eq!(m1.domain, "health");
assert_eq!(m1.source, "user");
assert_eq!(m1.emotional_state, None);
let m2 = db.get(&rids[2]).unwrap().unwrap();
assert_eq!(m2.domain, "personal");
assert_eq!(m2.source, "system");
assert_eq!(m2.emotional_state, Some("calm".to_string()));
}
#[test]
fn test_recall_with_response_structure() {
let db = YantrikDB::new(":memory:", 8).unwrap();
db.record(
"memory for recall response",
"episodic",
0.7,
0.0,
604800.0,
&empty_meta(),
&vec_seed(1.0, 8),
"default",
0.8,
"general",
"user",
None,
)
.unwrap();
let response = db
.recall_with_response(
&vec_seed(1.0, 8),
5,
None,
None,
false,
false,
None,
true,
None,
None,
None,
)
.unwrap();
assert!(!response.results.is_empty(), "results should not be empty");
assert!(
response.confidence >= 0.0 && response.confidence <= 1.0,
"confidence should be in [0, 1], got {}",
response.confidence
);
assert!(
!response.retrieval_summary.sources_used.is_empty(),
"sources_used should not be empty"
);
assert!(
response.retrieval_summary.candidate_count > 0,
"candidate_count should be > 0"
);
let _ = response.hints; }
#[test]
fn test_high_confidence_no_hints() {
let db = YantrikDB::new(":memory:", 8).unwrap();
let emb = vec_seed(5.0, 8);
for i in 0..5 {
db.record(
&format!("exact match memory {}", i),
"episodic",
0.9,
0.0,
604800.0,
&empty_meta(),
&emb,
"default",
0.8,
"general",
"user",
None,
)
.unwrap();
}
let response = db
.recall_with_response(
&emb, 5, None, None, false, false, None, true, None, None, None,
)
.unwrap();
assert!(!response.results.is_empty());
assert!(
response.confidence >= 0.60,
"Confidence should be >= 0.60 for exact match with full density, got {}",
response.confidence,
);
assert!(
response.hints.is_empty(),
"Hints should be empty for high-confidence recall, got {} hints",
response.hints.len(),
);
}
#[test]
fn test_low_confidence_has_hints() {
let db = YantrikDB::new(":memory:", 8).unwrap();
db.record(
"something about cats",
"episodic",
0.5,
0.0,
604800.0,
&empty_meta(),
&vec_seed(1.0, 8),
"default",
0.8,
"general",
"user",
None,
)
.unwrap();
let far_emb = vec_seed(100.0, 8);
let response = db
.recall_with_response(
&far_emb,
10,
None,
None,
false,
false,
Some("cats"), true,
None,
None,
None,
)
.unwrap();
assert!(
response.confidence < 0.60,
"Confidence should be < 0.60 for distant query, got {}",
response.confidence,
);
assert!(
!response.hints.is_empty(),
"Hints should be non-empty for low-confidence recall with short query_text",
);
}
#[test]
fn test_recall_refine_excludes_originals() {
let db = YantrikDB::new(":memory:", 8).unwrap();
let mut rids = Vec::new();
for i in 1..=5 {
let rid = db
.record(
&format!("memory number {}", i),
"episodic",
0.5,
0.0,
604800.0,
&empty_meta(),
&vec_seed(i as f32, 8),
"default",
0.8,
"general",
"user",
None,
)
.unwrap();
rids.push(rid);
}
let first_results = db
.recall(
&vec_seed(1.0, 8),
2,
None,
None,
false,
false,
None,
true,
None,
None,
None,
None,
None,
)
.unwrap();
assert_eq!(first_results.len(), 2);
let original_rids: Vec<String> = first_results.iter().map(|r| r.rid.clone()).collect();
let refined = db
.recall_refine(
&vec_seed(1.0, 8), &vec_seed(2.0, 8), &original_rids
.iter()
.map(|s| s.as_str())
.collect::<Vec<&str>>()
.iter()
.map(|s| s.to_string())
.collect::<Vec<String>>(),
3, None,
None,
None,
)
.unwrap();
for result in &refined.results {
assert!(
!original_rids.contains(&result.rid),
"Refined result should not contain original RID {}, but it does",
result.rid,
);
}
}
#[test]
fn test_recall_refine_returns_response() {
let db = YantrikDB::new(":memory:", 8).unwrap();
for i in 1..=4 {
db.record(
&format!("refine test {}", i),
"episodic",
0.5,
0.0,
604800.0,
&empty_meta(),
&vec_seed(i as f32, 8),
"default",
0.8,
"general",
"user",
None,
)
.unwrap();
}
let exclude: Vec<String> = vec![];
let response = db
.recall_refine(
&vec_seed(1.0, 8),
&vec_seed(2.0, 8),
&exclude,
3,
None,
None,
None,
)
.unwrap();
assert!(response.confidence >= 0.0 && response.confidence <= 1.0);
assert!(!response.retrieval_summary.sources_used.is_empty());
assert!(response.retrieval_summary.candidate_count > 0);
let _ = &response.hints;
}
#[test]
fn test_retrieval_summary_fields() {
let db = YantrikDB::new(":memory:", 8).unwrap();
for i in 1..=3 {
db.record(
&format!("summary test {}", i),
"episodic",
0.5,
0.0,
604800.0,
&empty_meta(),
&vec_seed(i as f32, 8),
"default",
0.8,
"general",
"user",
None,
)
.unwrap();
}
let response = db
.recall_with_response(
&vec_seed(1.0, 8),
5,
None,
None,
false,
false,
None,
true,
None,
None,
None,
)
.unwrap();
let summary = &response.retrieval_summary;
assert!(
summary.top_similarity > 0.0,
"top_similarity should be > 0, got {}",
summary.top_similarity
);
assert!(
summary.sources_used.contains(&"hnsw".to_string()),
"sources_used should contain 'hnsw', got {:?}",
summary.sources_used,
);
assert!(
summary.candidate_count > 0,
"candidate_count should be > 0, got {}",
summary.candidate_count
);
}
#[test]
fn test_recall_feedback_stores() {
let db = YantrikDB::new(":memory:", 8).unwrap();
let emb = vec_seed(1.0, 8);
let rid = db
.record(
"feedback target",
"episodic",
0.5,
0.0,
604800.0,
&empty_meta(),
&emb,
"default",
0.8,
"general",
"user",
None,
)
.unwrap();
db.recall_feedback(
Some("test query"),
Some(&emb),
&rid,
"relevant",
Some(0.85),
Some(1),
)
.unwrap();
let count: i64 = db
.conn()
.query_row(
"SELECT COUNT(*) FROM recall_feedback WHERE rid = ?1 AND feedback = 'relevant'",
params![rid],
|row| row.get(0),
)
.unwrap();
assert_eq!(count, 1, "Expected 1 feedback row, got {}", count);
}
#[test]
fn test_learned_weights_default() {
let db = YantrikDB::new(":memory:", 8).unwrap();
let weights = db.load_learned_weights().unwrap();
assert!(
(weights.w_sim - 0.50).abs() < 1e-6,
"w_sim default should be 0.50, got {}",
weights.w_sim
);
assert!(
(weights.w_decay - 0.20).abs() < 1e-6,
"w_decay default should be 0.20, got {}",
weights.w_decay
);
assert!(
(weights.w_recency - 0.30).abs() < 1e-6,
"w_recency default should be 0.30, got {}",
weights.w_recency
);
assert!(
(weights.gate_tau - 0.25).abs() < 1e-6,
"gate_tau default should be 0.25, got {}",
weights.gate_tau
);
assert!(
(weights.alpha_imp - 0.80).abs() < 1e-6,
"alpha_imp default should be 0.80, got {}",
weights.alpha_imp
);
assert_eq!(weights.generation, 0, "generation should start at 0");
}
#[test]
fn test_feedback_count_increments() {
let db = YantrikDB::new(":memory:", 8).unwrap();
let emb = vec_seed(1.0, 8);
let rid = db
.record(
"counting feedback",
"episodic",
0.5,
0.0,
604800.0,
&empty_meta(),
&emb,
"default",
0.8,
"general",
"user",
None,
)
.unwrap();
for i in 0..5 {
let feedback_type = if i % 2 == 0 { "relevant" } else { "irrelevant" };
db.recall_feedback(
Some("query"),
Some(&emb),
&rid,
feedback_type,
Some(0.5),
Some(i + 1),
)
.unwrap();
}
let count = db.feedback_count().unwrap();
assert_eq!(count, 5, "Expected feedback_count=5, got {}", count);
}
#[test]
fn test_learning_skipped_under_threshold() {
let db = YantrikDB::new(":memory:", 8).unwrap();
let emb = vec_seed(1.0, 8);
let rid = db
.record(
"learning test",
"episodic",
0.5,
0.0,
604800.0,
&empty_meta(),
&emb,
"default",
0.8,
"general",
"user",
None,
)
.unwrap();
for i in 0..10 {
db.recall_feedback(Some("q"), Some(&emb), &rid, "relevant", Some(0.5), Some(i))
.unwrap();
}
let result = db.run_learning().unwrap();
assert!(
!result,
"run_learning should return false with < 20 feedback items"
);
}
#[test]
fn test_learning_runs_with_enough_feedback() {
let db = YantrikDB::new(":memory:", 8).unwrap();
let emb = vec_seed(1.0, 8);
let rid = db
.record(
"learning convergence",
"episodic",
0.7,
0.0,
604800.0,
&empty_meta(),
&emb,
"default",
0.8,
"general",
"user",
None,
)
.unwrap();
for i in 0..25 {
let feedback_type = if i % 3 == 0 { "irrelevant" } else { "relevant" };
let score = 0.3 + (i as f64 * 0.02);
db.recall_feedback(
Some("learning query"),
Some(&emb),
&rid,
feedback_type,
Some(score),
Some(i + 1),
)
.unwrap();
}
let result = db.run_learning();
assert!(
result.is_ok(),
"run_learning should not error with 25 feedback items: {:?}",
result.err()
);
}
#[test]
fn test_think_includes_learning() {
let db = YantrikDB::new(":memory:", 8).unwrap();
let emb = vec_seed(1.0, 8);
let rid = db
.record(
"think learning integration",
"episodic",
0.7,
0.0,
604800.0,
&empty_meta(),
&emb,
"default",
0.8,
"general",
"user",
None,
)
.unwrap();
for i in 0..26 {
let feedback_type = if i % 4 == 0 { "irrelevant" } else { "relevant" };
db.recall_feedback(
Some("think query"),
Some(&emb),
&rid,
feedback_type,
Some(0.5),
Some(i + 1),
)
.unwrap();
}
let config = ThinkConfig::default();
let result = db.think(&config);
assert!(
result.is_ok(),
"think() should not error when learning has enough feedback: {:?}",
result.err()
);
}
#[test]
fn test_conflict_entity_substitution_org() {
let db = YantrikDB::new(":memory:", 8).unwrap();
let emb1 = vec_seed(1.0, 8);
let emb2 = vec_seed(1.1, 8);
db.relate("User", "Google", "works_at", 1.0).unwrap();
db.relate("User", "Meta", "works_at", 1.0).unwrap();
db.record(
"User works at Google as a senior engineer",
"episodic",
0.7,
0.0,
604800.0,
&empty_meta(),
&emb1,
"default",
0.8,
"work",
"user",
None,
)
.unwrap();
db.record(
"User works at Meta as a senior engineer",
"episodic",
0.7,
0.0,
604800.0,
&empty_meta(),
&emb2,
"default",
0.8,
"work",
"user",
None,
)
.unwrap();
let conflicts = crate::conflict::scan_conflicts(&db).unwrap();
assert!(!conflicts.is_empty(), "should detect works_at conflict");
assert_eq!(conflicts[0].conflict_type, "identity_fact");
}
#[test]
fn test_conflict_entity_substitution_tech() {
let db = YantrikDB::new(":memory:", 384).unwrap();
db.relate("API", "PostgreSQL", "uses", 1.0).unwrap();
db.relate("API", "MySQL", "uses", 1.0).unwrap();
let emb1 = vec_seed(2.0, 384);
let emb2 = vec_seed(2.05, 384);
db.record(
"The API service uses PostgreSQL for the database layer",
"semantic",
0.8,
0.0,
604800.0,
&empty_meta(),
&emb1,
"default",
0.8,
"architecture",
"user",
None,
)
.unwrap();
db.record(
"The API service uses MySQL for the database layer",
"semantic",
0.8,
0.0,
604800.0,
&empty_meta(),
&emb2,
"default",
0.8,
"architecture",
"user",
None,
)
.unwrap();
let conflicts = crate::conflict::scan_conflicts(&db).unwrap();
let entity_based = conflicts
.iter()
.filter(|c| c.detection_reason.contains("contradict"))
.collect::<Vec<_>>();
assert!(conflicts.len() >= 0);
}
#[test]
fn test_relate_infers_entity_types() {
let db = YantrikDB::new(":memory:", 8).unwrap();
db.relate("MyApp", "React", "built_with", 1.0).unwrap();
let entities = db.search_entities(Some("MyApp"), None, 1).unwrap();
assert_eq!(entities.len(), 1);
assert_eq!(entities[0].entity_type, "project");
let entities = db.search_entities(Some("React"), None, 1).unwrap();
assert_eq!(entities.len(), 1);
assert_eq!(entities[0].entity_type, "tech");
}
#[test]
fn test_relate_infers_infrastructure() {
let db = YantrikDB::new(":memory:", 8).unwrap();
db.relate("Backend", "AWS", "deployed_on", 1.0).unwrap();
let entities = db.search_entities(Some("AWS"), None, 1).unwrap();
assert_eq!(entities[0].entity_type, "infrastructure");
}
#[test]
fn test_recall_with_response_has_certainty_reasons() {
let db = YantrikDB::new(":memory:", 8).unwrap();
let emb = vec_seed(1.0, 8);
db.record(
"Important architecture decision about microservices",
"semantic",
0.8,
0.0,
604800.0,
&empty_meta(),
&emb,
"default",
0.8,
"work",
"user",
None,
)
.unwrap();
let response = db
.recall_with_response(
&emb,
5,
None,
None,
false,
false,
Some("architecture decision"),
false,
None,
None,
None,
)
.unwrap();
assert!(
!response.certainty_reasons.is_empty(),
"should have certainty reasons"
);
assert!(response.confidence >= 0.0 && response.confidence <= 1.0);
}
#[test]
fn test_recall_empty_db_low_confidence() {
let db = YantrikDB::new(":memory:", 8).unwrap();
let emb = vec_seed(1.0, 8);
let response = db
.recall_with_response(
&emb,
5,
None,
None,
false,
false,
Some("anything"),
false,
None,
None,
None,
)
.unwrap();
assert!(
response.confidence < 0.5,
"empty DB should have low confidence"
);
assert!(
response
.certainty_reasons
.iter()
.any(|r| r.contains("No") || r.contains("Sparse") || r.contains("Weak")),
"should explain low confidence"
);
}
#[test]
fn test_relationship_depth_basic() {
let db = YantrikDB::new(":memory:", 8).unwrap();
let emb = vec_seed(1.0, 8);
db.relate("Alice", "Bob", "knows", 1.0).unwrap();
db.relate("Alice", "ProjectX", "works_on", 1.0).unwrap();
db.record(
"Alice presented the quarterly report",
"episodic",
0.5,
0.3,
604800.0,
&empty_meta(),
&emb,
"default",
0.8,
"work",
"user",
None,
)
.unwrap();
db.record(
"Alice prefers async communication",
"semantic",
0.6,
0.0,
604800.0,
&empty_meta(),
&vec_seed(2.0, 8),
"default",
0.8,
"preference",
"user",
None,
)
.unwrap();
db.apply_pending_ops_once(100).unwrap();
let depth = db.relationship_depth("Alice", None).unwrap();
assert_eq!(depth.entity, "Alice");
assert_eq!(depth.entity_type, "person");
assert!(
depth.connection_count >= 2,
"Alice connected to Bob and ProjectX"
);
assert!(
depth.memories_mentioning >= 2,
"at least 2 memories mention Alice"
);
assert!(depth.depth_score > 0.0, "should have positive depth score");
assert!(depth.depth_score <= 1.0);
}
#[test]
fn test_relationship_depth_not_found() {
let db = YantrikDB::new(":memory:", 8).unwrap();
let result = db.relationship_depth("NonexistentEntity", None);
assert!(result.is_err(), "should error for unknown entity");
}
#[test]
fn test_record_and_surface_procedural() {
let db = YantrikDB::new(":memory:", 8).unwrap();
let emb = vec_seed(3.0, 8);
let rid = db
.record_procedural(
"Use Agent tool with Explore subtype for architectural questions in this codebase",
&emb,
"work",
"code search",
0.8,
"default",
)
.unwrap();
let mem = db.get(&rid).unwrap().unwrap();
assert_eq!(mem.memory_type, "procedural");
assert!((mem.importance - 0.8).abs() < 0.01);
let results = db
.surface_procedural(&emb, Some("how to search code"), Some("work"), 5, None)
.unwrap();
assert!(!results.is_empty(), "should surface the procedural memory");
assert_eq!(results[0].memory_type, "procedural");
}
#[test]
fn test_reinforce_procedural() {
let db = YantrikDB::new(":memory:", 8).unwrap();
let emb = vec_seed(4.0, 8);
let rid = db
.record_procedural(
"Always run tests before pushing",
&emb,
"work",
"git workflow",
0.5,
"default",
)
.unwrap();
let reinforced = db.reinforce_procedural(&rid, 1.0).unwrap();
assert!(reinforced);
let mem = db.get(&rid).unwrap().unwrap();
assert!(
mem.importance > 0.5,
"importance should increase after positive reinforcement"
);
}
#[test]
fn test_procedural_stats() {
let db = YantrikDB::new(":memory:", 8).unwrap();
db.record_procedural(
"proc 1",
&vec_seed(1.0, 8),
"work",
"task A",
0.7,
"default",
)
.unwrap();
db.record_procedural(
"proc 2",
&vec_seed(2.0, 8),
"work",
"task B",
0.9,
"default",
)
.unwrap();
db.record_procedural(
"proc 3",
&vec_seed(3.0, 8),
"health",
"exercise",
0.5,
"default",
)
.unwrap();
let stats = db.procedural_stats(None).unwrap();
assert!(
stats.len() >= 2,
"should have stats for work and health domains"
);
let work_stats = stats.iter().find(|(d, _, _)| d == "work");
assert!(work_stats.is_some());
let (_, count, _) = work_stats.unwrap();
assert_eq!(*count, 2);
}
#[test]
fn test_session_awareness_trigger() {
let db = YantrikDB::new(":memory:", 8).unwrap();
let emb = vec_seed(1.0, 8);
let sid = db
.session_start("default", "claude", &serde_json::json!({}))
.unwrap();
db.record(
"Worked on battle testing the MCP server",
"episodic",
0.7,
0.5,
604800.0,
&empty_meta(),
&emb,
"default",
0.8,
"work",
"user",
None,
)
.unwrap();
let _summary = db
.session_end(&sid, Some("Battle tested MCP server v0.2.8"))
.unwrap();
db.conn().execute(
"UPDATE sessions SET ended_at = ended_at - 86400 * 3, started_at = started_at - 86400 * 3 WHERE session_id = ?1",
params![sid],
).unwrap();
let config = ThinkConfig {
run_consolidation: false,
run_conflict_scan: false,
run_pattern_mining: false,
run_personality: false,
..Default::default()
};
let result = db.think(&config).unwrap();
let session_triggers: Vec<_> = result
.triggers
.iter()
.filter(|t| t.trigger_type == "session_awareness")
.collect();
assert!(
!session_triggers.is_empty(),
"should generate session awareness trigger after 3-day gap"
);
assert!(
session_triggers[0].reason.contains("hours"),
"reason should mention time gap"
);
}
use crate::engine::moves::{
adversarial_status, observability, posthoc_outcome, ClaimRef, RecordMoveEventInput,
SideEffectRef,
};
fn mk_move_input(move_type: &str, inputs: &[&str], outputs: &[&str]) -> RecordMoveEventInput {
RecordMoveEventInput {
move_type: move_type.to_string(),
operator_version: "v1".to_string(),
context_regime: None,
observability: observability::OBSERVED.to_string(),
inputs: inputs
.iter()
.enumerate()
.map(|(i, c)| ClaimRef {
claim_id: c.to_string(),
role: "input".to_string(),
ordinal: i as i64,
})
.collect(),
outputs: outputs
.iter()
.enumerate()
.map(|(i, c)| ClaimRef {
claim_id: c.to_string(),
role: "output".to_string(),
ordinal: i as i64,
})
.collect(),
..Default::default()
}
}
#[test]
fn test_m5b_schema_v23_tables_present() {
let db = YantrikDB::new(":memory:", 8).unwrap();
let conn = db.conn();
for table in [
"move_events",
"move_input_edge",
"move_output_edge",
"move_side_effect_edge",
"move_correction_event",
"move_adversarial_instance",
"move_type_registry",
"inference_basis_registry",
"move_composition_rule",
"move_type_profile",
] {
let exists: bool = conn
.query_row(
"SELECT 1 FROM sqlite_master WHERE type='table' AND name=?1",
rusqlite::params![table],
|_| Ok(true),
)
.unwrap_or(false);
assert!(exists, "V23 table {} must exist", table);
}
}
#[test]
fn test_m5b_registries_seeded_at_bootstrap() {
let db = YantrikDB::new(":memory:", 8).unwrap();
let conn = db.conn();
let mt_count: i64 = conn
.query_row(
"SELECT COUNT(*) FROM move_type_registry WHERE status = 'active'",
[],
|r| r.get(0),
)
.unwrap();
assert_eq!(mt_count, 13, "move_type_registry seed vocabulary count");
let ib_count: i64 = conn
.query_row(
"SELECT COUNT(*) FROM inference_basis_registry WHERE status = 'active'",
[],
|r| r.get(0),
)
.unwrap();
assert_eq!(
ib_count, 5,
"inference_basis_registry seed vocabulary count"
);
for mt in [
"analogy",
"decomposition",
"source_audit",
"hypothesis_generation",
] {
let found: bool = conn
.query_row(
"SELECT 1 FROM move_type_registry WHERE move_type = ?1",
rusqlite::params![mt],
|_| Ok(true),
)
.unwrap_or(false);
assert!(found, "seed vocabulary missing: {}", mt);
}
}
#[test]
fn test_m5b_record_move_event_basic() {
let db = YantrikDB::new(":memory:", 8).unwrap();
let move_id = db
.record_move_event(mk_move_input(
"analogy",
&["claim_a", "claim_b"],
&["claim_c"],
))
.unwrap();
assert!(!move_id.is_empty());
let ev = db.get_move_event(&move_id).unwrap().unwrap();
assert_eq!(ev.move_type, "analogy");
assert_eq!(ev.operator_version, "v1");
assert_eq!(ev.observability, "observed");
assert_eq!(ev.context_regime, "default");
assert!(ev.posthoc_outcome.is_none());
assert_eq!(ev.expected_evaluation_horizon_ms, Some(60_000));
let inputs = db.get_move_inputs(&move_id).unwrap();
assert_eq!(inputs.len(), 2);
assert_eq!(inputs[0].claim_id, "claim_a");
assert_eq!(inputs[1].claim_id, "claim_b");
let outputs = db.get_move_outputs(&move_id).unwrap();
assert_eq!(outputs.len(), 1);
assert_eq!(outputs[0].claim_id, "claim_c");
}
#[test]
fn test_m5b_record_move_event_with_side_effects() {
let db = YantrikDB::new(":memory:", 8).unwrap();
let mut inp = mk_move_input("quarantine", &["claim_suspect"], &[]);
inp.side_effects = vec![
SideEffectRef {
claim_id: "downstream_a".into(),
effect_kind: "quarantine".into(),
},
SideEffectRef {
claim_id: "downstream_b".into(),
effect_kind: "quarantine".into(),
},
];
let move_id = db.record_move_event(inp).unwrap();
let side_effects = db.get_move_side_effects(&move_id).unwrap();
assert_eq!(side_effects.len(), 2);
}
#[test]
fn test_m5b_list_moves_consuming_and_producing() {
let db = YantrikDB::new(":memory:", 8).unwrap();
let m1 = db
.record_move_event(mk_move_input("analogy", &["claim_x"], &["claim_y"]))
.unwrap();
let m2 = db
.record_move_event(mk_move_input("decomposition", &["claim_x"], &["claim_z"]))
.unwrap();
let consumers = db.list_moves_consuming_claim("claim_x", 10).unwrap();
assert_eq!(consumers.len(), 2);
let consumer_ids: std::collections::HashSet<_> =
consumers.iter().map(|m| m.move_id.clone()).collect();
assert!(consumer_ids.contains(&m1));
assert!(consumer_ids.contains(&m2));
let producers = db.list_moves_producing_claim("claim_y", 10).unwrap();
assert_eq!(producers.len(), 1);
assert_eq!(producers[0].move_id, m1);
}
#[test]
fn test_m5b_record_move_rejects_invalid_observability_fields() {
let db = YantrikDB::new(":memory:", 8).unwrap();
let mut inp = mk_move_input("analogy", &["a"], &["b"]);
inp.inference_confidence = Some(0.7);
let r = db.record_move_event(inp);
assert!(
r.is_err(),
"should reject inference_confidence on non-inferred"
);
let mut inp2 = mk_move_input("analogy", &["a"], &["b"]);
inp2.inference_basis = Some(vec!["structural_pattern_match".into()]);
let r2 = db.record_move_event(inp2);
assert!(r2.is_err(), "should reject inference_basis on non-inferred");
}
#[test]
fn test_m5b_inferred_move_carries_confidence() {
let db = YantrikDB::new(":memory:", 8).unwrap();
let mut inp = mk_move_input("analogy", &["a"], &["b"]);
inp.observability = observability::INFERRED.to_string();
inp.inference_confidence = Some(0.8);
inp.inference_basis = Some(vec!["structural_pattern_match".into()]);
let move_id = db.record_move_event(inp).unwrap();
let ev = db.get_move_event(&move_id).unwrap().unwrap();
assert_eq!(ev.observability, "inferred");
assert_eq!(ev.inference_confidence, Some(0.8));
assert!(ev
.inference_basis_json
.as_deref()
.unwrap_or("")
.contains("structural_pattern_match"));
}
#[test]
fn test_m5b_record_move_outcome_narrow_mutation() {
let db = YantrikDB::new(":memory:", 8).unwrap();
let move_id = db
.record_move_event(mk_move_input("analogy", &["a"], &["b"]))
.unwrap();
db.record_move_outcome(
&move_id,
posthoc_outcome::CORROBORATED,
Some(serde_json::json!({"predictive_gain": 0.3}).to_string()),
)
.unwrap();
let ev = db.get_move_event(&move_id).unwrap().unwrap();
assert_eq!(ev.posthoc_outcome.as_deref(), Some("corroborated"));
assert!(ev.posthoc_recorded_at.is_some());
assert!(ev.yield_json.contains("predictive_gain"));
let second = db.record_move_outcome(&move_id, posthoc_outcome::RETRACTED, None);
assert!(
second.is_err(),
"should reject overwriting an existing posthoc_outcome"
);
}
#[test]
fn test_m5b_record_move_outcome_rejects_invalid_label() {
let db = YantrikDB::new(":memory:", 8).unwrap();
let move_id = db
.record_move_event(mk_move_input("analogy", &["a"], &["b"]))
.unwrap();
let r = db.record_move_outcome(&move_id, "totally_made_up", None);
assert!(r.is_err());
}
#[test]
fn test_m5b_correction_never_mutates_original() {
let db = YantrikDB::new(":memory:", 8).unwrap();
let move_id = db
.record_move_event(mk_move_input("analogy", &["a"], &["b"]))
.unwrap();
let correction_id = db
.submit_move_correction(
&move_id,
Some("decomposition".to_string()),
None,
None,
"initial categorization was wrong".to_string(),
"curator_alice".to_string(),
)
.unwrap();
assert!(!correction_id.is_empty());
let original = db.get_move_event(&move_id).unwrap().unwrap();
assert_eq!(
original.move_type, "analogy",
"original row must not be mutated"
);
let canonical = db.get_move_event_canonical(&move_id).unwrap().unwrap();
assert_eq!(
canonical.move_type, "decomposition",
"canonical reflects latest correction"
);
let corrections = db.list_move_corrections(&move_id).unwrap();
assert_eq!(corrections.len(), 1);
assert_eq!(corrections[0].correction_id, correction_id);
}
#[test]
fn test_m5b_correction_latest_wins() {
let db = YantrikDB::new(":memory:", 8).unwrap();
let move_id = db
.record_move_event(mk_move_input("analogy", &["a"], &["b"]))
.unwrap();
db.submit_move_correction(
&move_id,
Some("decomposition".into()),
None,
None,
"first correction".into(),
"curator_a".into(),
)
.unwrap();
std::thread::sleep(std::time::Duration::from_millis(10));
db.submit_move_correction(
&move_id,
Some("ladder_up".into()),
None,
None,
"second correction, first was also wrong".into(),
"curator_b".into(),
)
.unwrap();
let canonical = db.get_move_event_canonical(&move_id).unwrap().unwrap();
assert_eq!(
canonical.move_type, "ladder_up",
"latest correction should win"
);
}
#[test]
fn test_m5b_correction_rejects_empty_reason_and_no_fields() {
let db = YantrikDB::new(":memory:", 8).unwrap();
let move_id = db
.record_move_event(mk_move_input("analogy", &["a"], &["b"]))
.unwrap();
let r = db.submit_move_correction(
&move_id,
None,
None,
None,
"reason".into(),
"curator".into(),
);
assert!(r.is_err());
let r2 = db.submit_move_correction(
&move_id,
Some("decomposition".into()),
None,
None,
"".into(),
"curator".into(),
);
assert!(r2.is_err());
}
#[test]
fn test_m5b_adversarial_candidate_lifecycle() {
let db = YantrikDB::new(":memory:", 8).unwrap();
let move_id = db
.record_move_event(mk_move_input("analogy", &["a"], &["b"]))
.unwrap();
let instance_id = db
.create_adversarial_candidate(
&move_id,
"contradiction",
Some("output claim was retracted within 24h".to_string()),
)
.unwrap();
let candidate = db.get_adversarial_instance(&instance_id).unwrap().unwrap();
assert_eq!(candidate.status, "candidate");
assert_eq!(candidate.discovered_via, "contradiction");
assert!(candidate.traced_root_cause.is_some());
assert!(candidate.generalized_lesson.is_none());
assert!(candidate.lesson_scope_json.is_none());
db.promote_adversarial_candidate(
&instance_id,
"analogy over cross-domain claims with dim≠ input modalities often produces false corroborations".into(),
serde_json::json!({"regimes": ["default"], "move_types": ["analogy"]}).to_string(),
"curator_alice".into(),
).unwrap();
let confirmed = db.get_adversarial_instance(&instance_id).unwrap().unwrap();
assert_eq!(confirmed.status, "confirmed");
assert!(confirmed.generalized_lesson.is_some());
assert!(confirmed.lesson_scope_json.is_some());
assert_eq!(
confirmed.curation_actor_id.as_deref(),
Some("curator_alice")
);
}
#[test]
fn test_m5b_adversarial_promote_rejects_non_candidate() {
let db = YantrikDB::new(":memory:", 8).unwrap();
let move_id = db
.record_move_event(mk_move_input("analogy", &["a"], &["b"]))
.unwrap();
let instance_id = db
.create_adversarial_candidate(&move_id, "contradiction", None)
.unwrap();
db.promote_adversarial_candidate(&instance_id, "lesson".into(), "{}".into(), "c".into())
.unwrap();
let r =
db.promote_adversarial_candidate(&instance_id, "lesson2".into(), "{}".into(), "c".into());
assert!(r.is_err(), "cannot promote non-candidate");
}
#[test]
fn test_m5b_adversarial_promote_requires_non_empty_lesson() {
let db = YantrikDB::new(":memory:", 8).unwrap();
let move_id = db
.record_move_event(mk_move_input("analogy", &["a"], &["b"]))
.unwrap();
let instance_id = db
.create_adversarial_candidate(&move_id, "retraction", None)
.unwrap();
let r = db.promote_adversarial_candidate(&instance_id, "".into(), "{}".into(), "c".into());
assert!(r.is_err());
let r2 = db.promote_adversarial_candidate(&instance_id, "lesson".into(), "".into(), "c".into());
assert!(r2.is_err());
}
#[test]
fn test_m5b_adversarial_reject_candidate() {
let db = YantrikDB::new(":memory:", 8).unwrap();
let move_id = db
.record_move_event(mk_move_input("analogy", &["a"], &["b"]))
.unwrap();
let instance_id = db
.create_adversarial_candidate(&move_id, "calibration_signal", None)
.unwrap();
db.reject_adversarial_candidate(&instance_id, "curator_bob".into())
.unwrap();
let rejected = db.get_adversarial_instance(&instance_id).unwrap().unwrap();
assert_eq!(rejected.status, "rejected");
let r = db.reject_adversarial_candidate(&instance_id, "curator_bob".into());
assert!(r.is_err());
}
#[test]
fn test_m5b_unknown_move_type_warns_but_does_not_reject() {
let db = YantrikDB::new(":memory:", 8).unwrap();
let inp = RecordMoveEventInput {
move_type: "entirely_novel_move_type_not_in_registry".into(),
operator_version: "v0".into(),
observability: observability::OBSERVED.into(),
..Default::default()
};
let r = db.record_move_event(inp);
assert!(r.is_ok(), "unknown move_type must not reject");
let ev = db.get_move_event(&r.unwrap()).unwrap().unwrap();
assert_eq!(ev.move_type, "entirely_novel_move_type_not_in_registry");
assert!(
ev.expected_evaluation_horizon_ms.is_none(),
"unregistered move_type has no default horizon"
);
}
#[test]
fn test_m5b_dependencies_stored_as_json() {
let db = YantrikDB::new(":memory:", 8).unwrap();
let m1 = db
.record_move_event(mk_move_input("analogy", &["a"], &["b"]))
.unwrap();
let mut inp = mk_move_input("decomposition", &["b"], &["c"]);
inp.dependencies = vec![m1.clone()];
let m2 = db.record_move_event(inp).unwrap();
let ev = db.get_move_event(&m2).unwrap().unwrap();
assert!(
ev.dependencies_json.contains(&m1),
"dependencies_json should reference the upstream move_id"
);
}
#[test]
fn test_m5b_append_only_preserves_original_after_correction() {
let db = YantrikDB::new(":memory:", 8).unwrap();
let move_id = db
.record_move_event(mk_move_input("analogy", &["a"], &["b"]))
.unwrap();
db.submit_move_correction(
&move_id,
Some("decomposition".into()),
Some("v2".into()),
None,
"reason".into(),
"curator".into(),
)
.unwrap();
let raw: (String, String) = db
.conn()
.query_row(
"SELECT move_type, operator_version FROM move_events WHERE move_id = ?1",
rusqlite::params![move_id],
|r| Ok((r.get(0)?, r.get(1)?)),
)
.unwrap();
assert_eq!(
raw.0, "analogy",
"move_events row must keep original move_type"
);
assert_eq!(
raw.1, "v1",
"move_events row must keep original operator_version"
);
}
#[test]
fn test_m5b_edge_tables_have_fk_integrity() {
let db = YantrikDB::new(":memory:", 8).unwrap();
let conn = db.conn();
conn.execute("PRAGMA foreign_keys = ON", []).unwrap();
let r = conn.execute(
"INSERT INTO move_input_edge (move_id, claim_id, input_role, ordinal) \
VALUES ('nonexistent_move', 'claim_x', 'input', 0)",
[],
);
assert!(r.is_err(), "FK constraint should reject orphan edge");
}
#[test]
fn test_m5b_adversarial_discovered_via_whitelist() {
let db = YantrikDB::new(":memory:", 8).unwrap();
let move_id = db
.record_move_event(mk_move_input("analogy", &["a"], &["b"]))
.unwrap();
let r = db.create_adversarial_candidate(&move_id, "invalid_source", None);
assert!(r.is_err(), "discovered_via CHECK enforces the whitelist");
}
#[test]
fn test_m5b_seed_registries_idempotent() {
let db = YantrikDB::new(":memory:", 8).unwrap();
let before: i64 = db
.conn()
.query_row("SELECT COUNT(*) FROM move_type_registry", [], |r| r.get(0))
.unwrap();
db.seed_move_registries().unwrap();
let after: i64 = db
.conn()
.query_row("SELECT COUNT(*) FROM move_type_registry", [], |r| r.get(0))
.unwrap();
assert_eq!(before, after, "INSERT OR IGNORE should not add duplicates");
}
#[test]
fn test_m8_axiom_registry_has_core_entries() {
let db = YantrikDB::new(":memory:", 8).unwrap();
let axioms = db.composition_axioms();
assert!(
axioms.len() >= 3,
"axiom registry should have at least 3 entries"
);
let names: Vec<&str> = axioms.iter().map(|a| a.name).collect();
assert!(names.contains(&"decompose_aggregate_identity"));
assert!(names.contains(&"negate_analogize_non_commutative"));
assert!(names.contains(&"source_audit_requires_external_ancestry"));
}
#[test]
fn test_m8_check_composition_non_commutative_match() {
let db = YantrikDB::new(":memory:", 8).unwrap();
let left = db
.record_move_event(mk_move_input("negate_and_test", &["a"], &["b"]))
.unwrap();
let right = db
.record_move_event(mk_move_input("analogy", &["b"], &["c"]))
.unwrap();
let matches = db.check_composition_axioms(&left, &right).unwrap();
assert_eq!(matches.len(), 1);
assert_eq!(matches[0].name, "negate_analogize_non_commutative");
}
#[test]
fn test_m8_check_composition_approx_identity_match() {
let db = YantrikDB::new(":memory:", 8).unwrap();
let d = db
.record_move_event(mk_move_input("decomposition", &["a"], &["b"]))
.unwrap();
let ag = db
.record_move_event(mk_move_input("aggregate_back", &["b"], &["c"]))
.unwrap();
let matches = db.check_composition_axioms(&d, &ag).unwrap();
assert_eq!(matches.len(), 1);
assert_eq!(matches[0].name, "decompose_aggregate_identity");
}
#[test]
fn test_m8_check_composition_no_match() {
let db = YantrikDB::new(":memory:", 8).unwrap();
let m1 = db
.record_move_event(mk_move_input("analogy", &["a"], &["b"]))
.unwrap();
let m2 = db
.record_move_event(mk_move_input("ladder_up", &["b"], &["c"]))
.unwrap();
let matches = db.check_composition_axioms(&m1, &m2).unwrap();
assert!(
matches.is_empty(),
"unrelated move pair should match no axiom"
);
}
#[test]
fn test_m8_source_audit_precondition_violation() {
let db = YantrikDB::new(":memory:", 8).unwrap();
seed_proposition(&db, "p_self_audit");
let conn = db.conn();
for (cid, self_gen, prop) in [
("self_input", 1, None),
("final_out", 1, Some("p_self_audit")),
] {
conn.execute(
"INSERT INTO claims (claim_id, src, dst, rel_type, created_at, \
extractor, polarity, namespace, proposition_id, regime_tag, \
self_generated, source_lineage, modality_signal, weight) \
VALUES (?1, ?2, ?3, 'rel', 0.0, 'ext', 1, 'default', ?4, 'default', \
?5, '[]', 'text', 1.0)",
rusqlite::params![
cid,
format!("src_{}", cid),
format!("dst_{}", cid),
prop,
self_gen
],
)
.unwrap();
}
drop(conn);
db.record_move_event(mk_move_input(
"hypothesis_generation",
&["self_input"],
&["final_out"],
))
.unwrap();
db.compute_write_tier_mobility("p_self_audit", "default")
.unwrap();
db.compute_background_mobility("p_self_audit", "default")
.unwrap();
let violations = db
.check_move_preconditions("source_audit", "p_self_audit", "default")
.unwrap();
assert!(
!violations.is_empty(),
"source_audit on ψ_a=1.0 should violate"
);
assert_eq!(
violations[0].axiom_name,
"source_audit_requires_external_ancestry"
);
assert!(violations[0].observation.contains("self_gen_ancestral"));
}
#[test]
fn test_m8_source_audit_passes_with_external_ancestry() {
let db = YantrikDB::new(":memory:", 8).unwrap();
seed_proposition(&db, "p_ext_audit");
let conn = db.conn();
for (cid, self_gen, prop) in [
("ext_src", 0, None),
("self_src", 1, None),
("final_out", 0, Some("p_ext_audit")),
] {
conn.execute(
"INSERT INTO claims (claim_id, src, dst, rel_type, created_at, \
extractor, polarity, namespace, proposition_id, regime_tag, \
self_generated, source_lineage, modality_signal, weight) \
VALUES (?1, ?2, ?3, 'rel', 0.0, 'ext', 1, 'default', ?4, 'default', \
?5, '[]', 'text', 1.0)",
rusqlite::params![
cid,
format!("src_{}", cid),
format!("dst_{}", cid),
prop,
self_gen
],
)
.unwrap();
}
drop(conn);
db.record_move_event(mk_move_input(
"decomposition",
&["ext_src", "self_src"],
&["final_out"],
))
.unwrap();
db.compute_write_tier_mobility("p_ext_audit", "default")
.unwrap();
db.compute_background_mobility("p_ext_audit", "default")
.unwrap();
let state = db
.get_mobility_state("p_ext_audit", "default")
.unwrap()
.unwrap();
assert!((state.self_gen_ancestral.unwrap() - 0.5).abs() < 1e-9);
let violations = db
.check_move_preconditions("source_audit", "p_ext_audit", "default")
.unwrap();
assert!(
violations.is_empty(),
"source_audit should be fine when external ancestry exists"
);
}
#[test]
fn test_m8_compression_requires_support() {
let db = YantrikDB::new(":memory:", 8).unwrap();
seed_proposition(&db, "p_no_support");
seed_contest_claim(
&db,
"p_no_support",
"c1",
"ext_a",
-1,
"[\"s\"]",
None,
"default",
None,
None,
);
db.compute_write_tier_mobility("p_no_support", "default")
.unwrap();
let violations = db
.check_move_preconditions("compression", "p_no_support", "default")
.unwrap();
assert!(
!violations.is_empty(),
"compression should violate on zero support_mass"
);
assert_eq!(violations[0].axiom_name, "compression_requires_support");
}
#[test]
fn test_m8_hypothesis_generation_blocked_on_present_tense_conflict() {
let db = YantrikDB::new(":memory:", 8).unwrap();
seed_proposition(&db, "p_contested");
seed_contest_claim(
&db,
"p_contested",
"c_sup",
"ext_a",
1,
"[\"s\"]",
None,
"default",
Some(0.0),
Some(20.0),
);
seed_contest_claim(
&db,
"p_contested",
"c_att",
"ext_b",
-1,
"[\"s\"]",
None,
"default",
Some(10.0),
Some(30.0),
);
db.compute_write_tier_mobility("p_contested", "default")
.unwrap();
db.compute_contest_state("p_contested", "default").unwrap();
let violations = db
.check_move_preconditions("hypothesis_generation", "p_contested", "default")
.unwrap();
assert!(
!violations.is_empty(),
"hypothesis_generation should be blocked by PRESENT_TENSE_CONFLICT"
);
assert_eq!(
violations[0].axiom_name,
"hypothesis_generation_skips_present_tense_conflict"
);
}
#[test]
fn test_m8_preconditions_ignore_unrelated_move_types() {
let db = YantrikDB::new(":memory:", 8).unwrap();
seed_proposition(&db, "p_any");
seed_contest_claim(
&db, "p_any", "c1", "ext_a", 1, "[\"s\"]", None, "default", None, None,
);
db.compute_write_tier_mobility("p_any", "default").unwrap();
let violations = db
.check_move_preconditions("analogy", "p_any", "default")
.unwrap();
assert!(violations.is_empty());
}
#[test]
fn test_m9_profile_counts_uses_and_resolutions() {
let db = YantrikDB::new(":memory:", 8).unwrap();
let m1 = db
.record_move_event(mk_move_input("analogy", &["a"], &["b"]))
.unwrap();
let m2 = db
.record_move_event(mk_move_input("analogy", &["c"], &["d"]))
.unwrap();
db.record_move_event(mk_move_input("analogy", &["e"], &["f"]))
.unwrap();
db.record_move_outcome(&m1, "corroborated", None).unwrap();
db.record_move_outcome(&m2, "retracted", None).unwrap();
let profile = db
.recompute_move_type_profile("analogy", "v1", "default")
.unwrap();
assert_eq!(profile.uses_count, 3);
assert_eq!(profile.corroborated_count, 1);
assert_eq!(profile.retracted_count, 1);
assert_eq!(profile.harmful_side_effect_count, 0);
assert_eq!(profile.resolved_count, 2);
assert!((profile.contradiction_introduction_rate.unwrap() - 0.5).abs() < 1e-9);
}
#[test]
fn test_m9_profile_round_trip_via_get() {
let db = YantrikDB::new(":memory:", 8).unwrap();
db.record_move_event(mk_move_input("analogy", &["a"], &["b"]))
.unwrap();
let computed = db
.recompute_move_type_profile("analogy", "v1", "default")
.unwrap();
let read = db
.get_move_type_profile("analogy", "v1", "default")
.unwrap()
.unwrap();
assert_eq!(read.uses_count, computed.uses_count);
assert_eq!(read.resolved_count, computed.resolved_count);
}
#[test]
fn test_m9_recompute_all_keys_present() {
let db = YantrikDB::new(":memory:", 8).unwrap();
db.record_move_event(mk_move_input("analogy", &["a"], &["b"]))
.unwrap();
db.record_move_event(mk_move_input("decomposition", &["c"], &["d"]))
.unwrap();
let count = db.recompute_all_move_type_profiles().unwrap();
assert_eq!(count, 2, "two distinct (type, version, regime) triples");
assert!(db
.get_move_type_profile("analogy", "v1", "default")
.unwrap()
.is_some());
assert!(db
.get_move_type_profile("decomposition", "v1", "default")
.unwrap()
.is_some());
}
#[test]
fn test_m9_profile_missing_returns_none() {
let db = YantrikDB::new(":memory:", 8).unwrap();
let p = db
.get_move_type_profile("analogy", "v1", "default")
.unwrap();
assert!(p.is_none());
}
#[test]
fn test_m9_contradiction_rate_none_when_no_resolutions() {
let db = YantrikDB::new(":memory:", 8).unwrap();
db.record_move_event(mk_move_input("analogy", &["a"], &["b"]))
.unwrap();
let profile = db
.recompute_move_type_profile("analogy", "v1", "default")
.unwrap();
assert_eq!(profile.resolved_count, 0);
assert!(profile.contradiction_introduction_rate.is_none());
}
#[test]
fn test_m10_retraction_auto_files_candidate() {
let db = YantrikDB::new(":memory:", 8).unwrap();
let move_id = db
.record_move_event(mk_move_input("analogy", &["a"], &["b"]))
.unwrap();
let before = db.list_adversarial_for_move(&move_id).unwrap();
assert!(before.is_empty());
db.record_move_outcome(&move_id, "retracted", None).unwrap();
let after = db.list_adversarial_for_move(&move_id).unwrap();
assert_eq!(after.len(), 1, "retraction must auto-file a candidate");
assert_eq!(after[0].status, "candidate");
assert_eq!(after[0].discovered_via, "retraction");
assert!(after[0].traced_root_cause.is_some());
assert!(after[0].generalized_lesson.is_none());
assert!(after[0].lesson_scope_json.is_none());
}
#[test]
fn test_m10_harmful_side_effect_auto_files_candidate() {
let db = YantrikDB::new(":memory:", 8).unwrap();
let move_id = db
.record_move_event(mk_move_input("analogy", &["a"], &["b"]))
.unwrap();
db.record_move_outcome(&move_id, "harmful_side_effect", None)
.unwrap();
let after = db.list_adversarial_for_move(&move_id).unwrap();
assert_eq!(after.len(), 1);
assert_eq!(after[0].discovered_via, "calibration_signal");
}
#[test]
fn test_m10_corroborated_outcome_does_not_file_candidate() {
let db = YantrikDB::new(":memory:", 8).unwrap();
let move_id = db
.record_move_event(mk_move_input("analogy", &["a"], &["b"]))
.unwrap();
db.record_move_outcome(&move_id, "corroborated", None)
.unwrap();
let after = db.list_adversarial_for_move(&move_id).unwrap();
assert!(
after.is_empty(),
"corroborated outcome must not file an adversarial candidate"
);
}
#[test]
fn test_m10_contest_flag_transition_auto_files_for_producing_move() {
let db = YantrikDB::new(":memory:", 8).unwrap();
seed_proposition(&db, "p_flag_auto");
seed_contest_claim(
&db,
"p_flag_auto",
"claim_out",
"ext_a",
1,
"[\"s\"]",
None,
"default",
None,
None,
);
let move_id = db
.record_move_event(mk_move_input(
"hypothesis_generation",
&["evidence"],
&["claim_out"],
))
.unwrap();
db.compute_write_tier_mobility("p_flag_auto", "default")
.unwrap();
db.compute_contest_state("p_flag_auto", "default").unwrap();
assert!(db.list_adversarial_for_move(&move_id).unwrap().is_empty());
seed_contest_claim(
&db,
"p_flag_auto",
"conflicting",
"ext_b",
-1,
"[\"s\"]",
None,
"default",
None,
None,
);
db.compute_write_tier_mobility("p_flag_auto", "default")
.unwrap();
db.compute_contest_state("p_flag_auto", "default").unwrap();
let after = db.list_adversarial_for_move(&move_id).unwrap();
assert!(
!after.is_empty(),
"SAME_SOURCE_CONFLICT should auto-file adversarial candidate for producing move"
);
assert_eq!(after[0].discovered_via, "contradiction");
assert_eq!(after[0].status, "candidate");
}
#[test]
fn test_m10_contest_flag_transition_dedups_on_repeat_recompute() {
let db = YantrikDB::new(":memory:", 8).unwrap();
seed_proposition(&db, "p_dedup");
seed_contest_claim(
&db,
"p_dedup",
"claim_out",
"ext_a",
1,
"[\"s\"]",
None,
"default",
None,
None,
);
let move_id = db
.record_move_event(mk_move_input("analogy", &["e"], &["claim_out"]))
.unwrap();
db.compute_write_tier_mobility("p_dedup", "default")
.unwrap();
db.compute_contest_state("p_dedup", "default").unwrap();
seed_contest_claim(
&db, "p_dedup", "conflict", "ext_b", -1, "[\"s\"]", None, "default", None, None,
);
db.compute_write_tier_mobility("p_dedup", "default")
.unwrap();
db.compute_contest_state("p_dedup", "default").unwrap();
db.compute_contest_state("p_dedup", "default").unwrap();
db.compute_contest_state("p_dedup", "default").unwrap();
db.compute_contest_state("p_dedup", "default").unwrap();
let candidates = db.list_adversarial_for_move(&move_id).unwrap();
let contradiction_candidates: Vec<_> = candidates
.iter()
.filter(|c| c.discovered_via == "contradiction")
.collect();
assert_eq!(
contradiction_candidates.len(),
1,
"repeat contest recompute must not duplicate contradiction candidates"
);
}
#[test]
fn test_m10_m9_feedback_loop_retraction_updates_profile() {
let db = YantrikDB::new(":memory:", 8).unwrap();
let m1 = db
.record_move_event(mk_move_input("analogy", &["a"], &["b"]))
.unwrap();
db.record_move_event(mk_move_input("analogy", &["c"], &["d"]))
.unwrap();
db.record_move_outcome(&m1, "retracted", None).unwrap();
let candidates = db.list_adversarial_for_move(&m1).unwrap();
assert_eq!(candidates.len(), 1);
let profile = db
.recompute_move_type_profile("analogy", "v1", "default")
.unwrap();
assert_eq!(profile.retracted_count, 1);
assert_eq!(profile.uses_count, 2);
}
#[test]
fn test_m6_temporal_coherence_stable_polarity_is_one() {
let db = YantrikDB::new(":memory:", 8).unwrap();
seed_proposition(&db, "p_coh");
seed_contest_claim(
&db, "p_coh", "c1", "ext_a", 1, "[\"s\"]", None, "default", None, None,
);
seed_contest_claim(
&db, "p_coh", "c2", "ext_b", 1, "[\"s\"]", None, "default", None, None,
);
seed_contest_claim(
&db, "p_coh", "c3", "ext_c", 1, "[\"s\"]", None, "default", None, None,
);
db.compute_write_tier_mobility("p_coh", "default").unwrap();
let state = db
.compute_background_mobility("p_coh", "default")
.unwrap()
.unwrap();
assert_eq!(state.temporal_coherence, Some(1.0));
}
#[test]
fn test_m6_temporal_coherence_flips_reduce_score() {
let db = YantrikDB::new(":memory:", 8).unwrap();
seed_proposition(&db, "p_flip");
let conn = db.conn();
let triples = [
("c_p1", "ext_a", 1, 1.0),
("c_a1", "ext_b", -1, 2.0),
("c_p2", "ext_c", 1, 3.0),
];
for (cid, ext, pol, ts) in triples {
conn.execute(
"INSERT INTO claims (claim_id, src, dst, rel_type, created_at, \
extractor, polarity, namespace, proposition_id, regime_tag, \
self_generated, source_lineage, modality_signal, weight) \
VALUES (?1, 'src_p_flip', 'dst_p_flip', 'rel_p_flip', ?2, ?3, ?4, \
'default', 'p_flip', 'default', 0, '[\"s\"]', 'text', 1.0)",
rusqlite::params![cid, ts, ext, pol],
)
.unwrap();
}
drop(conn);
db.compute_write_tier_mobility("p_flip", "default").unwrap();
let state = db
.compute_background_mobility("p_flip", "default")
.unwrap()
.unwrap();
let tau = state.temporal_coherence.unwrap();
assert!(
(tau - 0.0).abs() < 1e-9,
"expected τ=0 for maximally flipping polarity, got {}",
tau
);
}
#[test]
fn test_m6_temporal_coherence_single_claim_is_one_by_convention() {
let db = YantrikDB::new(":memory:", 8).unwrap();
seed_proposition(&db, "p_solo");
seed_contest_claim(
&db, "p_solo", "c1", "ext_a", 1, "[\"s\"]", None, "default", None, None,
);
db.compute_write_tier_mobility("p_solo", "default").unwrap();
let state = db
.compute_background_mobility("p_solo", "default")
.unwrap()
.unwrap();
assert_eq!(state.temporal_coherence, Some(1.0));
}
#[test]
fn test_m6_load_bearingness_counts_downstream_moves() {
let db = YantrikDB::new(":memory:", 8).unwrap();
seed_proposition(&db, "p_load");
seed_contest_claim(
&db,
"p_load",
"upstream_a",
"ext_a",
1,
"[\"s\"]",
None,
"default",
None,
None,
);
seed_contest_claim(
&db,
"p_load",
"upstream_b",
"ext_b",
1,
"[\"s\"]",
None,
"default",
None,
None,
);
db.compute_write_tier_mobility("p_load", "default").unwrap();
let before = db
.compute_background_mobility("p_load", "default")
.unwrap()
.unwrap();
assert_eq!(before.load_bearingness, Some(0.0));
db.record_move_event(mk_move_input("analogy", &["upstream_a"], &["derived_1"]))
.unwrap();
db.record_move_event(mk_move_input(
"decomposition",
&["upstream_a", "upstream_b"],
&["derived_2"],
))
.unwrap();
db.record_move_event(mk_move_input(
"analogy",
&["unrelated_claim"],
&["derived_3"],
))
.unwrap();
let after = db
.compute_background_mobility("p_load", "default")
.unwrap()
.unwrap();
assert_eq!(
after.load_bearingness,
Some(2.0),
"expected 2 downstream moves consuming this proposition's claims"
);
}
#[test]
fn test_m6_self_gen_ancestral_traces_backward() {
let db = YantrikDB::new(":memory:", 8).unwrap();
seed_proposition(&db, "p_anc");
let conn = db.conn();
for (cid, self_gen) in [
("evidence_a", 1),
("evidence_b", 0),
("intermediate_c", 0),
("final_d", 0),
] {
let prop = if cid == "final_d" {
Some("p_anc")
} else {
None
};
conn.execute(
"INSERT INTO claims (claim_id, src, dst, rel_type, created_at, \
extractor, polarity, namespace, proposition_id, regime_tag, \
self_generated, source_lineage, modality_signal, weight) \
VALUES (?1, ?2, ?3, 'rel_anc', 0.0, 'ext', 1, 'default', ?4, \
'default', ?5, '[]', 'text', 1.0)",
rusqlite::params![
cid,
format!("src_{}", cid),
format!("dst_{}", cid),
prop,
self_gen
],
)
.unwrap();
}
drop(conn);
let mut inp1 = mk_move_input(
"decomposition",
&["evidence_a", "evidence_b"],
&["intermediate_c"],
);
inp1.operator_version = "v1".to_string();
db.record_move_event(inp1).unwrap();
let mut inp2 = mk_move_input("ladder_up", &["intermediate_c"], &["final_d"]);
inp2.operator_version = "v1".to_string();
db.record_move_event(inp2).unwrap();
db.compute_write_tier_mobility("p_anc", "default").unwrap();
let state = db
.compute_background_mobility("p_anc", "default")
.unwrap()
.unwrap();
let psi_a = state.self_gen_ancestral.unwrap();
assert!(
(psi_a - 1.0 / 3.0).abs() < 1e-9,
"expected ψ_a = 1/3, got {}",
psi_a
);
}
#[test]
fn test_m6_self_gen_ancestral_no_moves_is_zero() {
let db = YantrikDB::new(":memory:", 8).unwrap();
seed_proposition(&db, "p_leaf");
seed_contest_claim(
&db, "p_leaf", "c1", "ext_a", 1, "[\"s\"]", None, "default", None, None,
);
db.compute_write_tier_mobility("p_leaf", "default").unwrap();
let state = db
.compute_background_mobility("p_leaf", "default")
.unwrap()
.unwrap();
assert_eq!(state.self_gen_ancestral, Some(0.0));
}
#[test]
fn test_m6_compute_background_missing_returns_none() {
let db = YantrikDB::new(":memory:", 8).unwrap();
let res = db.compute_background_mobility("nope", "default").unwrap();
assert!(res.is_none());
}
#[test]
fn test_m6_background_idempotent() {
let db = YantrikDB::new(":memory:", 8).unwrap();
seed_proposition(&db, "p_idem_bg");
seed_contest_claim(
&db,
"p_idem_bg",
"c1",
"ext_a",
1,
"[\"s\"]",
None,
"default",
None,
None,
);
db.compute_write_tier_mobility("p_idem_bg", "default")
.unwrap();
let first = db
.compute_background_mobility("p_idem_bg", "default")
.unwrap()
.unwrap();
let second = db
.compute_background_mobility("p_idem_bg", "default")
.unwrap()
.unwrap();
assert_eq!(first.temporal_coherence, second.temporal_coherence);
assert_eq!(first.load_bearingness, second.load_bearingness);
assert_eq!(first.self_gen_ancestral, second.self_gen_ancestral);
}
#[test]
fn test_m6_background_batch_scan() {
let db = YantrikDB::new(":memory:", 8).unwrap();
for pid in ["p_batch_1", "p_batch_2", "p_batch_3"] {
seed_proposition(&db, pid);
seed_contest_claim(
&db,
pid,
&format!("c_{}", pid),
"ext_a",
1,
"[\"s\"]",
None,
"default",
None,
None,
);
db.compute_write_tier_mobility(pid, "default").unwrap();
}
let pending_before: i64 = db
.conn()
.query_row(
"SELECT COUNT(*) FROM mobility_state WHERE temporal_coherence IS NULL",
[],
|r| r.get(0),
)
.unwrap();
assert_eq!(pending_before, 3);
let processed = db.recompute_background_mobility_batch(10).unwrap();
assert_eq!(processed, 3);
let pending_after: i64 = db
.conn()
.query_row(
"SELECT COUNT(*) FROM mobility_state WHERE temporal_coherence IS NULL",
[],
|r| r.get(0),
)
.unwrap();
assert_eq!(pending_after, 0);
}
#[test]
fn test_m6_background_marks_tier_components() {
let db = YantrikDB::new(":memory:", 8).unwrap();
seed_proposition(&db, "p_tier");
seed_contest_claim(
&db, "p_tier", "c1", "ext_a", 1, "[\"s\"]", None, "default", None, None,
);
db.compute_write_tier_mobility("p_tier", "default").unwrap();
let state = db
.compute_background_mobility("p_tier", "default")
.unwrap()
.unwrap();
for expected in [
"temporal_coherence",
"load_bearingness",
"self_gen_ancestral",
] {
assert!(
state.tier_bg_components.iter().any(|c| c == expected),
"tier_bg_components should include {}, got {:?}",
expected,
state.tier_bg_components
);
}
let state2 = db
.compute_background_mobility("p_tier", "default")
.unwrap()
.unwrap();
let tc_count = state2
.tier_bg_components
.iter()
.filter(|c| *c == "temporal_coherence")
.count();
assert_eq!(
tc_count, 1,
"repeated calls must not duplicate tier_bg_components entries"
);
}
#[test]
fn test_m5b_full_lifecycle_observed_to_retracted_with_adversarial() {
let db = YantrikDB::new(":memory:", 8).unwrap();
let move_id = db
.record_move_event(mk_move_input(
"hypothesis_generation",
&["evidence_a", "evidence_b"],
&["hypothesis_h"],
))
.unwrap();
db.record_move_outcome(
&move_id,
posthoc_outcome::RETRACTED,
Some(
serde_json::json!({"retraction_cause": "contradiction with high-weight claim"})
.to_string(),
),
)
.unwrap();
let auto_candidates = db.list_adversarial_for_move(&move_id).unwrap();
assert_eq!(auto_candidates.len(), 1);
assert_eq!(auto_candidates[0].status, adversarial_status::CANDIDATE);
assert_eq!(auto_candidates[0].discovered_via, "retraction");
db.promote_adversarial_candidate(
&auto_candidates[0].instance_id,
"hypothesis_generation from ≤2 evidence sources is prone to retraction".into(),
serde_json::json!({
"regimes": ["default"],
"move_types": ["hypothesis_generation"],
"input_signatures": {"min_inputs": 2}
})
.to_string(),
"curator_dana".into(),
)
.unwrap();
let ev = db.get_move_event(&move_id).unwrap().unwrap();
assert_eq!(ev.posthoc_outcome.as_deref(), Some("retracted"));
let instance = db
.get_adversarial_instance(&auto_candidates[0].instance_id)
.unwrap()
.unwrap();
assert_eq!(instance.status, adversarial_status::CONFIRMED);
let instances = db.list_adversarial_for_move(&move_id).unwrap();
assert_eq!(instances.len(), 1);
}
#[test]
fn test_insert_vector_makes_recall_find_it() {
let db = YantrikDB::new(":memory:", 8).unwrap();
let _ = db
.record(
"seed",
"episodic",
0.5,
0.0,
604800.0,
&empty_meta(),
&vec_seed(0.1, 8),
"default",
0.8,
"general",
"user",
None,
)
.unwrap();
let synthetic_rid = "test-synthetic-rid-1";
let synthetic_emb = vec_seed(0.9, 8);
db.insert_vector(synthetic_rid, &synthetic_emb).unwrap();
let results = db
.recall(
&synthetic_emb,
5,
None,
None,
false,
false,
None,
true,
None,
None,
None,
None,
None,
)
.unwrap();
let stats = db.stats(None).unwrap();
assert!(
stats.vec_index_entries >= 2,
"vec_index_entries should be at least 2 (record + insert_vector); got {}",
stats.vec_index_entries
);
assert!(results.len() <= 5);
}
#[test]
fn test_insert_vector_idempotent_on_same_rid() {
let db = YantrikDB::new(":memory:", 8).unwrap();
let rid = "idempotency-test";
let emb = vec_seed(0.5, 8);
db.insert_vector(rid, &emb).unwrap();
db.insert_vector(rid, &emb).unwrap();
}
#[test]
fn test_encrypt_embedding_pub_unencrypted_returns_input_unchanged() {
let db = YantrikDB::new(":memory:", 8).unwrap();
let raw: Vec<u8> = (0u8..32).collect();
let out = db.encrypt_embedding_pub(&raw).unwrap();
assert_eq!(out, raw, "no-encryption path must return input unchanged");
}
#[test]
fn test_encrypt_embedding_pub_with_encryption_returns_ciphertext() {
let master_key = [0xAB; 32];
let db = YantrikDB::new_encrypted(":memory:", 8, &master_key).unwrap();
let raw: Vec<u8> = (0u8..32).collect();
let out = db.encrypt_embedding_pub(&raw).unwrap();
assert_ne!(
out, raw,
"encrypted path must produce ciphertext, not plaintext"
);
assert!(
out.len() > raw.len(),
"encrypted blob should be longer than plaintext (nonce + tag overhead)"
);
}
#[test]
fn record_with_rid_basic_succeeds() {
let db = YantrikDB::new(":memory:", 64).unwrap();
let emb = vec_seed(1.0, 64);
db.record_with_rid(
"rid_test_1",
"the quick brown fox",
"episodic",
0.5,
0.0,
604800.0,
&empty_meta(),
&emb,
"default",
0.8,
"general",
"user",
None,
1_700_000_000_000_000,
&[],
"test-model.v1",
None,
)
.expect("record_with_rid succeeds");
let row = db.get("rid_test_1").unwrap().unwrap();
assert_eq!(row.rid, "rid_test_1");
assert_eq!(row.text, "the quick brown fox");
assert_eq!(row.memory_type, "episodic");
}
#[test]
fn record_with_rid_persists_v25_columns() {
let db = YantrikDB::new(":memory:", 64).unwrap();
let emb = vec_seed(2.0, 64);
db.record_with_rid(
"rid_v25",
"test v25 columns",
"semantic",
0.5,
0.0,
604800.0,
&empty_meta(),
&emb,
"default",
0.8,
"general",
"user",
None,
1_700_000_000_000_000,
&[],
"bge-base-en-v1.5",
None,
)
.unwrap();
let conn = db.read_conn();
let (cum, model): (i64, Option<String>) = conn
.query_row(
"SELECT created_at_unix_micros, embedding_model FROM memories WHERE rid = ?1",
rusqlite::params!["rid_v25"],
|row| Ok((row.get(0)?, row.get(1)?)),
)
.unwrap();
assert_eq!(cum, 1_700_000_000_000_000);
assert_eq!(model.as_deref(), Some("bge-base-en-v1.5"));
}
#[test]
fn record_with_rid_is_idempotent_on_replay() {
let db = YantrikDB::new(":memory:", 64).unwrap();
let emb = vec_seed(3.0, 64);
let entities = ["Alice", "Acme"];
let entity_refs: Vec<&str> = entities.iter().copied().collect();
for _ in 0..3 {
db.record_with_rid(
"rid_idem",
"Alice works at Acme",
"episodic",
0.5,
0.0,
604800.0,
&empty_meta(),
&emb,
"default",
0.8,
"general",
"user",
None,
1_700_000_001_000_000,
&entity_refs,
"test-model.v1",
None,
)
.expect("idempotent re-apply");
}
db.apply_pending_ops_once(100).unwrap();
let conn = db.read_conn();
let memory_count: i64 = conn
.query_row(
"SELECT COUNT(*) FROM memories WHERE rid = ?1",
rusqlite::params!["rid_idem"],
|row| row.get(0),
)
.unwrap();
assert_eq!(memory_count, 1, "memories has exactly one row");
let me_count: i64 = conn
.query_row(
"SELECT COUNT(*) FROM memory_entities WHERE memory_rid = ?1",
rusqlite::params!["rid_idem"],
|row| row.get(0),
)
.unwrap();
assert_eq!(
me_count, 2,
"memory_entities has 2 rows (Alice, Acme), no doubles"
);
let mc: i64 = conn
.query_row(
"SELECT mention_count FROM entities WHERE name = 'Alice'",
[],
|row| row.get(0),
)
.unwrap();
assert_eq!(mc, 1, "mention_count not bumped on replay");
let op_count: i64 = conn
.query_row(
"SELECT COUNT(*) FROM oplog WHERE op_type = 'record_with_rid' AND target_rid = ?1",
rusqlite::params!["rid_idem"],
|row| row.get(0),
)
.unwrap();
assert_eq!(op_count, 1, "oplog has exactly one record_with_rid entry");
}
#[test]
fn record_with_rid_rejects_dimension_mismatch() {
let db = YantrikDB::new(":memory:", 64).unwrap();
let bad = vec![0.0f32; 32]; let err = db
.record_with_rid(
"rid_bad_dim",
"x",
"episodic",
0.5,
0.0,
604800.0,
&empty_meta(),
&bad,
"default",
0.8,
"general",
"user",
None,
1_700_000_002_000_000,
&[],
"test-model.v1",
None,
)
.expect_err("must reject");
match err {
crate::error::YantrikDbError::EmbeddingDimensionMismatch { expected, got } => {
assert_eq!(expected, 64);
assert_eq!(got, 32);
}
other => panic!("expected EmbeddingDimensionMismatch, got {other:?}"),
}
assert!(db.get("rid_bad_dim").unwrap().is_none());
}
#[test]
fn record_with_rid_uses_caller_supplied_timestamp() {
let db = YantrikDB::new(":memory:", 64).unwrap();
let emb = vec_seed(4.0, 64);
let caller_ts: i64 = 1_700_000_005_000_000;
db.record_with_rid(
"rid_ts",
"test ts",
"episodic",
0.5,
0.0,
604800.0,
&empty_meta(),
&emb,
"default",
0.8,
"general",
"user",
None,
caller_ts,
&[],
"test-model.v1",
None,
)
.unwrap();
let conn = db.read_conn();
let (cat_real, cat_micros): (f64, i64) = conn
.query_row(
"SELECT created_at, created_at_unix_micros FROM memories WHERE rid = ?1",
rusqlite::params!["rid_ts"],
|row| Ok((row.get(0)?, row.get(1)?)),
)
.unwrap();
assert_eq!(cat_micros, caller_ts);
let expected_real = (caller_ts as f64) / 1_000_000.0;
assert!(
(cat_real - expected_real).abs() < 1e-6,
"created_at REAL should reflect caller timestamp: got {} expected {}",
cat_real,
expected_real
);
}
#[test]
fn record_with_rid_makes_recall_find_it() {
let db = YantrikDB::new(":memory:", 64).unwrap();
let emb = vec_seed(7.0, 64);
db.record_with_rid(
"rid_recall",
"memory inserted via record_with_rid",
"episodic",
0.7,
0.0,
604800.0,
&empty_meta(),
&emb,
"default",
0.8,
"general",
"user",
None,
1_700_000_006_000_000,
&[],
"test-model.v1",
None,
)
.unwrap();
let results = db
.recall(
&emb, 5, None, None, false, false, None, true, None, None, None, None, None,
)
.unwrap();
assert!(
results.iter().any(|r| r.rid == "rid_recall"),
"rid_recall should appear in recall results"
);
}
#[test]
fn tombstone_with_rid_basic_succeeds() {
let db = YantrikDB::new(":memory:", 64).unwrap();
let emb = vec_seed(1.0, 64);
let rid = db
.record(
"to tombstone",
"episodic",
0.5,
0.0,
604800.0,
&empty_meta(),
&emb,
"default",
0.8,
"general",
"user",
None,
)
.unwrap();
db.tombstone_with_rid(
&rid,
"default",
Some("test reason"),
1_700_000_010_000_000,
None,
)
.expect("tombstone_with_rid succeeds");
let mem = db.get(&rid).unwrap().unwrap();
assert_eq!(mem.consolidation_status, "tombstoned");
}
#[test]
fn tombstone_with_rid_persists_reason() {
let db = YantrikDB::new(":memory:", 64).unwrap();
let emb = vec_seed(2.0, 64);
let rid = db
.record(
"memory with reason",
"episodic",
0.5,
0.0,
604800.0,
&empty_meta(),
&emb,
"default",
0.8,
"general",
"user",
None,
)
.unwrap();
db.tombstone_with_rid(
&rid,
"default",
Some("user requested deletion"),
1_700_000_011_000_000,
None,
)
.unwrap();
let conn = db.read_conn();
let reason: Option<String> = conn
.query_row(
"SELECT tombstone_reason FROM memories WHERE rid = ?1",
rusqlite::params![&rid],
|row| row.get(0),
)
.unwrap();
assert_eq!(reason.as_deref(), Some("user requested deletion"));
}
#[test]
fn tombstone_with_rid_idempotent_on_replay() {
let db = YantrikDB::new(":memory:", 64).unwrap();
let emb = vec_seed(3.0, 64);
let rid = db
.record(
"idempotent tombstone",
"episodic",
0.5,
0.0,
604800.0,
&empty_meta(),
&emb,
"default",
0.8,
"general",
"user",
None,
)
.unwrap();
for _ in 0..3 {
db.tombstone_with_rid(&rid, "default", Some("replay"), 1_700_000_012_000_000, None)
.expect("idempotent re-apply");
}
let conn = db.read_conn();
let op_count: i64 = conn
.query_row(
"SELECT COUNT(*) FROM oplog WHERE op_type = 'forget' AND target_rid = ?1",
rusqlite::params![&rid],
|row| row.get(0),
)
.unwrap();
assert_eq!(
op_count, 1,
"oplog has exactly one forget entry despite 3 calls"
);
}
#[test]
fn tombstone_with_rid_idempotent_on_missing() {
let db = YantrikDB::new(":memory:", 64).unwrap();
db.tombstone_with_rid(
"rid_never_existed",
"default",
None,
1_700_000_013_000_000,
None,
)
.expect("must be Ok(()) on missing rid");
let conn = db.read_conn();
let op_count: i64 = conn
.query_row(
"SELECT COUNT(*) FROM oplog WHERE target_rid = ?1",
rusqlite::params!["rid_never_existed"],
|row| row.get(0),
)
.unwrap();
assert_eq!(op_count, 0);
}
#[test]
fn tombstone_with_rid_uses_caller_supplied_timestamp() {
let db = YantrikDB::new(":memory:", 64).unwrap();
let emb = vec_seed(4.0, 64);
let rid = db
.record(
"ts test",
"episodic",
0.5,
0.0,
604800.0,
&empty_meta(),
&emb,
"default",
0.8,
"general",
"user",
None,
)
.unwrap();
let caller_ts: i64 = 1_700_000_999_000_000;
db.tombstone_with_rid(&rid, "default", None, caller_ts, None)
.unwrap();
let conn = db.read_conn();
let updated_at: f64 = conn
.query_row(
"SELECT updated_at FROM memories WHERE rid = ?1",
rusqlite::params![&rid],
|row| row.get(0),
)
.unwrap();
let expected = (caller_ts as f64) / 1_000_000.0;
assert!(
(updated_at - expected).abs() < 1e-6,
"updated_at should reflect caller ts: got {} expected {}",
updated_at,
expected
);
}
#[test]
fn forget_still_works_after_refactor() {
let db = YantrikDB::new(":memory:", 64).unwrap();
let emb = vec_seed(5.0, 64);
let rid = db
.record(
"forget test",
"episodic",
0.5,
0.0,
604800.0,
&empty_meta(),
&emb,
"default",
0.8,
"general",
"user",
None,
)
.unwrap();
let first = db.forget(&rid).unwrap();
assert!(first, "first forget on live row returns true");
let second = db.forget(&rid).unwrap();
assert!(
!second,
"second forget on already-tombstoned row returns false"
);
let missing = db.forget("rid_never_existed").unwrap();
assert!(!missing, "forget on missing rid returns false");
}
#[test]
fn tombstone_with_rid_hides_from_recall() {
let db = YantrikDB::new(":memory:", 64).unwrap();
let emb = vec_seed(6.0, 64);
let rid = db
.record(
"hide me",
"episodic",
0.5,
0.0,
604800.0,
&empty_meta(),
&emb,
"default",
0.8,
"general",
"user",
None,
)
.unwrap();
let r = db
.recall(
&emb, 5, None, None, false, false, None, true, None, None, None, None, None,
)
.unwrap();
assert!(r.iter().any(|x| x.rid == rid), "visible before tombstone");
db.tombstone_with_rid(&rid, "default", None, 1_700_000_014_000_000, None)
.unwrap();
let r2 = db
.recall(
&emb, 5, None, None, false, false, None, true, None, None, None, None, None,
)
.unwrap();
assert!(!r2.iter().any(|x| x.rid == rid), "hidden after tombstone");
}
#[test]
fn upsert_entity_edge_with_id_basic_succeeds() {
let db = YantrikDB::new(":memory:", 64).unwrap();
db.upsert_entity_edge_with_id(
"edge_1",
"Alice",
"Acme",
"works_at",
0.9,
"default",
1_700_000_020_000_000,
None,
)
.expect("upsert succeeds");
let conn = db.read_conn();
let (cid, src, dst, rel, weight): (String, String, String, String, f64) = conn
.query_row(
"SELECT claim_id, src, dst, rel_type, weight FROM claims WHERE claim_id = ?1",
rusqlite::params!["edge_1"],
|row| {
Ok((
row.get(0)?,
row.get(1)?,
row.get(2)?,
row.get(3)?,
row.get(4)?,
))
},
)
.unwrap();
assert_eq!(cid, "edge_1");
assert_eq!(src, "Alice");
assert_eq!(dst, "Acme");
assert_eq!(rel, "works_at");
assert!((weight - 0.9).abs() < 1e-6);
}
#[test]
fn upsert_entity_edge_with_id_is_idempotent_on_replay() {
let db = YantrikDB::new(":memory:", 64).unwrap();
for _ in 0..3 {
db.upsert_entity_edge_with_id(
"edge_idem",
"Bob",
"Beta Corp",
"founded",
0.8,
"default",
1_700_000_021_000_000,
None,
)
.expect("idempotent");
}
let conn = db.read_conn();
let count: i64 = conn
.query_row(
"SELECT COUNT(*) FROM claims WHERE claim_id = ?1",
rusqlite::params!["edge_idem"],
|row| row.get(0),
)
.unwrap();
assert_eq!(count, 1, "exactly one claim row regardless of replay");
let op_count: i64 = conn.query_row(
"SELECT COUNT(*) FROM oplog WHERE op_type = 'upsert_entity_edge_with_id' AND target_rid = ?1",
rusqlite::params!["edge_idem"],
|row| row.get(0),
).unwrap();
assert_eq!(op_count, 1, "exactly one oplog entry regardless of replay");
}
#[test]
fn upsert_entity_edge_uses_caller_supplied_timestamp() {
let db = YantrikDB::new(":memory:", 64).unwrap();
let caller_ts: i64 = 1_700_000_555_000_000;
db.upsert_entity_edge_with_id(
"edge_ts", "X", "Y", "knows", 0.5, "default", caller_ts, None,
)
.unwrap();
let conn = db.read_conn();
let created_at: f64 = conn
.query_row(
"SELECT created_at FROM claims WHERE claim_id = ?1",
rusqlite::params!["edge_ts"],
|row| row.get(0),
)
.unwrap();
let expected = (caller_ts as f64) / 1_000_000.0;
assert!(
(created_at - expected).abs() < 1e-6,
"created_at REAL reflects caller ts: got {} expected {}",
created_at,
expected
);
}
#[test]
fn upsert_entity_edge_creates_entities() {
let db = YantrikDB::new(":memory:", 64).unwrap();
db.upsert_entity_edge_with_id(
"edge_ent",
"Charlie",
"Delta Inc",
"ceo_of",
1.0,
"default",
1_700_000_022_000_000,
None,
)
.unwrap();
let conn = db.read_conn();
let charlie: i64 = conn
.query_row(
"SELECT COUNT(*) FROM entities WHERE name = 'Charlie'",
[],
|row| row.get(0),
)
.unwrap();
let delta: i64 = conn
.query_row(
"SELECT COUNT(*) FROM entities WHERE name = 'Delta Inc'",
[],
|row| row.get(0),
)
.unwrap();
assert_eq!(charlie, 1);
assert_eq!(delta, 1);
}
#[test]
fn delete_entity_edge_with_id_basic_succeeds() {
let db = YantrikDB::new(":memory:", 64).unwrap();
db.upsert_entity_edge_with_id(
"edge_del",
"A",
"B",
"knows",
0.5,
"default",
1_700_000_023_000_000,
None,
)
.unwrap();
db.delete_entity_edge_with_id("edge_del", "default", 1_700_000_024_000_000, None)
.expect("delete succeeds");
let conn = db.read_conn();
let tombstoned: i64 = conn
.query_row(
"SELECT tombstoned FROM claims WHERE claim_id = ?1",
rusqlite::params!["edge_del"],
|row| row.get(0),
)
.unwrap();
assert_eq!(tombstoned, 1);
}
#[test]
fn delete_entity_edge_with_id_idempotent_on_missing() {
let db = YantrikDB::new(":memory:", 64).unwrap();
db.delete_entity_edge_with_id("edge_never", "default", 1_700_000_025_000_000, None)
.expect("missing edge: ok");
let conn = db.read_conn();
let op_count: i64 = conn.query_row(
"SELECT COUNT(*) FROM oplog WHERE op_type = 'delete_entity_edge_with_id' AND target_rid = ?1",
rusqlite::params!["edge_never"],
|row| row.get(0),
).unwrap();
assert_eq!(op_count, 0, "no oplog noise for missing edge delete");
}
#[test]
fn delete_entity_edge_with_id_idempotent_on_replay() {
let db = YantrikDB::new(":memory:", 64).unwrap();
db.upsert_entity_edge_with_id(
"edge_del2",
"P",
"Q",
"knows",
0.5,
"default",
1_700_000_026_000_000,
None,
)
.unwrap();
for _ in 0..3 {
db.delete_entity_edge_with_id("edge_del2", "default", 1_700_000_027_000_000, None)
.expect("idempotent");
}
let conn = db.read_conn();
let op_count: i64 = conn.query_row(
"SELECT COUNT(*) FROM oplog WHERE op_type = 'delete_entity_edge_with_id' AND target_rid = ?1",
rusqlite::params!["edge_del2"],
|row| row.get(0),
).unwrap();
assert_eq!(
op_count, 1,
"exactly one delete oplog entry across 3 replays"
);
}
#[test]
fn record_with_rid_uses_caller_supplied_seq_and_bumps_visible() {
let db = YantrikDB::new(":memory:", 64).unwrap();
let emb = vec_seed(11.0, 64);
db.record_with_rid(
"rid_seq_supplied",
"x",
"episodic",
0.5,
0.0,
604800.0,
&empty_meta(),
&emb,
"alpha",
0.8,
"general",
"user",
None,
1_700_000_100_000_000,
&[],
"test-model.v1",
Some(1_000_000),
)
.unwrap();
assert_eq!(
db.visible_seq_for("alpha"),
1_000_000,
"visible_seq[alpha] equals caller-supplied seq"
);
db.record_with_rid(
"rid_after_ratchet",
"y",
"episodic",
0.5,
0.0,
604800.0,
&empty_meta(),
&vec_seed(12.0, 64),
"alpha",
0.8,
"general",
"user",
None,
1_700_000_101_000_000,
&[],
"test-model.v1",
None,
)
.unwrap();
assert!(
db.visible_seq_for("alpha") > 1_000_000,
"engine-allocated seq is > ratcheted high-water"
);
}
#[test]
fn tombstone_with_rid_bumps_visible_seq_even_when_rid_missing() {
let db = YantrikDB::new(":memory:", 64).unwrap();
assert_eq!(db.visible_seq_for("beta"), 0);
db.tombstone_with_rid(
"rid_unknown_locally",
"beta",
None,
1_700_000_200_000_000,
Some(2_000_000),
)
.unwrap();
assert_eq!(
db.visible_seq_for("beta"),
2_000_000,
"tombstone_with_rid bumps visible_seq[beta] even on missing rid"
);
}
#[test]
fn upsert_entity_edge_with_id_bumps_visible_seq() {
let db = YantrikDB::new(":memory:", 64).unwrap();
db.upsert_entity_edge_with_id(
"edge_seq",
"X",
"Y",
"knows",
0.5,
"gamma",
1_700_000_300_000_000,
Some(3_000_000),
)
.unwrap();
assert_eq!(db.visible_seq_for("gamma"), 3_000_000);
db.upsert_entity_edge_with_id(
"edge_seq",
"X",
"Y",
"knows",
0.5,
"gamma",
1_700_000_300_000_000,
Some(3_000_000),
)
.unwrap();
assert_eq!(
db.visible_seq_for("gamma"),
3_000_000,
"same-seq replay does not regress watermark"
);
db.upsert_entity_edge_with_id(
"edge_seq2",
"P",
"Q",
"knows",
0.5,
"gamma",
1_700_000_301_000_000,
Some(3_500_000),
)
.unwrap();
assert_eq!(db.visible_seq_for("gamma"), 3_500_000);
}
#[test]
fn delete_entity_edge_with_id_bumps_visible_seq_even_when_edge_missing() {
let db = YantrikDB::new(":memory:", 64).unwrap();
db.delete_entity_edge_with_id(
"edge_never",
"delta",
1_700_000_400_000_000,
Some(4_000_000),
)
.unwrap();
assert_eq!(db.visible_seq_for("delta"), 4_000_000);
}
#[test]
fn issue_8_tombstoned_memories_excluded_from_recall() {
let db = YantrikDB::new(":memory:", 64).unwrap();
let emb = vec_seed(42.0, 64);
let rid = db
.record(
"memory to forget for issue 8 repro",
"episodic",
0.5,
0.0,
604800.0,
&empty_meta(),
&emb,
"default",
0.8,
"general",
"user",
None,
)
.unwrap();
let before = db
.recall(
&emb, 5, None, None, false, false, None, true, None, None, None, None, None,
)
.unwrap();
assert!(
before.iter().any(|r| r.rid == rid),
"memory must be findable before forget"
);
db.forget(&rid).unwrap();
let after = db
.recall(
&emb, 5, None, None, false, false, None, true, None, None, None, None, None,
)
.unwrap();
assert!(
!after.iter().any(|r| r.rid == rid),
"issue #8: tombstoned memory must NOT appear in recall results"
);
}
#[test]
fn issue_8_tombstoned_persists_across_engine_reopen() {
let dir = tempfile::tempdir().unwrap();
let path = dir.path().join("issue8.db");
let path_str = path.to_str().unwrap();
let emb = vec_seed(44.0, 64);
let rid;
{
let db = YantrikDB::new(path_str, 64).unwrap();
rid = db
.record(
"memory survives reopen test",
"episodic",
0.5,
0.0,
604800.0,
&empty_meta(),
&emb,
"default",
0.8,
"general",
"user",
None,
)
.unwrap();
db.forget(&rid).unwrap();
}
{
let db2 = YantrikDB::new(path_str, 64).unwrap();
let after = db2
.recall(
&emb, 5, None, None, false, false, None, true, None, None, None, None, None,
)
.unwrap();
assert!(
!after.iter().any(|r| r.rid == rid),
"tombstoned memory must stay hidden across engine reopen"
);
}
}
#[test]
fn visible_seq_starts_at_zero_for_new_namespace() {
let db = YantrikDB::new(":memory:", 64).unwrap();
assert_eq!(db.visible_seq_for("never_used"), 0);
}
#[test]
fn record_bumps_visible_seq_for_namespace() {
let db = YantrikDB::new(":memory:", 64).unwrap();
let before = db.visible_seq_for("default");
let _ = db
.record(
"test",
"episodic",
0.5,
0.0,
604800.0,
&empty_meta(),
&vec_seed(1.0, 64),
"default",
0.8,
"general",
"user",
None,
)
.unwrap();
let after = db.visible_seq_for("default");
assert!(after > before, "record() must bump visible_seq[default]");
}
#[test]
fn visible_seq_isolated_per_namespace() {
let db = YantrikDB::new(":memory:", 64).unwrap();
db.record(
"ns_a memory",
"episodic",
0.5,
0.0,
604800.0,
&empty_meta(),
&vec_seed(1.0, 64),
"ns_a",
0.8,
"general",
"user",
None,
)
.unwrap();
let seq_a = db.visible_seq_for("ns_a");
let seq_b = db.visible_seq_for("ns_b");
assert!(seq_a > 0);
assert_eq!(seq_b, 0, "ns_b unaffected by writes to ns_a");
}
#[test]
fn wait_for_visible_seq_succeeds_when_already_reached() {
let db = YantrikDB::new(":memory:", 64).unwrap();
db.record(
"set watermark",
"episodic",
0.5,
0.0,
604800.0,
&empty_meta(),
&vec_seed(1.0, 64),
"default",
0.8,
"general",
"user",
None,
)
.unwrap();
let current = db.visible_seq_for("default");
db.wait_for_visible_seq("default", current, std::time::Duration::from_millis(100))
.expect("already-reached watermark");
}
#[test]
fn wait_for_visible_seq_times_out_on_unreachable() {
let db = YantrikDB::new(":memory:", 64).unwrap();
let err = db
.wait_for_visible_seq("never", 9999, std::time::Duration::from_millis(50))
.expect_err("must timeout");
match err {
crate::error::YantrikDbError::RyWaitTimeout {
namespace,
requested_seq,
observed_seq,
waited_ms,
} => {
assert_eq!(namespace, "never");
assert_eq!(requested_seq, 9999);
assert_eq!(observed_seq, 0);
assert_eq!(waited_ms, 50);
}
other => panic!("expected RyWaitTimeout, got {other:?}"),
}
}
#[test]
fn wait_for_visible_seq_wakes_on_concurrent_write() {
use std::sync::Arc;
use std::thread;
use std::time::Duration;
let db = Arc::new(YantrikDB::new(":memory:", 64).unwrap());
let db_w = Arc::clone(&db);
let waiter =
thread::spawn(move || db_w.wait_for_visible_seq("default", 1, Duration::from_secs(2)));
let db_writer = Arc::clone(&db);
let writer = thread::spawn(move || {
thread::sleep(Duration::from_millis(50));
db_writer
.record(
"wake the waiter",
"episodic",
0.5,
0.0,
604800.0,
&empty_meta(),
&vec_seed(1.0, 64),
"default",
0.8,
"general",
"user",
None,
)
.unwrap();
});
writer.join().unwrap();
let result = waiter.join().unwrap();
assert!(result.is_ok(), "waiter should be notified by the write");
}
#[test]
fn recall_with_seq_returns_results_when_seq_reached() {
let db = YantrikDB::new(":memory:", 64).unwrap();
let emb = vec_seed(1.0, 64);
let _ = db
.record(
"ryw test",
"episodic",
0.5,
0.0,
604800.0,
&empty_meta(),
&emb,
"default",
0.8,
"general",
"user",
None,
)
.unwrap();
let current = db.visible_seq_for("default");
let r = db
.recall_with_seq(
&emb,
5,
None,
None,
false,
false,
None,
true,
Some("default"),
None,
None,
current,
std::time::Duration::from_millis(100),
)
.unwrap();
assert!(
!r.is_empty(),
"recall_with_seq returns results once seq reached"
);
}
#[test]
fn recall_with_seq_times_out_on_unreachable() {
let db = YantrikDB::new(":memory:", 64).unwrap();
let err = db
.recall_with_seq(
&vec_seed(1.0, 64),
5,
None,
None,
false,
false,
None,
true,
Some("default"),
None,
None,
9999,
std::time::Duration::from_millis(50),
)
.expect_err("must timeout");
assert!(matches!(
err,
crate::error::YantrikDbError::RyWaitTimeout { .. }
));
}
#[cfg(feature = "bundled-embedder")]
#[test]
fn bundled_embedder_auto_attaches_on_default_dim() {
use crate::embedder::BUNDLED_EMBEDDER_DIM;
let db = YantrikDB::new(":memory:", BUNDLED_EMBEDDER_DIM).unwrap();
assert!(
db.has_embedder(),
"default-build YantrikDB::new with bundled dim must auto-attach BundledEmbedder"
);
}
#[cfg(feature = "bundled-embedder")]
#[test]
fn with_default_constructor_attaches_bundled_embedder() {
let db = YantrikDB::with_default(":memory:").unwrap();
assert!(
db.has_embedder(),
"with_default must auto-attach BundledEmbedder"
);
}
#[cfg(feature = "bundled-embedder")]
#[test]
fn bundled_embedder_does_not_attach_on_mismatched_dim() {
let db = YantrikDB::new(":memory:", 384).unwrap();
assert!(
!db.has_embedder(),
"dim mismatch should NOT auto-attach (avoids silent corruption)"
);
}
#[cfg(feature = "bundled-embedder")]
#[test]
fn bundled_embedder_record_text_round_trip() {
let db = YantrikDB::with_default(":memory:").unwrap();
let _rid = db
.record_text(
"Alice met Acme yesterday",
"episodic",
0.5,
0.0,
604800.0,
&empty_meta(),
"default",
0.8,
"general",
"user",
None,
)
.expect("record_text should work without explicit set_embedder");
let results = db.recall_text("Alice", 5).expect("recall_text should work");
assert!(!results.is_empty(), "recall finds the recorded memory");
assert!(
results[0].text.contains("Alice"),
"potion-2M finds the recorded memory; got: {:?}",
results[0].text
);
}
#[cfg(feature = "bundled-embedder")]
#[test]
fn explicit_set_embedder_overrides_bundled() {
struct DummyEmbedder;
impl crate::types::Embedder for DummyEmbedder {
fn embed(
&self,
_t: &str,
) -> std::result::Result<Vec<f32>, Box<dyn std::error::Error + Send + Sync>> {
let mut v = vec![0.0; 64];
v[0] = 0.7777;
Ok(v)
}
fn dim(&self) -> usize {
64
}
}
let mut db = YantrikDB::with_default(":memory:").unwrap();
assert!(db.has_embedder(), "starts with bundled");
db.set_embedder(Box::new(DummyEmbedder)).unwrap();
let v = db.embed("anything").unwrap();
assert!(
(v[0] - 0.7777).abs() < 1e-6,
"DummyEmbedder's sentinel must be visible — set_embedder overrode bundled"
);
}
#[test]
fn migration_replay_does_not_trip_on_already_present_column() {
use tempfile::NamedTempFile;
let tmp = NamedTempFile::new().unwrap();
let path = tmp.path().to_str().unwrap();
{
let _db = YantrikDB::new(path, 8).unwrap();
}
{
let conn = rusqlite::Connection::open(path).unwrap();
conn.execute(
"INSERT OR REPLACE INTO meta (key, value) VALUES ('schema_version', '23')",
[],
)
.unwrap();
}
let db = YantrikDB::new(path, 8)
.expect("v0.7.3 idempotent migration runner must heal rewound-meta deployments");
db.record(
"post-heal smoke",
"episodic",
0.5,
0.0,
604800.0,
&empty_meta(),
&vec_seed(1.0, 8),
"default",
0.8,
"general",
"user",
None,
)
.unwrap();
}
#[test]
fn migration_replay_does_not_trip_on_alter_table_against_view() {
use tempfile::NamedTempFile;
let tmp = NamedTempFile::new().unwrap();
let path = tmp.path().to_str().unwrap();
{
let _db = YantrikDB::new(path, 8).unwrap();
}
{
let conn = rusqlite::Connection::open(path).unwrap();
let kind: String = conn
.query_row(
"SELECT type FROM sqlite_master WHERE name = 'edges'",
[],
|row| row.get(0),
)
.unwrap();
assert_eq!(
kind, "view",
"fixture precondition: at current schema, edges should be a backward-compat view"
);
}
{
let conn = rusqlite::Connection::open(path).unwrap();
conn.execute(
"INSERT OR REPLACE INTO meta (key, value) VALUES ('schema_version', '14')",
[],
)
.unwrap();
}
let db = YantrikDB::new(path, 8)
.expect("v0.7.8 idempotent runner must heal rewound-meta DBs that hit ALTER-on-view");
db.record(
"post-heal smoke (issue 10)",
"episodic",
0.5,
0.0,
604800.0,
&empty_meta(),
&vec_seed(1.0, 8),
"default",
0.8,
"general",
"user",
None,
)
.unwrap();
}
#[test]
fn migration_meta_stamp_does_not_downgrade() {
use tempfile::NamedTempFile;
let tmp = NamedTempFile::new().unwrap();
let path = tmp.path().to_str().unwrap();
{
let conn = rusqlite::Connection::open(path).unwrap();
conn.execute_batch(
"CREATE TABLE IF NOT EXISTS meta (key TEXT PRIMARY KEY, value TEXT NOT NULL);
INSERT OR REPLACE INTO meta (key, value) VALUES ('schema_version', '999');",
)
.unwrap();
}
let _db = YantrikDB::new(path, 8).unwrap();
let stamped: String = {
let conn = rusqlite::Connection::open(path).unwrap();
conn.query_row(
"SELECT value FROM meta WHERE key = 'schema_version'",
[],
|row| row.get(0),
)
.unwrap()
};
assert_eq!(
stamped, "999",
"MAX-stamp invariant: open() must never rewind meta.schema_version below the on-disk value"
);
}
fn table_columns(conn: &rusqlite::Connection, table: &str) -> Vec<String> {
let mut stmt = conn
.prepare(&format!("PRAGMA table_info({table})"))
.unwrap();
stmt.query_map([], |row| row.get::<_, String>(1))
.unwrap()
.collect::<std::result::Result<Vec<_>, _>>()
.unwrap()
}
fn index_exists(conn: &rusqlite::Connection, index: &str) -> bool {
let count: i64 = conn
.query_row(
"SELECT COUNT(*) FROM sqlite_master WHERE type='index' AND name=?1",
params![index],
|row| row.get(0),
)
.unwrap();
count == 1
}
#[test]
fn schema_v26_fresh_install_has_provenance_columns_and_indexes() {
let db = YantrikDB::new(":memory:", 8).unwrap();
let conn = db.conn();
let cols = table_columns(&conn, "memories");
for required in [
"prior_rid",
"resolution_kind",
"dismissal_reason",
"confidence_at_write",
] {
assert!(
cols.iter().any(|c| c == required),
"v26: fresh-install memories table missing column {required}, got: {cols:?}"
);
}
assert!(
index_exists(&conn, "idx_memories_prior_rid"),
"v26: fresh-install missing partial index idx_memories_prior_rid"
);
assert!(
index_exists(&conn, "idx_memories_resolution_kind"),
"v26: fresh-install missing partial index idx_memories_resolution_kind"
);
}
#[test]
fn schema_v26_migration_from_v25_adds_columns_and_normalizes_source() {
use tempfile::NamedTempFile;
let tmp = NamedTempFile::new().unwrap();
let path = tmp.path().to_str().unwrap();
{
let db = YantrikDB::new(path, 8).unwrap();
let conn = db.conn();
for (rid, src) in [
("01900000-0000-7000-8000-000000000001", "user"),
("01900000-0000-7000-8000-000000000002", "inference"),
("01900000-0000-7000-8000-000000000003", "legacy-freetext"),
] {
conn.execute(
"INSERT INTO memories (rid, type, text, created_at, updated_at, last_access, source) \
VALUES (?1, 'episodic', 'test', 0.0, 0.0, 0.0, ?2)",
params![rid, src],
)
.unwrap();
}
}
{
let conn = rusqlite::Connection::open(path).unwrap();
conn.execute(
"INSERT OR REPLACE INTO meta (key, value) VALUES ('schema_version', '25')",
[],
)
.unwrap();
conn.execute("DROP INDEX IF EXISTS idx_memories_prior_rid", [])
.unwrap();
conn.execute("DROP INDEX IF EXISTS idx_memories_resolution_kind", [])
.unwrap();
}
let db = YantrikDB::new(path, 8)
.expect("v26 migration must run cleanly against a rewound-meta v25 DB");
let conn = db.conn();
let cols = table_columns(&conn, "memories");
for required in [
"prior_rid",
"resolution_kind",
"dismissal_reason",
"confidence_at_write",
] {
assert!(
cols.iter().any(|c| c == required),
"v26 migration: missing column {required} after re-open"
);
}
assert!(
index_exists(&conn, "idx_memories_prior_rid"),
"v26 migration: missing partial index idx_memories_prior_rid"
);
assert!(
index_exists(&conn, "idx_memories_resolution_kind"),
"v26 migration: missing partial index idx_memories_resolution_kind"
);
let normalized: String = conn
.query_row(
"SELECT source FROM memories WHERE rid = '01900000-0000-7000-8000-000000000003'",
[],
|row| row.get(0),
)
.unwrap();
assert_eq!(
normalized, "user",
"v26 migration must coerce non-enum source to 'user'"
);
let preserved: String = conn
.query_row(
"SELECT source FROM memories WHERE rid = '01900000-0000-7000-8000-000000000002'",
[],
|row| row.get(0),
)
.unwrap();
assert_eq!(
preserved, "inference",
"v26 migration must preserve enum-valid source values"
);
let log: String = conn
.query_row(
"SELECT value FROM meta WHERE key = 'source_normalization_log_v26'",
[],
|row| row.get(0),
)
.unwrap();
assert!(
log.contains("normalized 1 rows"),
"v26 migration must log normalization count, got: {log}"
);
}
#[test]
fn schema_v26_migration_replay_is_idempotent() {
use tempfile::NamedTempFile;
let tmp = NamedTempFile::new().unwrap();
let path = tmp.path().to_str().unwrap();
{
let _db = YantrikDB::new(path, 8).unwrap();
}
{
let conn = rusqlite::Connection::open(path).unwrap();
conn.execute(
"INSERT OR REPLACE INTO meta (key, value) VALUES ('schema_version', '25')",
[],
)
.unwrap();
}
let db = YantrikDB::new(path, 8)
.expect("v26 migration runner must heal rewound-meta deployments on a v26-schema DB");
db.record(
"post-v26-heal smoke",
"episodic",
0.5,
0.0,
604800.0,
&empty_meta(),
&vec_seed(1.0, 8),
"default",
0.8,
"general",
"user",
None,
)
.unwrap();
}
#[test]
fn record_in_normal_state_takes_sync_path() {
let db = YantrikDB::new(":memory:", 8).unwrap();
assert_eq!(
db.write_router.state(),
crate::engine::write_router::RouterState::Normal
);
let rid = db
.record(
"sync path test",
"episodic",
0.5,
0.0,
86400.0,
&empty_meta(),
&vec_seed(1.0, 8),
"default",
0.9,
"general",
"user",
None,
)
.unwrap();
let count: i64 = db
.conn()
.query_row(
"SELECT COUNT(*) FROM memories WHERE rid = ?1",
params![rid],
|row| row.get(0),
)
.unwrap();
assert_eq!(
count, 1,
"sync-path write must land in memories immediately"
);
}
#[test]
fn record_in_queueing_state_routes_to_oplog_does_not_touch_memories() {
let db = YantrikDB::new(":memory:", 8).unwrap();
db.write_router.switch_to_queueing();
assert_eq!(
db.write_router.state(),
crate::engine::write_router::RouterState::Queueing
);
let mem_before: i64 = db
.conn()
.query_row("SELECT COUNT(*) FROM memories", [], |row| row.get(0))
.unwrap();
let oplog_before: i64 = db
.conn()
.query_row(
"SELECT COUNT(*) FROM oplog WHERE applied = 0 AND op_type = 'record'",
[],
|row| row.get(0),
)
.unwrap();
let rid = db
.record(
"queued path test",
"episodic",
0.5,
0.0,
86400.0,
&empty_meta(),
&vec_seed(1.0, 8),
"default",
0.9,
"general",
"user",
None,
)
.unwrap();
let mem_after: i64 = db
.conn()
.query_row("SELECT COUNT(*) FROM memories", [], |row| row.get(0))
.unwrap();
assert_eq!(
mem_after, mem_before,
"queued path must NOT write to memories table during reembed cutover \
(brainstorm-2/3 invariant 1: queued-after-barrier writes are replayed \
by post-swap materializer, not committed to old generation)"
);
let oplog_after: i64 = db
.conn()
.query_row(
"SELECT COUNT(*) FROM oplog WHERE applied = 0 AND op_type = 'record' \
AND target_rid = ?1",
params![rid],
|row| row.get(0),
)
.unwrap();
assert_eq!(
oplog_after,
oplog_before + 1,
"queued path must write the record op to oplog with applied=0"
);
let applied_gen: Option<i64> = db
.conn()
.query_row(
"SELECT applied_generation FROM oplog WHERE target_rid = ?1",
params![rid],
|row| row.get(0),
)
.unwrap();
assert!(
applied_gen.is_none(),
"applied_generation must be NULL for queued ops; got Some({applied_gen:?})"
);
db.write_router.switch_to_normal();
}
#[test]
fn record_guard_drops_inflight_counter_panic_safe_via_raii() {
let db = std::sync::Arc::new(YantrikDB::new(":memory:", 8).unwrap());
assert_eq!(db.write_router.inflight(), 0);
let db_panic = std::sync::Arc::clone(&db);
let result = std::panic::catch_unwind(std::panic::AssertUnwindSafe(|| {
let _guard = db_panic
.write_router
.try_enter_sync_writer()
.expect("Normal state must yield guard");
assert_eq!(db_panic.write_router.inflight(), 1);
panic!("simulated mid-write panic");
}));
assert!(result.is_err(), "panic must propagate up");
assert_eq!(
db.write_router.inflight(),
0,
"panic-safe inflight counter (RAII Drop) — required for reembed cutover correctness"
);
}
mod mode_test_embedders {
use crate::types::Embedder;
pub struct FakeEmbedder {
pub dim: usize,
pub fp: Option<String>,
pub name: Option<String>,
pub sentinel: f32,
}
impl Embedder for FakeEmbedder {
fn embed(
&self,
_text: &str,
) -> std::result::Result<Vec<f32>, Box<dyn std::error::Error + Send + Sync>> {
let mut v = vec![0.0_f32; self.dim];
if !v.is_empty() {
v[0] = self.sentinel;
}
Ok(v)
}
fn dim(&self) -> usize {
self.dim
}
fn fingerprint(&self) -> Option<String> {
self.fp.clone()
}
fn name(&self) -> Option<String> {
self.name.clone()
}
}
}
#[test]
fn set_embedder_test_1_same_dim_different_digest_on_populated_db_rejected() {
use mode_test_embedders::FakeEmbedder;
let mut db = YantrikDB::new(":memory:", 64).unwrap();
db.set_embedder(Box::new(FakeEmbedder {
dim: 64,
fp: Some("sha256:embedder_A".to_string()),
name: Some("embedder_A".to_string()),
sentinel: 1.0,
}))
.unwrap();
let _ = db
.record(
"first memory",
"semantic",
0.5,
0.0,
86400.0,
&empty_meta(),
&vec_seed(1.0, 64),
"default",
0.9,
"general",
"user",
None,
)
.unwrap();
let err = db
.set_embedder(Box::new(FakeEmbedder {
dim: 64,
fp: Some("sha256:embedder_B".to_string()),
name: Some("embedder_B".to_string()),
sentinel: 2.0,
}))
.unwrap_err();
assert!(
matches!(
err,
crate::error::YantrikDbError::ChangeEmbedderDigestRequiresReembed { .. }
),
"same-dim-different-digest on Known-provenance populated DB must \
return ChangeEmbedderDigestRequiresReembed (silent-corruption \
prevention invariant from brainstorm-3); got {err:?}"
);
}
#[test]
fn set_embedder_test_2_different_dim_on_populated_db_rejected() {
use mode_test_embedders::FakeEmbedder;
let mut db = YantrikDB::new(":memory:", 64).unwrap();
db.set_embedder(Box::new(FakeEmbedder {
dim: 64,
fp: Some("sha256:fp64".to_string()),
name: None,
sentinel: 1.0,
}))
.unwrap();
let _ = db
.record(
"m",
"semantic",
0.5,
0.0,
86400.0,
&empty_meta(),
&vec_seed(1.0, 64),
"default",
0.9,
"general",
"user",
None,
)
.unwrap();
let err = db
.set_embedder(Box::new(FakeEmbedder {
dim: 128,
fp: Some("sha256:fp128".to_string()),
name: None,
sentinel: 2.0,
}))
.unwrap_err();
assert!(
matches!(
err,
crate::error::YantrikDbError::ChangeEmbedderDimensionRequiresReembed { .. }
),
"dim change on populated DB must return \
ChangeEmbedderDimensionRequiresReembed; got {err:?}"
);
}
#[test]
fn set_embedder_test_3_empty_db_with_fingerprint_upgrades_provenance_to_known() {
use mode_test_embedders::FakeEmbedder;
let mut db = YantrikDB::new(":memory:", 64).unwrap();
assert!(matches!(
db.search_state.load().index_embedding,
crate::engine::reembed::EmbeddingProvenance::ExternalOrUnknown { .. }
));
db.set_embedder(Box::new(FakeEmbedder {
dim: 64,
fp: Some("sha256:initial".to_string()),
name: Some("initial".to_string()),
sentinel: 1.0,
}))
.unwrap();
let s = db.search_state.load_full();
match &s.index_embedding {
crate::engine::reembed::EmbeddingProvenance::Known { name, digest, dim } => {
assert_eq!(name.as_deref(), Some("initial"));
assert_eq!(digest, "sha256:initial");
assert_eq!(*dim, 64);
}
other => panic!("expected Known provenance after attach on empty DB, got {other:?}"),
}
}
#[test]
fn set_embedder_test_4_empty_db_no_fingerprint_stays_external_or_unknown() {
use mode_test_embedders::FakeEmbedder;
let mut db = YantrikDB::new(":memory:", 64).unwrap();
db.set_embedder(Box::new(FakeEmbedder {
dim: 64,
fp: None,
name: None,
sentinel: 1.0,
}))
.unwrap();
let s = db.search_state.load_full();
assert!(
matches!(
s.index_embedding,
crate::engine::reembed::EmbeddingProvenance::ExternalOrUnknown { dim: 64 }
),
"no-fingerprint embedder on empty DB must keep ExternalOrUnknown provenance, \
got {:?}",
s.index_embedding
);
assert!(s.has_runtime_embedder());
}
#[test]
fn set_embedder_test_5_same_digest_replacement_does_not_bump_generation() {
use mode_test_embedders::FakeEmbedder;
let mut db = YantrikDB::new(":memory:", 64).unwrap();
db.set_embedder(Box::new(FakeEmbedder {
dim: 64,
fp: Some("sha256:same".to_string()),
name: Some("same".to_string()),
sentinel: 1.0,
}))
.unwrap();
let gen_before = db.search_state.load().generation;
db.set_embedder(Box::new(FakeEmbedder {
dim: 64,
fp: Some("sha256:same".to_string()),
name: Some("same".to_string()),
sentinel: 2.0,
}))
.unwrap();
let gen_after = db.search_state.load().generation;
assert_eq!(
gen_after, gen_before,
"same-digest replacement must NOT bump generation (no coherent-bundle change)"
);
let v = db.embed("anything").unwrap();
assert!(
(v[0] - 2.0).abs() < 1e-6,
"runtime Arc must have been replaced; expected sentinel 2.0, got {}",
v[0]
);
}
#[test]
fn set_embedder_test_6_external_or_unknown_compat_attach_does_not_claim_provenance() {
use mode_test_embedders::FakeEmbedder;
let mut db = YantrikDB::new(":memory:", 64).unwrap();
let _ = db
.record(
"external vec",
"semantic",
0.5,
0.0,
86400.0,
&empty_meta(),
&vec_seed(1.0, 64),
"default",
0.9,
"general",
"user",
None,
)
.unwrap();
assert!(matches!(
db.search_state.load().index_embedding,
crate::engine::reembed::EmbeddingProvenance::ExternalOrUnknown { .. }
));
db.set_embedder(Box::new(FakeEmbedder {
dim: 64,
fp: Some("sha256:attached".to_string()),
name: Some("attached".to_string()),
sentinel: 1.0,
}))
.unwrap();
let s = db.search_state.load_full();
assert!(
matches!(
s.index_embedding,
crate::engine::reembed::EmbeddingProvenance::ExternalOrUnknown { .. }
),
"compat-attach must NOT upgrade ExternalOrUnknown provenance to Known; \
existing vectors weren't built with this embedder. Got {:?}",
s.index_embedding
);
assert!(s.has_runtime_embedder());
assert_eq!(
s.runtime_embedder_digest.as_deref(),
Some("sha256:attached")
);
}
#[test]
fn set_embedder_test_7_has_embedder_derives_from_search_state() {
let mut db = YantrikDB::new(":memory:", 384).unwrap();
assert!(!db.has_embedder(), "fresh engine: no embedder");
use mode_test_embedders::FakeEmbedder;
db.set_embedder(Box::new(FakeEmbedder {
dim: 384,
fp: Some("sha256:x".to_string()),
name: None,
sentinel: 0.42,
}))
.unwrap();
assert!(db.has_embedder());
let v = db.embed("anything").unwrap();
assert!((v[0] - 0.42).abs() < 1e-6);
}
#[test]
fn record_text_revalidates_generation_and_retries_after_swap() {
use std::sync::mpsc::channel;
use std::sync::{Arc, Mutex};
use std::time::Duration;
let (started_tx, started_rx) = channel::<()>();
let (release_tx, release_rx) = channel::<()>();
let call_count = Arc::new(std::sync::atomic::AtomicUsize::new(0));
struct SharedBlocking {
dim: usize,
fp: Option<String>,
name: Option<String>,
sentinel: f32,
started_tx: Mutex<Option<std::sync::mpsc::Sender<()>>>,
release_rx: Mutex<Option<std::sync::mpsc::Receiver<()>>>,
call_count: Arc<std::sync::atomic::AtomicUsize>,
}
impl crate::types::Embedder for SharedBlocking {
fn embed(
&self,
_text: &str,
) -> std::result::Result<Vec<f32>, Box<dyn std::error::Error + Send + Sync>> {
let n = self
.call_count
.fetch_add(1, std::sync::atomic::Ordering::SeqCst);
if n == 0 {
if let Some(tx) = self.started_tx.lock().unwrap().take() {
let _ = tx.send(());
}
if let Some(rx) = self.release_rx.lock().unwrap().take() {
let _ = rx.recv();
}
}
let mut v = vec![0.0_f32; self.dim];
if !v.is_empty() {
v[0] = self.sentinel;
}
Ok(v)
}
fn dim(&self) -> usize {
self.dim
}
fn fingerprint(&self) -> Option<String> {
self.fp.clone()
}
fn name(&self) -> Option<String> {
self.name.clone()
}
}
let mut db = YantrikDB::new(":memory:", 64).unwrap();
db.set_embedder(Box::new(SharedBlocking {
dim: 64,
fp: Some("sha256:initial".to_string()),
name: Some("blocking-initial".to_string()),
sentinel: 0.42,
started_tx: Mutex::new(Some(started_tx)),
release_rx: Mutex::new(Some(release_rx)),
call_count: Arc::clone(&call_count),
}))
.unwrap();
let gen_before = db.search_state.load().generation;
let arc_db = Arc::new(db);
let worker_db = Arc::clone(&arc_db);
let worker = std::thread::spawn(move || {
worker_db
.record_text(
"hello",
"episodic",
0.5,
0.0,
604800.0,
&serde_json::json!({}),
"default",
0.8,
"general",
"user",
None,
)
.unwrap()
});
started_rx
.recv_timeout(Duration::from_secs(5))
.expect("embed should have started within 5s");
let old_state = arc_db.search_state.load_full();
let new_state = crate::engine::reembed::SearchState {
index_embedding: crate::engine::reembed::EmbeddingProvenance::Known {
name: Some("blocking-rotated".to_string()),
digest: "sha256:rotated".to_string(),
dim: 64,
},
embedder: old_state.embedder.clone(),
runtime_embedder_name: Some("blocking-rotated".to_string()),
runtime_embedder_digest: Some("sha256:rotated".to_string()),
generation: old_state.generation + 1,
covers_through_seq: old_state.covers_through_seq,
hnsw_m: old_state.hnsw_m,
hnsw_ef_construction: old_state.hnsw_ef_construction,
hnsw_ef_search: old_state.hnsw_ef_search,
vec_index: Arc::clone(&old_state.vec_index),
};
arc_db.search_state.store(Arc::new(new_state));
release_tx.send(()).unwrap();
let rid = worker.join().expect("record_text must complete");
let n_calls = call_count.load(std::sync::atomic::Ordering::SeqCst);
assert!(
n_calls >= 2,
"record_text must re-embed after SearchState swap; got {n_calls} calls"
);
assert!(!rid.is_empty(), "record_text returns a valid rid");
let gen_after = arc_db.search_state.load().generation;
assert!(
gen_after > gen_before,
"test must observe a generation advance: before={gen_before} after={gen_after}"
);
let conn = arc_db.conn();
let count: i64 = conn
.query_row(
"SELECT COUNT(*) FROM memories WHERE rid = ?1",
[&rid],
|r| r.get(0),
)
.unwrap();
assert_eq!(count, 1, "the retried record_text must be durably stored");
}
#[test]
fn record_text_routes_to_queued_when_router_is_queueing() {
use mode_test_embedders::FakeEmbedder;
let mut db = YantrikDB::new(":memory:", 64).unwrap();
db.set_embedder(Box::new(FakeEmbedder {
dim: 64,
fp: Some("sha256:initial".to_string()),
name: Some("initial".to_string()),
sentinel: 0.5,
}))
.unwrap();
db.write_router.switch_to_queueing();
let rid = db
.record_text(
"hello-queued",
"episodic",
0.5,
0.0,
604800.0,
&serde_json::json!({}),
"default",
0.8,
"general",
"user",
None,
)
.unwrap();
assert!(!rid.is_empty());
let conn = db.conn();
let memories_count: i64 = conn
.query_row(
"SELECT COUNT(*) FROM memories WHERE rid = ?1",
[&rid],
|r| r.get(0),
)
.unwrap();
assert_eq!(
memories_count, 0,
"queued path must NOT write to memories table"
);
let oplog_row: (i64, Option<String>) = conn
.query_row(
"SELECT applied, embedding_model FROM oplog WHERE target_rid = ?1 AND op_type = 'record'",
[&rid],
|r| Ok((r.get(0)?, r.get(1)?)),
)
.unwrap();
assert_eq!(oplog_row.0, 0, "queued op must be applied=0");
assert_eq!(
oplog_row.1.as_deref(),
Some("initial"),
"queued op carries the current runtime embedder name for post-swap re-encode"
);
}
#[test]
fn log_op_stamps_applied_generation_from_active_search_state() {
let db = YantrikDB::new(":memory:", 64).unwrap();
let initial_generation: i64 = db.search_state.load().generation as i64;
let op_id = db
.log_op("test_event", None, &serde_json::json!({"x": 1}), None)
.unwrap();
let applied_generation: Option<i64> = {
let conn = db.conn();
conn.query_row(
"SELECT applied_generation FROM oplog WHERE op_id = ?1",
[&op_id],
|r| r.get(0),
)
.unwrap()
};
assert_eq!(
applied_generation,
Some(initial_generation),
"log_op must stamp applied_generation with the current SearchState generation"
);
let old_state = db.search_state.load_full();
let bumped = crate::engine::reembed::SearchState {
index_embedding: old_state.index_embedding.clone(),
embedder: old_state.embedder.clone(),
runtime_embedder_name: old_state.runtime_embedder_name.clone(),
runtime_embedder_digest: old_state.runtime_embedder_digest.clone(),
generation: old_state.generation + 1,
covers_through_seq: old_state.covers_through_seq,
hnsw_m: old_state.hnsw_m,
hnsw_ef_construction: old_state.hnsw_ef_construction,
hnsw_ef_search: old_state.hnsw_ef_search,
vec_index: std::sync::Arc::clone(&old_state.vec_index),
};
db.search_state.store(std::sync::Arc::new(bumped));
let op_id2 = db
.log_op(
"test_event_after_bump",
None,
&serde_json::json!({"x": 2}),
None,
)
.unwrap();
let applied_generation2: Option<i64> = {
let conn = db.conn();
conn.query_row(
"SELECT applied_generation FROM oplog WHERE op_id = ?1",
[&op_id2],
|r| r.get(0),
)
.unwrap()
};
assert_eq!(
applied_generation2,
Some(initial_generation + 1),
"log_op picks up the new generation after search_state.store"
);
}
#[test]
fn set_embedder_test_8_atomic_publication_no_partial_state() {
use mode_test_embedders::FakeEmbedder;
use std::sync::Arc;
let mut db = YantrikDB::new(":memory:", 64).unwrap();
db.set_embedder(Box::new(FakeEmbedder {
dim: 64,
fp: Some("sha256:initial".to_string()),
name: None,
sentinel: 1.0,
}))
.unwrap();
let state = db.search_state.load_full();
assert_eq!(
state.embedder.is_some(),
state.runtime_embedder_digest.is_some(),
"embedder Some <=> digest Some must hold for any consistent snapshot"
);
for sentinel in [2.0_f32, 3.0, 4.0, 5.0] {
db.set_embedder(Box::new(FakeEmbedder {
dim: 64,
fp: Some("sha256:initial".to_string()),
name: None,
sentinel,
}))
.unwrap();
let s = db.search_state.load_full();
assert_eq!(
s.embedder.is_some(),
s.runtime_embedder_digest.is_some(),
"consistency must hold across replacements (no partial state)"
);
let _arc_held: Arc<crate::engine::reembed::SearchState> = s;
}
}
#[test]
fn search_state_initial_on_fresh_engine() {
let db = YantrikDB::new(":memory:", 384).unwrap();
let state = db.search_state.load_full();
assert_eq!(
state.dim(),
384,
"initial dim must match constructor parameter"
);
assert!(matches!(
state.index_embedding,
crate::engine::reembed::EmbeddingProvenance::ExternalOrUnknown { dim: 384 }
));
assert!(
!state.has_runtime_embedder(),
"fresh engine must have no runtime embedder until set_embedder*"
);
assert_eq!(state.generation, 0);
assert_eq!(state.covers_through_seq, 0);
}
#[test]
fn schema_v27_fresh_install_has_reembed_surfaces() {
let db = YantrikDB::new(":memory:", 8).unwrap();
let conn = db.conn();
let memories_cols = table_columns(&conn, "memories");
for required in ["embedding_new", "embedding_new_model"] {
assert!(
memories_cols.iter().any(|c| c == required),
"v27: fresh-install memories table missing column {required}, got: {memories_cols:?}"
);
}
let oplog_cols = table_columns(&conn, "oplog");
assert!(
oplog_cols.iter().any(|c| c == "embedding_model"),
"v27: fresh-install oplog table missing column embedding_model, got: {oplog_cols:?}"
);
assert!(
oplog_cols.iter().any(|c| c == "applied_generation"),
"v27: fresh-install oplog table missing column applied_generation \
(brainstorm-2 correction \u{2014} per-generation application tracking \
replaces boolean `applied` as truth), got: {oplog_cols:?}"
);
let events_cols = table_columns(&conn, "reembed_events");
for required in ["generation", "phase", "timestamp", "payload_json"] {
assert!(
events_cols.iter().any(|c| c == required),
"v27: fresh-install reembed_events missing column {required}, got: {events_cols:?}"
);
}
for required_idx in [
"idx_reembed_events_generation",
"idx_oplog_applied_generation",
] {
assert!(
index_exists(&conn, required_idx),
"v27: fresh-install missing index {required_idx}"
);
}
}
#[test]
fn schema_v27_migration_from_v26_is_additive_only() {
use tempfile::NamedTempFile;
let tmp = NamedTempFile::new().unwrap();
let path = tmp.path().to_str().unwrap();
let planted_rid = "01900000-0000-7000-8000-00000000c027";
{
let db = YantrikDB::new(path, 8).unwrap();
let conn = db.conn();
conn.execute(
"INSERT INTO memories (rid, type, text, embedding, created_at, updated_at, last_access, source) \
VALUES (?1, 'episodic', 'planted under v27 schema', X'01020304', 0.0, 0.0, 0.0, 'user')",
params![planted_rid],
)
.unwrap();
}
{
let conn = rusqlite::Connection::open(path).unwrap();
conn.execute(
"INSERT OR REPLACE INTO meta (key, value) VALUES ('schema_version', '26')",
[],
)
.unwrap();
conn.execute("DROP TABLE IF EXISTS reembed_events", [])
.unwrap();
conn.execute("DROP INDEX IF EXISTS idx_reembed_events_generation", [])
.unwrap();
}
let db = YantrikDB::new(path, 8)
.expect("v27 migration must run cleanly against a rewound-meta v26 DB");
let conn = db.conn();
let (preserved_text, embedding_new): (String, Option<Vec<u8>>) = conn
.query_row(
"SELECT text, embedding_new FROM memories WHERE rid = ?1",
params![planted_rid],
|row| Ok((row.get::<_, String>(0)?, row.get::<_, Option<Vec<u8>>>(1)?)),
)
.unwrap();
assert_eq!(
preserved_text, "planted under v27 schema",
"v27 migration must NOT mutate existing memory data"
);
assert!(
embedding_new.is_none(),
"v27 migration must leave embedding_new as NULL on pre-existing rows"
);
assert!(
index_exists(&conn, "idx_reembed_events_generation"),
"v27 migration must recreate idx_reembed_events_generation"
);
conn.execute(
"INSERT INTO reembed_events (generation, phase, timestamp, payload_json) \
VALUES (?1, ?2, ?3, ?4)",
params![1_i64, "Probing", 0.0_f64, "{}"],
)
.unwrap();
let count: i64 = conn
.query_row(
"SELECT COUNT(*) FROM reembed_events WHERE generation = 1",
[],
|row| row.get(0),
)
.unwrap();
assert_eq!(
count, 1,
"reembed_events table must be writable post-migration"
);
}
#[test]
fn schema_v27_migration_replay_is_idempotent() {
use tempfile::NamedTempFile;
let tmp = NamedTempFile::new().unwrap();
let path = tmp.path().to_str().unwrap();
{
let _db = YantrikDB::new(path, 8).unwrap();
}
{
let conn = rusqlite::Connection::open(path).unwrap();
conn.execute(
"INSERT OR REPLACE INTO meta (key, value) VALUES ('schema_version', '26')",
[],
)
.unwrap();
}
let db = YantrikDB::new(path, 8)
.expect("v27 migration runner must heal rewound-meta deployments on a v27-schema DB");
db.record(
"post-v27-heal smoke",
"episodic",
0.5,
0.0,
604800.0,
&empty_meta(),
&vec_seed(1.0, 8),
"default",
0.8,
"general",
"user",
None,
)
.unwrap();
}
#[test]
fn try_publish_search_state_rejects_stale_generation() {
let db = YantrikDB::new(":memory:", 64).unwrap();
let initial = db.search_state.load_full();
assert_eq!(initial.generation, 0, "fresh engine starts at gen 0");
let stale_proposal = crate::engine::reembed::SearchState {
index_embedding: initial.index_embedding.clone(),
embedder: initial.embedder.clone(),
runtime_embedder_name: initial.runtime_embedder_name.clone(),
runtime_embedder_digest: initial.runtime_embedder_digest.clone(),
generation: 0,
covers_through_seq: 0,
hnsw_m: initial.hnsw_m,
hnsw_ef_construction: initial.hnsw_ef_construction,
hnsw_ef_search: initial.hnsw_ef_search,
vec_index: std::sync::Arc::clone(&initial.vec_index),
};
let advanced = crate::engine::reembed::SearchState {
index_embedding: initial.index_embedding.clone(),
embedder: initial.embedder.clone(),
runtime_embedder_name: initial.runtime_embedder_name.clone(),
runtime_embedder_digest: initial.runtime_embedder_digest.clone(),
generation: 2,
covers_through_seq: 0,
hnsw_m: initial.hnsw_m,
hnsw_ef_construction: initial.hnsw_ef_construction,
hnsw_ef_search: initial.hnsw_ef_search,
vec_index: std::sync::Arc::clone(&initial.vec_index),
};
db.search_state.store(std::sync::Arc::new(advanced));
let err = db
.try_publish_search_state(stale_proposal)
.expect_err("stale-generation publish must be rejected");
match err {
crate::error::YantrikDbError::SearchStatePublishStaleGeneration {
current_generation,
attempted_generation,
} => {
assert_eq!(current_generation, 2);
assert_eq!(attempted_generation, 0);
}
other => panic!("unexpected error variant: {other:?}"),
}
assert_eq!(
db.search_state.load().generation,
2,
"rejected publish must leave search_state untouched"
);
}
#[test]
fn try_publish_search_state_accepts_equal_generation_publish() {
let db = YantrikDB::new(":memory:", 64).unwrap();
let initial = db.search_state.load_full();
let same_gen = crate::engine::reembed::SearchState {
index_embedding: initial.index_embedding.clone(),
embedder: initial.embedder.clone(),
runtime_embedder_name: Some("rotated-name".to_string()),
runtime_embedder_digest: initial.runtime_embedder_digest.clone(),
generation: initial.generation,
covers_through_seq: initial.covers_through_seq,
hnsw_m: initial.hnsw_m,
hnsw_ef_construction: initial.hnsw_ef_construction,
hnsw_ef_search: initial.hnsw_ef_search,
vec_index: std::sync::Arc::clone(&initial.vec_index),
};
db.try_publish_search_state(same_gen)
.expect("equal-generation publish must be accepted");
assert_eq!(
db.search_state.load().runtime_embedder_name.as_deref(),
Some("rotated-name"),
"the equal-generation publish must have landed"
);
}
#[test]
fn try_publish_search_state_accepts_strictly_advancing_generation() {
let db = YantrikDB::new(":memory:", 64).unwrap();
let initial = db.search_state.load_full();
let advanced = crate::engine::reembed::SearchState {
index_embedding: initial.index_embedding.clone(),
embedder: initial.embedder.clone(),
runtime_embedder_name: initial.runtime_embedder_name.clone(),
runtime_embedder_digest: initial.runtime_embedder_digest.clone(),
generation: initial.generation + 1,
covers_through_seq: initial.covers_through_seq,
hnsw_m: initial.hnsw_m,
hnsw_ef_construction: initial.hnsw_ef_construction,
hnsw_ef_search: initial.hnsw_ef_search,
vec_index: std::sync::Arc::clone(&initial.vec_index),
};
db.try_publish_search_state(advanced)
.expect("strictly-advancing-generation publish must be accepted");
assert_eq!(
db.search_state.load().generation,
initial.generation + 1,
"the advanced publish must have landed"
);
}
#[test]
fn schema_v28_fresh_install_has_embedding_generation_and_active_generation() {
let db = YantrikDB::new(":memory:", 8).unwrap();
let conn = db.conn();
let cols = table_columns(&conn, "memories");
assert!(
cols.iter().any(|c| c == "embedding_generation"),
"v28: fresh install must add memories.embedding_generation column, got: {cols:?}"
);
assert!(
index_exists(&conn, "idx_memories_embedding_generation"),
"v28: fresh install must create idx_memories_embedding_generation"
);
let active_gen: Option<String> = conn
.query_row(
"SELECT value FROM meta WHERE key = 'active_generation'",
[],
|r| r.get(0),
)
.ok();
assert_eq!(
active_gen.as_deref(),
Some("0"),
"v28: fresh install must seed meta.active_generation = '0'"
);
let schema_version: String = conn
.query_row(
"SELECT value FROM meta WHERE key = 'schema_version'",
[],
|r| r.get(0),
)
.unwrap();
assert_eq!(
schema_version,
crate::base::schema::SCHEMA_VERSION.to_string(),
"fresh install stamps SCHEMA_VERSION"
);
}
#[test]
fn schema_v28_migration_from_v27_is_additive_and_idempotent() {
use tempfile::NamedTempFile;
let tmp = NamedTempFile::new().unwrap();
let path = tmp.path().to_str().unwrap();
let planted_rid = "01900000-0000-7000-8000-00000000c028";
{
let db = YantrikDB::new(path, 8).unwrap();
let conn = db.conn();
conn.execute(
"INSERT INTO memories (rid, type, text, embedding, created_at, updated_at, last_access, source, embedding_generation) \
VALUES (?1, 'episodic', 'planted under v28 schema', X'01020304', 0.0, 0.0, 0.0, 'user', 42)",
params![planted_rid],
)
.unwrap();
}
{
let conn = rusqlite::Connection::open(path).unwrap();
conn.execute(
"INSERT OR REPLACE INTO meta (key, value) VALUES ('schema_version', '27')",
[],
)
.unwrap();
}
let db =
YantrikDB::new(path, 8).expect("v28 migration runner must heal rewound-meta deployments");
let conn = db.conn();
let (text, gen): (String, i64) = conn
.query_row(
"SELECT text, embedding_generation FROM memories WHERE rid = ?1",
[&planted_rid],
|r| Ok((r.get(0)?, r.get(1)?)),
)
.unwrap();
assert_eq!(text, "planted under v28 schema");
assert_eq!(gen, 42, "migration must not mutate existing row data");
let schema_version: String = conn
.query_row(
"SELECT value FROM meta WHERE key = 'schema_version'",
[],
|r| r.get(0),
)
.unwrap();
assert_eq!(
schema_version,
crate::base::schema::SCHEMA_VERSION.to_string()
);
let active_gen: String = conn
.query_row(
"SELECT value FROM meta WHERE key = 'active_generation'",
[],
|r| r.get(0),
)
.unwrap();
assert_eq!(active_gen, "0");
}
#[test]
fn record_stamps_embedding_generation_from_search_state() {
let db = YantrikDB::new(":memory:", 8).unwrap();
let rid = db
.record(
"stamped",
"episodic",
0.5,
0.0,
604800.0,
&empty_meta(),
&vec_seed(1.0, 8),
"default",
0.8,
"general",
"user",
None,
)
.unwrap();
let conn = db.conn();
let stamped: i64 = conn
.query_row(
"SELECT embedding_generation FROM memories WHERE rid = ?1",
[&rid],
|r| r.get(0),
)
.unwrap();
assert_eq!(stamped, 0, "fresh engine: state.generation = 0, stamp = 0");
drop(conn);
let old_state = db.search_state.load_full();
let advanced = crate::engine::reembed::SearchState {
index_embedding: old_state.index_embedding.clone(),
embedder: old_state.embedder.clone(),
runtime_embedder_name: old_state.runtime_embedder_name.clone(),
runtime_embedder_digest: old_state.runtime_embedder_digest.clone(),
generation: 7,
covers_through_seq: old_state.covers_through_seq,
hnsw_m: old_state.hnsw_m,
hnsw_ef_construction: old_state.hnsw_ef_construction,
hnsw_ef_search: old_state.hnsw_ef_search,
vec_index: std::sync::Arc::clone(&old_state.vec_index),
};
db.try_publish_search_state(advanced).unwrap();
let rid2 = db
.record(
"stamped at gen 7",
"episodic",
0.5,
0.0,
604800.0,
&empty_meta(),
&vec_seed(2.0, 8),
"default",
0.8,
"general",
"user",
None,
)
.unwrap();
let conn = db.conn();
let stamped2: i64 = conn
.query_row(
"SELECT embedding_generation FROM memories WHERE rid = ?1",
[&rid2],
|r| r.get(0),
)
.unwrap();
assert_eq!(
stamped2, 7,
"record after generation advance must stamp the new generation"
);
}
#[test]
fn open_reads_durable_active_generation_into_search_state() {
use tempfile::NamedTempFile;
let tmp = NamedTempFile::new().unwrap();
let path = tmp.path().to_str().unwrap();
{
let db = YantrikDB::new(path, 8).unwrap();
assert_eq!(db.search_state.load().generation, 0);
}
{
let conn = rusqlite::Connection::open(path).unwrap();
conn.execute(
"INSERT OR REPLACE INTO meta (key, value) VALUES ('active_generation', '3')",
[],
)
.unwrap();
}
let db = YantrikDB::new(path, 8).unwrap();
assert_eq!(
db.search_state.load().generation,
3,
"open() must read meta.active_generation into SearchState.generation"
);
}
#[test]
fn set_embedder_routes_through_try_publish_search_state() {
use mode_test_embedders::FakeEmbedder;
let mut db = YantrikDB::new(":memory:", 64).unwrap();
let gen_before = db.search_state.load().generation;
db.set_embedder(Box::new(FakeEmbedder {
dim: 64,
fp: Some("sha256:check".to_string()),
name: Some("check".to_string()),
sentinel: 0.1,
}))
.unwrap();
let gen_after = db.search_state.load().generation;
assert_eq!(
gen_before, gen_after,
"set_embedder must preserve generation (only Phase-2 reembed advances it)"
);
assert!(db.has_embedder(), "set_embedder must have published");
}
#[test]
fn search_state_publish_is_atomic_under_concurrent_reads() {
use std::sync::atomic::{AtomicBool, Ordering};
use std::sync::Arc;
use std::thread;
let db = Arc::new(YantrikDB::new(":memory:", 64).unwrap());
let initial = db.search_state.load_full();
let stop = Arc::new(AtomicBool::new(false));
let mut handles = Vec::new();
for _ in 0..4 {
let db_c = Arc::clone(&db);
let stop_c = Arc::clone(&stop);
handles.push(thread::spawn(move || {
let mut observations: Vec<(u64, usize, u32)> = Vec::new();
while !stop_c.load(Ordering::Relaxed) {
let s = db_c.search_state.load_full();
observations.push((s.generation, s.dim(), s.hnsw_m));
}
observations
}));
}
for n in 1..=50u64 {
let prev = db.search_state.load_full();
let next = crate::engine::reembed::SearchState {
index_embedding: prev.index_embedding.clone(),
embedder: prev.embedder.clone(),
runtime_embedder_name: prev.runtime_embedder_name.clone(),
runtime_embedder_digest: prev.runtime_embedder_digest.clone(),
generation: prev.generation + 1,
covers_through_seq: prev.covers_through_seq + n,
hnsw_m: if n % 2 == 0 { 16 } else { 32 },
hnsw_ef_construction: prev.hnsw_ef_construction,
hnsw_ef_search: prev.hnsw_ef_search,
vec_index: std::sync::Arc::clone(&prev.vec_index),
};
db.try_publish_search_state(next).unwrap();
}
stop.store(true, Ordering::Relaxed);
let baseline = (initial.generation, initial.dim(), initial.hnsw_m);
for handle in handles {
let observations = handle.join().unwrap();
for (gen, dim, hnsw_m) in &observations {
let consistent = (*gen, *dim, *hnsw_m) == baseline
|| (*gen >= 1
&& *dim == initial.dim()
&& ((*gen % 2 == 0 && *hnsw_m == 16) || (*gen % 2 == 1 && *hnsw_m == 32)));
assert!(
consistent,
"torn SearchState observation: gen={gen} dim={dim} hnsw_m={hnsw_m} \
(expected even-gen→16 or odd-gen→32, or baseline)"
);
}
}
}
#[test]
fn open_recovery_discards_staging_when_sql_swap_uncommitted() {
use tempfile::NamedTempFile;
let tmp = NamedTempFile::new().unwrap();
let path = tmp.path().to_str().unwrap();
let planted_rid = "01900000-0000-7000-8000-00000000d017";
{
let db = YantrikDB::new(path, 8).unwrap();
db.record(
"pre-reembed row",
"episodic",
0.5,
0.0,
604800.0,
&empty_meta(),
&vec_seed(0.5, 8),
"default",
0.8,
"general",
"user",
None,
)
.unwrap();
let conn = db.conn();
conn.execute(
"UPDATE memories SET embedding_new = X'AABBCCDD', \
embedding_new_model = 'simulated-target' WHERE rowid = 1",
[],
)
.unwrap();
conn.execute(
"INSERT OR REPLACE INTO meta (key, value) VALUES ('reembed_state', ?1)",
params![serde_json::json!({
"generation": 5,
"phase": "Encoding",
"old_embedder": "old",
"new_embedder_name": "simulated-target",
})
.to_string()],
)
.unwrap();
let _ = planted_rid;
}
let db = YantrikDB::new(path, 8).unwrap();
let conn = db.conn();
let still_in_flight: Option<String> = conn
.query_row(
"SELECT value FROM meta WHERE key = 'reembed_state'",
[],
|r| r.get(0),
)
.ok();
assert!(
still_in_flight.is_none(),
"Layer 7 must clear meta.reembed_state on uncommitted-swap recovery; got: {still_in_flight:?}"
);
let staged: i64 = conn
.query_row(
"SELECT COUNT(*) FROM memories WHERE embedding_new IS NOT NULL OR \
embedding_new_model IS NOT NULL",
[],
|r| r.get(0),
)
.unwrap();
assert_eq!(
staged, 0,
"staging columns must be NULL after discard recovery"
);
let active: String = conn
.query_row(
"SELECT value FROM meta WHERE key = 'active_generation'",
[],
|r| r.get(0),
)
.unwrap();
assert_eq!(active, "0");
assert_eq!(db.search_state.load().generation, 0);
let aborted_recovery: i64 = conn
.query_row(
"SELECT COUNT(*) FROM reembed_events WHERE phase = 'Aborted' AND generation = 5 \
AND payload_json LIKE '%discarded_staging%'",
[],
|r| r.get(0),
)
.unwrap();
assert_eq!(aborted_recovery, 1, "Aborted recovery event must be logged");
}
#[test]
fn open_recovery_durable_swap_resumes_at_new_generation() {
use tempfile::NamedTempFile;
let tmp = NamedTempFile::new().unwrap();
let path = tmp.path().to_str().unwrap();
{
let db = YantrikDB::new(path, 8).unwrap();
let _ = db.record(
"pre-reembed row",
"episodic",
0.5,
0.0,
604800.0,
&empty_meta(),
&vec_seed(0.5, 8),
"default",
0.8,
"general",
"user",
None,
);
let conn = db.conn();
conn.execute(
"INSERT OR REPLACE INTO meta (key, value) VALUES ('active_generation', '3')",
[],
)
.unwrap();
conn.execute(
"UPDATE memories SET embedding_generation = 3 WHERE embedding IS NOT NULL",
[],
)
.unwrap();
conn.execute(
"INSERT OR REPLACE INTO meta (key, value) VALUES ('reembed_state', ?1)",
params![serde_json::json!({
"generation": 3,
"phase": "Swapping",
"old_embedder": "old",
"new_embedder_name": "new",
})
.to_string()],
)
.unwrap();
}
let db = YantrikDB::new(path, 8).unwrap();
let conn = db.conn();
let still_in_flight: Option<String> = conn
.query_row(
"SELECT value FROM meta WHERE key = 'reembed_state'",
[],
|r| r.get(0),
)
.ok();
assert!(
still_in_flight.is_none(),
"Layer 7 must clear meta.reembed_state on durable-swap recovery"
);
assert_eq!(
db.search_state.load().generation,
3,
"SearchState rebuilds at durable active generation"
);
let active: String = conn
.query_row(
"SELECT value FROM meta WHERE key = 'active_generation'",
[],
|r| r.get(0),
)
.unwrap();
assert_eq!(active, "3");
let completed_recovery: i64 = conn
.query_row(
"SELECT COUNT(*) FROM reembed_events WHERE phase = 'Completed' AND generation = 3 \
AND payload_json LIKE '%completed_durable%'",
[],
|r| r.get(0),
)
.unwrap();
assert_eq!(
completed_recovery, 1,
"Completed recovery event must be logged"
);
}
#[test]
fn open_with_uncommitted_staging_columns_stays_at_old_generation() {
use tempfile::NamedTempFile;
let tmp = NamedTempFile::new().unwrap();
let path = tmp.path().to_str().unwrap();
let planted_rid = {
let db = YantrikDB::new(path, 8).unwrap();
db.record(
"pre-reembed row",
"episodic",
0.5,
0.0,
604800.0,
&empty_meta(),
&vec_seed(0.5, 8),
"default",
0.8,
"general",
"user",
None,
)
.unwrap()
};
{
let conn = rusqlite::Connection::open(path).unwrap();
conn.execute(
"UPDATE memories SET embedding_new = X'AABBCCDD', \
embedding_new_model = 'simulated-new-embedder' WHERE rid = ?1",
params![planted_rid],
)
.unwrap();
}
let db = YantrikDB::new(path, 8).unwrap();
assert_eq!(
db.search_state.load().generation,
0,
"open() must not promote partial staging into the active generation"
);
let conn = db.conn();
let staged_present: bool = conn
.query_row(
"SELECT embedding_new IS NOT NULL FROM memories WHERE rid = ?1",
[&planted_rid],
|r| r.get(0),
)
.unwrap();
assert!(
staged_present,
"staged column survives the open (Phase-2 resume logic decides what to do with it)"
);
let active_present: bool = conn
.query_row(
"SELECT embedding IS NOT NULL FROM memories WHERE rid = ?1",
[&planted_rid],
|r| r.get(0),
)
.unwrap();
assert!(
active_present,
"pre-reembed active embedding bytes preserved"
);
let row_gen: i64 = conn
.query_row(
"SELECT embedding_generation FROM memories WHERE rid = ?1",
[&planted_rid],
|r| r.get(0),
)
.unwrap();
assert_eq!(
row_gen, 0,
"row's stamped generation unchanged by partial staging"
);
}
#[test]
fn covers_through_seq_is_durably_carried_on_published_search_state() {
let db = YantrikDB::new(":memory:", 8).unwrap();
let initial = db.search_state.load_full();
assert_eq!(initial.covers_through_seq, 0, "fresh engine: covers 0");
let next = crate::engine::reembed::SearchState {
index_embedding: initial.index_embedding.clone(),
embedder: initial.embedder.clone(),
runtime_embedder_name: initial.runtime_embedder_name.clone(),
runtime_embedder_digest: initial.runtime_embedder_digest.clone(),
generation: initial.generation + 1,
covers_through_seq: 12345,
hnsw_m: initial.hnsw_m,
hnsw_ef_construction: initial.hnsw_ef_construction,
hnsw_ef_search: initial.hnsw_ef_search,
vec_index: std::sync::Arc::clone(&initial.vec_index),
};
db.try_publish_search_state(next).unwrap();
assert_eq!(
db.search_state.load().covers_through_seq,
12345,
"published covers_through_seq must be readable from the active SearchState"
);
let bumped = crate::engine::reembed::SearchState {
index_embedding: initial.index_embedding.clone(),
embedder: initial.embedder.clone(),
runtime_embedder_name: initial.runtime_embedder_name.clone(),
runtime_embedder_digest: initial.runtime_embedder_digest.clone(),
generation: initial.generation + 2,
covers_through_seq: 98765,
hnsw_m: initial.hnsw_m,
hnsw_ef_construction: initial.hnsw_ef_construction,
hnsw_ef_search: initial.hnsw_ef_search,
vec_index: std::sync::Arc::clone(&initial.vec_index),
};
db.try_publish_search_state(bumped).unwrap();
assert_eq!(
db.search_state.load().covers_through_seq,
98765,
"covers_through_seq advances per swap"
);
}
#[test]
fn record_text_round_trip_through_queue_path_under_reembed() {
use mode_test_embedders::FakeEmbedder;
let mut db = YantrikDB::new(":memory:", 64).unwrap();
db.set_embedder(Box::new(FakeEmbedder {
dim: 64,
fp: Some("sha256:queued-test".to_string()),
name: Some("queued-test-embedder".to_string()),
sentinel: 0.7,
}))
.unwrap();
db.write_router.switch_to_queueing();
let rid = db
.record_text(
"queued-round-trip-text",
"episodic",
0.5,
0.0,
604800.0,
&serde_json::json!({"k": "v"}),
"default",
0.8,
"general",
"user",
None,
)
.unwrap();
let conn = db.conn();
let mem_count: i64 = conn
.query_row(
"SELECT COUNT(*) FROM memories WHERE rid = ?1",
[&rid],
|r| r.get(0),
)
.unwrap();
assert_eq!(mem_count, 0, "queued write does not touch memories table");
let (op_type, applied, applied_generation, embedding_model, payload): (
String,
i64,
Option<i64>,
Option<String>,
String,
) = conn
.query_row(
"SELECT op_type, applied, applied_generation, embedding_model, payload \
FROM oplog WHERE target_rid = ?1 AND op_type = 'record'",
[&rid],
|r| Ok((r.get(0)?, r.get(1)?, r.get(2)?, r.get(3)?, r.get(4)?)),
)
.unwrap();
assert_eq!(op_type, "record");
assert_eq!(applied, 0, "queued op is applied=0");
assert_eq!(
applied_generation, None,
"queued op has applied_generation=NULL (post-swap materializer fills under new gen)"
);
assert_eq!(
embedding_model.as_deref(),
Some("queued-test-embedder"),
"queued op carries the active runtime embedder name"
);
let v: serde_json::Value = serde_json::from_str(&payload).unwrap();
assert_eq!(
v["text"].as_str(),
Some("queued-round-trip-text"),
"payload preserves the original text for re-encode"
);
}
#[test]
fn boundary_audit_pattern_detects_synthetic_violation() {
let synthetic_violation =
" let sql = \"SELECT rid, embedding FROM memories WHERE rid = ?1\";";
let lower = synthetic_violation.to_ascii_lowercase();
let patterns = [
"select embedding ",
"select embedding,",
"select embedding\"",
"select embedding\\",
", embedding ",
", embedding,",
", embedding\"",
", embedding\\",
", embedding\n",
];
let any_match = patterns.iter().any(|p| lower.contains(p));
assert!(
any_match,
"boundary audit pattern must catch the synthetic raw-SQL-embedding pattern; \
if this asserts, the audit in durable_embeddings.rs is letting violations slip"
);
let allowlist_safe = " let sql = \"SELECT rid, embedding_hash FROM memories\";";
let lower_safe = allowlist_safe.to_ascii_lowercase();
let safe_match = patterns.iter().any(|p| lower_safe.contains(p));
assert!(
!safe_match,
"audit must NOT flag the allowlist-safe `embedding_hash` pattern; \
the audit over-rejects which would prevent legitimate refactors"
);
}
#[test]
fn record_with_rid_backpressure_does_not_leak_orphan_memories_row() {
let db = YantrikDB::new(":memory:", 8).unwrap();
let _ = db;
let db = YantrikDB::new(":memory:", 8).unwrap();
let dim = db.embedding_dim();
let mut last_backpressure_rid: Option<String> = None;
for i in 0..400 {
let embedding: Vec<f32> = (0..dim).map(|j| ((i + j) as f32) * 0.001).collect();
let attempted_rid = format!("orphan-test-rid-{i}");
let res = db.record_with_rid(
&attempted_rid,
&format!("orphan-test-text-{i}"),
"episodic",
0.5,
0.0,
604800.0,
&serde_json::json!({}),
&embedding,
"default",
0.8,
"general",
"user",
None,
(i as i64) * 1_000_000,
&[],
"test-embedder",
None,
);
if let Err(crate::error::YantrikDbError::Backpressure { .. }) = res {
last_backpressure_rid = Some(attempted_rid);
break;
}
}
let rid = last_backpressure_rid
.expect("test infrastructure: pump must produce a Backpressure within 400 writes");
let conn = db.conn();
let exists: i64 = conn
.query_row(
"SELECT COUNT(*) FROM memories WHERE rid = ?1",
params![&rid],
|r| r.get(0),
)
.unwrap();
assert_eq!(
exists, 0,
"v0.7.19: compensating DELETE must remove memories row when Backpressure fires; \
leaving the row is the trader's 23k-orphan pattern"
);
let oplog_count: i64 = conn
.query_row(
"SELECT COUNT(*) FROM oplog WHERE target_rid = ?1 AND op_type = 'record_with_rid'",
params![&rid],
|r| r.get(0),
)
.unwrap();
assert_eq!(
oplog_count, 0,
"Backpressure path correctly leaves no oplog entry"
);
}
#[test]
fn record_with_rid_backpressure_does_not_overcount_session_memory_count() {
let db = YantrikDB::new(":memory:", 8).unwrap();
let sid = db
.session_start("default", "claude", &serde_json::json!({}))
.unwrap();
let dim = db.embedding_dim();
let mut saw_backpressure = false;
for i in 0..400 {
let embedding: Vec<f32> = (0..dim).map(|j| ((i + j) as f32) * 0.001).collect();
let res = db.record_with_rid(
&format!("sess-count-rid-{i}"),
&format!("sess-count-text-{i}"),
"episodic",
0.5,
0.0,
604800.0,
&serde_json::json!({}),
&embedding,
"default",
0.8,
"general",
"user",
None,
(i as i64) * 1_000_000,
&[],
"test-embedder",
None,
);
if let Err(crate::error::YantrikDbError::Backpressure { .. }) = res {
saw_backpressure = true;
break;
}
}
assert!(
saw_backpressure,
"test infrastructure: pump must produce a Backpressure within 400 writes"
);
let conn = db.conn();
let actual_rows: i64 = conn
.query_row(
"SELECT COUNT(*) FROM memories WHERE session_id = ?1",
params![&sid],
|r| r.get(0),
)
.unwrap();
let stored_count: i64 = conn
.query_row(
"SELECT memory_count FROM sessions WHERE session_id = ?1",
params![&sid],
|r| r.get(0),
)
.unwrap();
assert_eq!(
stored_count, actual_rows,
"v0.7.23: session memory_count ({stored_count}) must match the number of \
memories rows actually linked to the session ({actual_rows}); a higher \
count means the backpressure compensation failed to reverse the bump"
);
}
#[test]
fn record_coerces_blank_namespace_to_default() {
let db = YantrikDB::new(":memory:", 8).unwrap();
let rid = db
.record(
"alpha",
"semantic",
0.5,
0.0,
604800.0,
&empty_meta(),
&vec_seed(1.0, 8),
"",
0.8,
"d",
"user",
None,
)
.unwrap();
assert_eq!(db.get(&rid).unwrap().unwrap().namespace, "default");
let rid_ws = db
.record(
"beta",
"semantic",
0.5,
0.0,
604800.0,
&empty_meta(),
&vec_seed(2.0, 8),
" ",
0.8,
"d",
"user",
None,
)
.unwrap();
assert_eq!(db.get(&rid_ws).unwrap().unwrap().namespace, "default");
let rid_ns = db
.record(
"gamma",
"semantic",
0.5,
0.0,
604800.0,
&empty_meta(),
&vec_seed(3.0, 8),
"acme",
0.8,
"d",
"user",
None,
)
.unwrap();
assert_eq!(db.get(&rid_ns).unwrap().unwrap().namespace, "acme");
}
#[test]
fn record_with_rid_coerces_blank_namespace_to_default() {
let db = YantrikDB::new(":memory:", 8).unwrap();
db.record_with_rid(
"blank-ns-rid",
"alpha",
"semantic",
0.5,
0.0,
604800.0,
&empty_meta(),
&vec_seed(1.0, 8),
"",
0.8,
"d",
"user",
None,
1_000_000,
&[],
"test-embedder",
None,
)
.unwrap();
assert_eq!(
db.get("blank-ns-rid").unwrap().unwrap().namespace,
"default",
"record_with_rid must coerce blank namespace to default (server write path)"
);
}
#[test]
fn list_records_enumerates_by_kind_with_keyset_cursor() {
let db = YantrikDB::new(":memory:", 8).unwrap();
let meta = |k: &str, d: &str| serde_json::json!({ "kind": k, "drive_id": d });
let mut reply_rids = Vec::new();
for i in 0..5 {
reply_rids.push(
db.record(
&format!("reply {i}"),
"semantic",
0.5,
0.0,
604800.0,
&meta("operator_reply_v1", "D1"),
&vec_seed(i as f32, 8),
"ns",
0.8,
"d",
"user",
None,
)
.unwrap(),
);
db.record(
&format!("thought {i}"),
"semantic",
0.5,
0.0,
604800.0,
&meta("focal_thought", "D2"),
&vec_seed((i + 100) as f32, 8),
"ns",
0.8,
"d",
"user",
None,
)
.unwrap();
}
let (p1, c1) = db
.list_records(
Some("ns"),
Some("operator_reply_v1"),
None,
None,
None,
None,
3,
"asc",
)
.unwrap();
assert_eq!(p1.len(), 3);
assert!(p1.iter().all(|m| m.metadata["kind"] == "operator_reply_v1"));
assert!(
p1[0].rid < p1[1].rid && p1[1].rid < p1[2].rid,
"rid ascending"
);
let c1 = c1.expect("full page yields a cursor");
let (p2, c2) = db
.list_records(
Some("ns"),
Some("operator_reply_v1"),
None,
None,
None,
Some(&c1),
3,
"asc",
)
.unwrap();
assert_eq!(p2.len(), 2);
assert!(c2.is_none(), "partial page yields no cursor");
let mut got: Vec<String> = p1.iter().chain(p2.iter()).map(|m| m.rid.clone()).collect();
got.sort();
let mut want = reply_rids.clone();
want.sort();
assert_eq!(
got, want,
"keyset pages cover the full kind=X set exactly once"
);
let (desc, _) = db
.list_records(
Some("ns"),
Some("operator_reply_v1"),
None,
None,
None,
None,
10,
"desc",
)
.unwrap();
assert_eq!(desc.len(), 5);
assert!(desc[0].rid > desc[4].rid, "desc = newest first");
let (d1, _) = db
.list_records(Some("ns"), None, Some("D1"), None, None, None, 100, "asc")
.unwrap();
assert_eq!(d1.len(), 5);
assert!(d1.iter().all(|m| m.metadata["drive_id"] == "D1"));
let (none, cn) = db
.list_records(Some("ns"), Some("nope"), None, None, None, None, 10, "asc")
.unwrap();
assert!(none.is_empty() && cn.is_none());
}
#[test]
fn list_records_tolerates_non_json_metadata() {
let db = YantrikDB::new(":memory:", 8).unwrap();
let good = db
.record(
"good",
"semantic",
0.5,
0.0,
604800.0,
&serde_json::json!({ "kind": "k1" }),
&vec_seed(1.0, 8),
"ns",
0.8,
"d",
"user",
None,
)
.unwrap();
let bad = db
.record(
"bad",
"semantic",
0.5,
0.0,
604800.0,
&empty_meta(),
&vec_seed(2.0, 8),
"ns",
0.8,
"d",
"user",
None,
)
.unwrap();
db.conn()
.execute(
"UPDATE memories SET metadata = 'not::json::{' WHERE rid = ?1",
params![bad],
)
.unwrap();
let (page, _) = db
.list_records(Some("ns"), Some("k1"), None, None, None, None, 10, "asc")
.unwrap();
assert_eq!(page.len(), 1);
assert_eq!(page[0].rid, good);
let (all, _) = db
.list_records(Some("ns"), None, None, None, None, None, 10, "asc")
.unwrap();
assert_eq!(all.len(), 2);
}
#[test]
fn migrate_v31_to_v32_indexes_existing_rows_and_tolerates_bad_json() {
let conn = rusqlite::Connection::open_in_memory().unwrap();
conn.execute_batch(
"CREATE TABLE memories (rid TEXT PRIMARY KEY, metadata TEXT, consolidation_status TEXT);\
INSERT INTO memories VALUES ('r1', '{\"kind\":\"operator_reply_v1\",\"drive_id\":\"D1\"}', 'active');\
INSERT INTO memories VALUES ('r2', 'not-valid-json-{', 'active');",
)
.unwrap();
conn.execute_batch(crate::schema::MIGRATE_V31_TO_V32)
.unwrap();
let kind: Option<String> = conn
.query_row("SELECT kind FROM memories WHERE rid='r1'", [], |r| r.get(0))
.unwrap();
assert_eq!(kind.as_deref(), Some("operator_reply_v1"));
let drive: Option<String> = conn
.query_row("SELECT drive_id FROM memories WHERE rid='r1'", [], |r| {
r.get(0)
})
.unwrap();
assert_eq!(drive.as_deref(), Some("D1"));
let bad: Option<String> = conn
.query_row("SELECT kind FROM memories WHERE rid='r2'", [], |r| r.get(0))
.unwrap();
assert_eq!(bad, None);
let n: i64 = conn
.query_row(
"SELECT COUNT(*) FROM memories WHERE kind='operator_reply_v1'",
[],
|r| r.get(0),
)
.unwrap();
assert_eq!(n, 1);
}
#[test]
fn think_extract_attribute_claims_detects_free_text_value_update() {
let db = YantrikDB::new(":memory:", 8).unwrap();
db.record(
"Brand color is blue #1F4E79.",
"semantic",
0.9,
0.0,
604800.0,
&empty_meta(),
&vec_seed(1.0, 8),
"u",
0.8,
"brand_color",
"user",
None,
)
.unwrap();
db.record(
"Brand color is now green #2E7D32.",
"semantic",
0.95,
0.0,
604800.0,
&empty_meta(),
&vec_seed(2.0, 8),
"u",
0.8,
"brand_color",
"user",
None,
)
.unwrap();
let off = ThinkConfig {
run_consolidation: false,
run_personality: false,
..Default::default()
};
assert_eq!(
db.think(&off).unwrap().conflicts_found,
0,
"baseline: free-text attribute-value update is not detected (the reported gap)"
);
let on = ThinkConfig {
run_consolidation: false,
run_personality: false,
extract_attribute_claims: true,
..Default::default()
};
let res = db.think(&on).unwrap();
assert!(
res.conflicts_found >= 1,
"v0.7.23: attribute-value conflict should be detected, got {}",
res.conflicts_found
);
let conn = db.conn();
let n: i64 = conn
.query_row(
"SELECT COUNT(*) FROM edges WHERE src = 'brand color' AND rel_type = 'is'",
[],
|r| r.get(0),
)
.unwrap();
assert_eq!(
n, 2,
"two attribute-value claims (blue, green) should exist for (brand color, is)"
);
}
#[test]
fn schema_v29_fresh_install_has_replication_apply_log_table() {
let db = YantrikDB::new(":memory:", 8).unwrap();
let conn = db.conn();
let cols = table_columns(&conn, "replication_apply_log");
for required in ["rid", "op_type", "source_actor", "applied_at"] {
assert!(
cols.iter().any(|c| c == required),
"v29: replication_apply_log missing column {required}, got: {cols:?}"
);
}
assert!(
index_exists(&conn, "idx_replication_apply_log_source_actor"),
"v29: idx_replication_apply_log_source_actor must exist"
);
let schema_version: String = conn
.query_row(
"SELECT value FROM meta WHERE key = 'schema_version'",
[],
|r| r.get(0),
)
.unwrap();
let v: i32 = schema_version.parse().unwrap();
assert!(
v >= 29,
"schema_version should be at least 29 (when replication_apply_log landed), got {v}",
);
}
#[test]
fn correct_preserves_rid_and_created_at() {
let db = YantrikDB::new(":memory:", 8).unwrap();
let rid = db
.record(
"original text",
"episodic",
0.5,
0.0,
604800.0,
&empty_meta(),
&vec_seed(1.0, 8),
"default",
0.8,
"general",
"user",
None,
)
.unwrap();
let original = db.get(&rid).unwrap().unwrap();
let original_created_at = original.created_at;
std::thread::sleep(std::time::Duration::from_millis(10));
let result = db
.correct(
&rid,
Some("corrected text"),
None,
None,
None,
"test correction",
)
.unwrap();
assert_eq!(result.corrected_rid, rid, "rid must be preserved");
assert_eq!(result.original_rid, rid);
assert!(!result.original_tombstoned);
assert_eq!(result.revision_num, 1);
let updated = db.get(&rid).unwrap().unwrap();
assert_eq!(updated.text, "corrected text");
assert!(
(updated.created_at - original_created_at).abs() < 1e-9,
"created_at must be preserved (was {}, became {})",
original_created_at,
updated.created_at,
);
}
#[test]
fn correct_writes_revision_row_with_prior_state() {
let db = YantrikDB::new(":memory:", 8).unwrap();
let rid = db
.record(
"v0",
"episodic",
0.5,
-0.2,
604800.0,
&serde_json::json!({"key": "before"}),
&vec_seed(1.0, 8),
"default",
0.8,
"general",
"user",
None,
)
.unwrap();
let _ = db
.correct(
&rid,
Some("v1"),
Some(&serde_json::json!({"key": "after"})),
Some(0.7),
Some(0.3),
"first correction",
)
.unwrap();
let history = db.history(&rid).unwrap();
assert_eq!(history.len(), 1, "one revision expected");
let rev = &history[0];
assert_eq!(rev.revision_num, 1);
assert_eq!(rev.prior_text, "v0");
assert!((rev.prior_importance - 0.5).abs() < 1e-9);
assert!((rev.prior_valence + 0.2).abs() < 1e-9);
assert_eq!(rev.reason, "first correction");
assert_eq!(
rev.prior_metadata.get("key").and_then(|v| v.as_str()),
Some("before")
);
}
#[test]
fn correct_multiple_revisions_increment_revision_num() {
let db = YantrikDB::new(":memory:", 8).unwrap();
let rid = db
.record(
"v0",
"episodic",
0.5,
0.0,
604800.0,
&empty_meta(),
&vec_seed(1.0, 8),
"default",
0.8,
"general",
"user",
None,
)
.unwrap();
let r1 = db
.correct(&rid, Some("v1"), None, None, None, "first")
.unwrap();
let r2 = db
.correct(&rid, Some("v2"), None, None, None, "second")
.unwrap();
let r3 = db
.correct(&rid, Some("v3"), None, None, None, "third")
.unwrap();
assert_eq!(r1.revision_num, 1);
assert_eq!(r2.revision_num, 2);
assert_eq!(r3.revision_num, 3);
let history = db.history(&rid).unwrap();
assert_eq!(history.len(), 3, "three revisions expected");
assert_eq!(history[0].prior_text, "v0");
assert_eq!(history[1].prior_text, "v1");
assert_eq!(history[2].prior_text, "v2");
let final_state = db.get(&rid).unwrap().unwrap();
assert_eq!(final_state.text, "v3");
}
#[test]
fn correct_rejects_empty_reason() {
let db = YantrikDB::new(":memory:", 8).unwrap();
let rid = db
.record(
"text",
"episodic",
0.5,
0.0,
604800.0,
&empty_meta(),
&vec_seed(1.0, 8),
"default",
0.8,
"general",
"user",
None,
)
.unwrap();
let err = db
.correct(&rid, Some("new"), None, None, None, "")
.expect_err("empty reason must be rejected");
match err {
crate::error::YantrikDbError::InvalidInput(msg) => {
assert!(msg.contains("reason"), "got: {msg}");
}
other => panic!("expected InvalidInput, got {other:?}"),
}
let err2 = db
.correct(&rid, Some("new"), None, None, None, " \t\n ")
.expect_err("whitespace-only reason must be rejected");
assert!(matches!(
err2,
crate::error::YantrikDbError::InvalidInput(_)
));
}
#[test]
fn correct_rejects_no_mutation_fields() {
let db = YantrikDB::new(":memory:", 8).unwrap();
let rid = db
.record(
"text",
"episodic",
0.5,
0.0,
604800.0,
&empty_meta(),
&vec_seed(1.0, 8),
"default",
0.8,
"general",
"user",
None,
)
.unwrap();
let err = db
.correct(&rid, None, None, None, None, "no fields supplied")
.expect_err("no-op correction must be rejected");
assert!(matches!(err, crate::error::YantrikDbError::InvalidInput(_)));
}
#[test]
fn correct_preserves_inbound_graph_edges() {
let db = YantrikDB::new(":memory:", 8).unwrap();
let rid_subject = db
.record(
"subject memory",
"episodic",
0.5,
0.0,
604800.0,
&empty_meta(),
&vec_seed(1.0, 8),
"default",
0.8,
"general",
"user",
None,
)
.unwrap();
db.relate("anchor_entity", &rid_subject, "tags", 1.0)
.unwrap();
let edges_before = db.get_edges(&rid_subject).unwrap();
assert!(!edges_before.is_empty(), "edge should exist before correct");
db.correct(
&rid_subject,
Some("updated subject"),
None,
None,
None,
"test link integrity",
)
.unwrap();
let edges_after = db.get_edges(&rid_subject).unwrap();
assert_eq!(
edges_before.len(),
edges_after.len(),
"inbound edges must be preserved across correct(); \
this is the central v0.7.20 win over the v0.7.19 tombstone semantics"
);
}
#[test]
fn correct_history_empty_for_never_corrected_record() {
let db = YantrikDB::new(":memory:", 8).unwrap();
let rid = db
.record(
"never corrected",
"episodic",
0.5,
0.0,
604800.0,
&empty_meta(),
&vec_seed(1.0, 8),
"default",
0.8,
"general",
"user",
None,
)
.unwrap();
let history = db.history(&rid).unwrap();
assert!(
history.is_empty(),
"never-corrected record has no revisions"
);
}
#[test]
fn schema_v30_fresh_install_has_record_revisions_table() {
let db = YantrikDB::new(":memory:", 8).unwrap();
let conn = db.conn();
let count: i64 = conn
.query_row(
"SELECT COUNT(*) FROM sqlite_master \
WHERE type = 'table' AND name = 'record_revisions'",
[],
|row| row.get(0),
)
.unwrap();
assert_eq!(
count, 1,
"record_revisions table must exist on fresh install"
);
let schema_version: String = conn
.query_row(
"SELECT value FROM meta WHERE key = 'schema_version'",
[],
|r| r.get(0),
)
.unwrap();
let v: i32 = schema_version.parse().unwrap();
assert!(
v >= crate::base::schema::SCHEMA_VERSION,
"schema_version must be at least {}, got {v}",
crate::base::schema::SCHEMA_VERSION,
);
}
#[test]
fn schema_v31_fresh_install_has_record_links_table() {
let db = YantrikDB::new(":memory:", 8).unwrap();
let conn = db.conn();
let table: i64 = conn
.query_row(
"SELECT COUNT(*) FROM sqlite_master \
WHERE type = 'table' AND name = 'record_links'",
[],
|row| row.get(0),
)
.unwrap();
assert_eq!(table, 1, "record_links table must exist on fresh install");
let idx: i64 = conn
.query_row(
"SELECT COUNT(*) FROM sqlite_master WHERE type = 'index' \
AND name IN ('idx_record_links_source', 'idx_record_links_target')",
[],
|row| row.get(0),
)
.unwrap();
assert_eq!(idx, 2, "both record_links covering indexes must exist");
assert!(crate::base::schema::SCHEMA_VERSION >= 31);
}
fn seed_for_certainty_test(db: &YantrikDB, n: usize) -> Vec<String> {
let mut rids = Vec::with_capacity(n);
for i in 0..n {
let certainty = (i as f64 + 1.0) / (n as f64);
let rid = db
.record(
&format!("certainty test memory {i}"),
"semantic",
0.5,
0.0,
604800.0,
&empty_meta(),
&vec_seed(1.0 + (i as f32) * 0.001, 8),
"default",
certainty,
"general",
"user",
None,
)
.unwrap();
rids.push(rid);
std::thread::sleep(std::time::Duration::from_millis(2));
}
rids
}
#[test]
fn recall_certainty_min_filters_low_certainty() {
let db = YantrikDB::new(":memory:", 8).unwrap();
let _rids = seed_for_certainty_test(&db, 5);
let results = db
.recall(
&vec_seed(1.0, 8),
10,
None,
None,
false,
false,
None,
true,
None,
None,
None,
Some(0.5), None,
)
.unwrap();
assert!(
!results.is_empty(),
"expected at least one high-certainty result"
);
for r in &results {
assert!(
r.certainty >= 0.5,
"certainty_min=0.5 must filter out cert={}: rid={}",
r.certainty,
r.rid
);
}
}
#[test]
fn recall_order_certainty_returns_results_in_certainty_desc() {
let db = YantrikDB::new(":memory:", 8).unwrap();
let _rids = seed_for_certainty_test(&db, 5);
let results = db
.recall(
&vec_seed(1.0, 8),
10,
None,
None,
false,
false,
None,
true,
None,
None,
None,
None,
Some("certainty"),
)
.unwrap();
assert!(
results.len() >= 2,
"need at least 2 results to verify ordering, got {}",
results.len()
);
for w in results.windows(2) {
assert!(
w[0].certainty >= w[1].certainty,
"order=certainty must be non-increasing; got [{}, {}]",
w[0].certainty,
w[1].certainty,
);
}
}
#[test]
fn recall_order_recency_returns_results_in_created_at_desc() {
let db = YantrikDB::new(":memory:", 8).unwrap();
let _rids = seed_for_certainty_test(&db, 5);
let results = db
.recall(
&vec_seed(1.0, 8),
10,
None,
None,
false,
false,
None,
true,
None,
None,
None,
None,
Some("recency"),
)
.unwrap();
assert!(
results.len() >= 2,
"need at least 2 results to verify ordering, got {}",
results.len()
);
for w in results.windows(2) {
assert!(
w[0].created_at >= w[1].created_at,
"order=recency must be non-increasing on created_at; got [{}, {}]",
w[0].created_at,
w[1].created_at,
);
}
}
#[test]
fn recall_order_invalid_string_returns_invalid_input_error() {
let db = YantrikDB::new(":memory:", 8).unwrap();
let _rids = seed_for_certainty_test(&db, 3);
let err = db
.recall(
&vec_seed(1.0, 8),
10,
None,
None,
false,
false,
None,
true,
None,
None,
None,
None,
Some("most_relevant"), )
.expect_err("invalid order string must return error");
match err {
crate::error::YantrikDbError::InvalidInput(msg) => {
assert!(
msg.contains("order") && msg.contains("most_relevant"),
"error message should name the bad order value, got: {msg}"
);
}
other => panic!("expected InvalidInput, got {other:?}"),
}
}
#[test]
fn recall_default_order_is_relevance_unchanged() {
let db = YantrikDB::new(":memory:", 8).unwrap();
let _rids = seed_for_certainty_test(&db, 5);
let results = db
.recall(
&vec_seed(1.0, 8),
10,
None,
None,
false,
false,
None,
true,
None,
None,
None,
None,
None, )
.unwrap();
assert!(
results.len() >= 2,
"need at least 2 results to verify default order, got {}",
results.len()
);
for w in results.windows(2) {
assert!(
w[0].score >= w[1].score,
"default order must be score-desc (relevance); got [{}, {}]",
w[0].score,
w[1].score,
);
}
}