use kimetsu_core::KimetsuResult;
use kimetsu_core::ids::new_id;
use kimetsu_core::memory::{MemoryScope, normalize_memory_text};
use rusqlite::{Connection, OptionalExtension, params};
use serde::{Deserialize, Serialize};
use time::OffsetDateTime;
use time::format_description::well_known::Rfc3339;
use crate::embeddings::{Embedder, cosine_similarity, decode_embedding};
pub fn conflict_detection_enabled(config_value: bool) -> bool {
match std::env::var("KIMETSU_DETECT_CONFLICTS") {
Ok(raw) => {
let v = raw.trim().to_ascii_lowercase();
if v.is_empty() {
config_value
} else {
!matches!(v.as_str(), "0" | "false" | "off" | "no")
}
}
Err(_) => config_value,
}
}
pub const DEFAULT_CONFLICT_THRESHOLD: f32 = 0.8;
pub const DEFAULT_TOP_K: u32 = 3;
pub const NEAR_TIE_BAND: f32 = 0.15;
pub fn resolve_conflicts_enabled(config_value: bool) -> bool {
match std::env::var("KIMETSU_RESOLVE_CONFLICTS") {
Ok(raw) => {
let v = raw.trim().to_ascii_lowercase();
if v.is_empty() {
config_value
} else {
!matches!(v.as_str(), "0" | "false" | "off" | "no")
}
}
Err(_) => config_value,
}
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub enum ResolutionOutcome {
AutoResolvedNewWon,
AutoResolvedExistingWon,
NearTieQueued,
}
pub fn resolution_score(confidence: f32, created_at_rfc3339: &str) -> f32 {
let age_days = match OffsetDateTime::parse(created_at_rfc3339, &Rfc3339) {
Ok(ts) => {
let now = OffsetDateTime::now_utc();
let secs = (now - ts).whole_seconds().max(0);
secs as f64 / 86_400.0
}
Err(_) => 0.0, };
const HALF_LIFE_DAYS: f64 = 30.0;
let recency_weight = (-std::f64::consts::LN_2 / HALF_LIFE_DAYS * age_days).exp() as f32;
(confidence.clamp(0.0, 1.0) * recency_weight).clamp(0.0, 1.0)
}
#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct ConflictHit {
pub existing_memory_id: String,
pub existing_kind: String,
pub existing_text: String,
pub similarity: f32,
}
#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct ConflictReport {
pub conflict_id: String,
pub new_memory_id: String,
pub new_text: String,
pub existing_memory_id: String,
pub existing_text: String,
pub scope: String,
pub kind: String,
pub similarity: f32,
pub detected_at: String,
pub resolved_at: Option<String>,
pub resolution: Option<String>,
}
pub fn find_potential_conflicts(
conn: &Connection,
scope: &MemoryScope,
new_text: &str,
embedder: &dyn Embedder,
top_k: u32,
threshold: f32,
) -> KimetsuResult<Vec<ConflictHit>> {
find_potential_conflicts_with_vec(
conn, scope, new_text, None, embedder, None, top_k, threshold,
)
}
#[allow(clippy::too_many_arguments)]
pub(crate) fn find_potential_conflicts_with_vec(
conn: &Connection,
scope: &MemoryScope,
new_text: &str,
precomputed_vec: Option<&[f32]>,
embedder: &dyn Embedder,
exclude_id: Option<&str>,
top_k: u32,
threshold: f32,
) -> KimetsuResult<Vec<ConflictHit>> {
if embedder.is_noop() {
return Ok(Vec::new());
}
let new_vec: Vec<f32>;
let query_vec: &[f32] = if let Some(v) = precomputed_vec {
v
} else {
new_vec = embedder
.embed(new_text)
.map_err(|e| format!("embedder failed during conflict scan: {e}"))?;
if new_vec.len() != embedder.dim() {
return Err(format!(
"embedder {} returned {} dims, expected {}",
embedder.model_id(),
new_vec.len(),
embedder.dim()
)
.into());
}
&new_vec
};
let new_normalized = normalize_memory_text(new_text);
let scope_label = scope.to_string();
let active_model = embedder.model_id();
#[cfg_attr(not(feature = "embeddings"), allow(unused_variables))]
let pool_size = (top_k * 8).max(64) as i64;
#[cfg(feature = "embeddings")]
{
let handle = crate::ann::handle_for_query(conn, query_vec.len(), active_model)?;
let ann_rowids: Vec<i64> = handle
.read()
.unwrap_or_else(|p| p.into_inner())
.search(query_vec, pool_size as usize)?
.into_iter()
.map(|(rowid, _)| rowid)
.collect();
if !ann_rowids.is_empty() {
let placeholders: String = ann_rowids
.iter()
.enumerate()
.map(|(i, _)| format!("?{}", i + 1))
.collect::<Vec<_>>()
.join(", ");
let sql = format!(
"SELECT memory_id, kind, text, normalized_text, embedding, embedding_model
FROM memories
WHERE invalidated_at IS NULL
AND scope = '{scope_label}'
AND embedding_model = '{active_model}'
AND rowid IN ({placeholders})"
);
let mut stmt = conn.prepare(&sql)?;
let params_vec: Vec<&dyn rusqlite::ToSql> = ann_rowids
.iter()
.map(|n| n as &dyn rusqlite::ToSql)
.collect();
let rows_iter = stmt.query_map(params_vec.as_slice(), |row| {
Ok((
row.get::<_, String>(0)?,
row.get::<_, String>(1)?,
row.get::<_, String>(2)?,
row.get::<_, String>(3)?,
row.get::<_, Vec<u8>>(4)?,
))
})?;
let mut hits: Vec<ConflictHit> = Vec::new();
for row in rows_iter {
let (existing_id, kind, text, normalized, bytes) = row?;
if normalized == new_normalized {
continue;
}
if let Some(excl) = exclude_id {
if existing_id == excl {
continue;
}
}
let Ok(existing_vec) = decode_embedding(&bytes, Some(query_vec.len())) else {
continue;
};
let sim = cosine_similarity(query_vec, &existing_vec);
if sim >= threshold {
hits.push(ConflictHit {
existing_memory_id: existing_id,
existing_kind: kind,
existing_text: text,
similarity: sim,
});
}
}
hits.sort_by(|a, b| {
b.similarity
.partial_cmp(&a.similarity)
.unwrap_or(std::cmp::Ordering::Equal)
});
hits.truncate(top_k as usize);
return Ok(hits);
}
}
find_potential_conflicts_sql(
conn,
&scope_label,
&new_normalized,
query_vec,
active_model,
exclude_id,
top_k,
threshold,
)
}
#[allow(clippy::too_many_arguments)]
fn find_potential_conflicts_sql(
conn: &Connection,
scope_label: &str,
new_normalized: &str,
query_vec: &[f32],
active_model: &str,
exclude_id: Option<&str>,
top_k: u32,
threshold: f32,
) -> KimetsuResult<Vec<ConflictHit>> {
let mut stmt = conn.prepare(
"
SELECT memory_id, kind, text, normalized_text, embedding
FROM memories
WHERE scope = ?1
AND invalidated_at IS NULL
AND embedding IS NOT NULL
AND embedding_model = ?2
",
)?;
let rows = stmt.query_map(params![scope_label, active_model], |row| {
Ok((
row.get::<_, String>(0)?,
row.get::<_, String>(1)?,
row.get::<_, String>(2)?,
row.get::<_, String>(3)?,
row.get::<_, Vec<u8>>(4)?,
))
})?;
let mut hits: Vec<ConflictHit> = Vec::new();
for row in rows {
let (existing_id, kind, text, normalized, bytes) = row?;
if normalized == new_normalized {
continue;
}
if let Some(excl) = exclude_id {
if existing_id == excl {
continue;
}
}
let Ok(existing_vec) = decode_embedding(&bytes, Some(query_vec.len())) else {
continue;
};
let sim = cosine_similarity(query_vec, &existing_vec);
if sim >= threshold {
hits.push(ConflictHit {
existing_memory_id: existing_id,
existing_kind: kind,
existing_text: text,
similarity: sim,
});
}
}
hits.sort_by(|a, b| {
b.similarity
.partial_cmp(&a.similarity)
.unwrap_or(std::cmp::Ordering::Equal)
});
hits.truncate(top_k as usize);
Ok(hits)
}
pub fn record_conflict(
conn: &Connection,
new_memory_id: &str,
scope: &MemoryScope,
kind: &str,
hit: &ConflictHit,
) -> KimetsuResult<String> {
let existing: Option<String> = conn
.query_row(
"
SELECT conflict_id
FROM memory_conflicts
WHERE new_memory_id = ?1 AND existing_memory_id = ?2
",
params![new_memory_id, hit.existing_memory_id],
|row| row.get::<_, String>(0),
)
.optional()?;
if let Some(id) = existing {
return Ok(id);
}
let conflict_id = new_id().to_string();
let detected_at = OffsetDateTime::now_utc()
.format(&time::format_description::well_known::Rfc3339)
.map_err(|e| format!("timestamp format: {e}"))?;
conn.execute(
"
INSERT INTO memory_conflicts (
conflict_id, new_memory_id, existing_memory_id,
scope, kind, similarity, detected_at
)
VALUES (?1, ?2, ?3, ?4, ?5, ?6, ?7)
",
params![
conflict_id,
new_memory_id,
hit.existing_memory_id,
scope.to_string(),
kind,
hit.similarity as f64,
detected_at,
],
)?;
Ok(conflict_id)
}
pub fn detect_and_record(
conn: &Connection,
new_memory_id: &str,
scope: &MemoryScope,
kind: &str,
text: &str,
embedder: &dyn Embedder,
) -> usize {
detect_and_record_with_vec(conn, new_memory_id, scope, kind, text, None, embedder)
}
pub(crate) fn detect_and_record_with_vec(
conn: &Connection,
new_memory_id: &str,
scope: &MemoryScope,
kind: &str,
text: &str,
precomputed_vec: Option<&[f32]>,
embedder: &dyn Embedder,
) -> usize {
let hits = match find_potential_conflicts_with_vec(
conn,
scope,
text,
precomputed_vec,
embedder,
Some(new_memory_id),
DEFAULT_TOP_K,
DEFAULT_CONFLICT_THRESHOLD,
) {
Ok(h) => h,
Err(e) => {
eprintln!("kimetsu-brain: conflict scan skipped: {e}");
return 0;
}
};
let mut recorded = 0usize;
for hit in &hits {
match record_conflict(conn, new_memory_id, scope, kind, hit) {
Ok(_) => recorded += 1,
Err(e) => {
eprintln!(
"kimetsu-brain: failed to record conflict {} <-> {}: {e}",
new_memory_id, hit.existing_memory_id
);
}
}
}
recorded
}
#[allow(clippy::too_many_arguments)]
pub(crate) fn detect_record_and_resolve_with_vec(
conn: &Connection,
new_memory_id: &str,
scope: &MemoryScope,
kind: &str,
text: &str,
precomputed_vec: Option<&[f32]>,
embedder: &dyn Embedder,
new_confidence: f32,
new_created_at: &str,
) -> (usize, usize) {
let hits = match find_potential_conflicts_with_vec(
conn,
scope,
text,
precomputed_vec,
embedder,
Some(new_memory_id),
DEFAULT_TOP_K,
DEFAULT_CONFLICT_THRESHOLD,
) {
Ok(h) => h,
Err(e) => {
eprintln!("kimetsu-brain: conflict scan skipped: {e}");
return (0, 0);
}
};
let _ = (new_confidence, new_created_at);
let mut queued = 0;
for hit in &hits {
match record_conflict(conn, new_memory_id, scope, kind, hit) {
Ok(_) => queued += 1,
Err(e) => eprintln!("kimetsu-brain: could not queue related claims: {e}"),
}
}
(0, queued)
}
pub fn list_unresolved_conflicts(
conn: &Connection,
limit: u32,
) -> KimetsuResult<Vec<ConflictReport>> {
let mut stmt = conn.prepare(
"
SELECT c.conflict_id, c.new_memory_id, mn.text, c.existing_memory_id,
me.text, c.scope, c.kind, c.similarity, c.detected_at,
c.resolved_at, c.resolution
FROM memory_conflicts c
LEFT JOIN memories mn ON mn.memory_id = c.new_memory_id
LEFT JOIN memories me ON me.memory_id = c.existing_memory_id
WHERE c.resolved_at IS NULL
ORDER BY c.detected_at DESC
LIMIT ?1
",
)?;
let rows = stmt.query_map(params![limit], |row| {
Ok(ConflictReport {
conflict_id: row.get(0)?,
new_memory_id: row.get(1)?,
new_text: row.get::<_, Option<String>>(2)?.unwrap_or_default(),
existing_memory_id: row.get(3)?,
existing_text: row.get::<_, Option<String>>(4)?.unwrap_or_default(),
scope: row.get(5)?,
kind: row.get(6)?,
similarity: row.get::<_, f64>(7)? as f32,
detected_at: row.get(8)?,
resolved_at: row.get(9)?,
resolution: row.get(10)?,
})
})?;
let mut out = Vec::new();
for row in rows {
out.push(row?);
}
Ok(out)
}
pub fn resolve_conflict(
conn: &Connection,
conflict_id: &str,
resolution: &str,
) -> KimetsuResult<bool> {
let resolution = resolution.trim();
if !matches!(resolution, "kept_new" | "kept_existing" | "kept_both") {
return Err(format!(
"invalid conflict resolution {resolution:?}; expected kept_new | kept_existing | kept_both"
)
.into());
}
let mut changed = false;
crate::projector::with_write_txn(conn, |conn| {
let metadata: Option<(String,String,String,String,f64,String)> = conn.query_row(
"SELECT new_memory_id,existing_memory_id,scope,kind,similarity,detected_at FROM memory_conflicts WHERE conflict_id=?1 AND resolved_at IS NULL",
[conflict_id], |r|Ok((r.get(0)?,r.get(1)?,r.get(2)?,r.get(3)?,r.get(4)?,r.get(5)?))).optional()?;
let Some((new_id, existing_id, scope, kind, similarity, detected_at)) = metadata else {
return Ok(());
};
let event = kimetsu_core::event::Event::new(
kimetsu_core::ids::RunId::new(),
"conflict.resolved",
serde_json::json!({
"conflict_id":conflict_id,"new_memory_id":new_id,"existing_memory_id":existing_id,
"scope":scope,"kind":kind,"similarity":similarity,"detected_at":detected_at,"resolution":resolution
}),
);
crate::projector::apply_event(conn, &event)?;
changed = true;
Ok(())
})?;
Ok(changed)
}
pub(crate) fn project_resolution(
conn: &Connection,
event: &kimetsu_core::event::Event,
) -> KimetsuResult<()> {
let field = |key| {
event
.payload
.get(key)
.and_then(|v| v.as_str())
.ok_or_else(|| format!("conflict.resolved missing {key}"))
};
let id = field("conflict_id")?;
let new_id = field("new_memory_id")?;
let existing_id = field("existing_memory_id")?;
let resolution = field("resolution")?;
if new_id == existing_id || !matches!(resolution, "kept_new" | "kept_existing" | "kept_both") {
return Err("invalid conflict pair or resolution".into());
}
let pair: Option<(String, String)> = conn
.query_row(
"SELECT new_memory_id,existing_memory_id FROM memory_conflicts WHERE conflict_id=?1",
[id],
|r| Ok((r.get(0)?, r.get(1)?)),
)
.optional()?;
if pair.is_some_and(|(a, b)| a != new_id || b != existing_id) {
return Err("conflict pair mismatch".into());
}
let ts = event
.ts
.format(&time::format_description::well_known::Rfc3339)?;
conn.execute("INSERT OR IGNORE INTO memory_conflicts(conflict_id,new_memory_id,existing_memory_id,scope,kind,similarity,detected_at) VALUES(?1,?2,?3,?4,?5,?6,?7)",
params![id,new_id,existing_id,field("scope")?,field("kind")?,event.payload["similarity"].as_f64().unwrap_or(0.0),field("detected_at")?])?;
let loser = match resolution {
"kept_new" => Some(existing_id),
"kept_existing" => Some(new_id),
_ => None,
};
if let Some(loser) = loser {
conn.execute("UPDATE memories SET invalidated_at=COALESCE(invalidated_at,?2),invalidated_reason=CASE WHEN invalidated_reason IS NULL OR invalidated_reason IN ('forgotten','forgotten/archived','forgotten_archived') THEN ?3 ELSE invalidated_reason END WHERE memory_id=?1",params![loser,ts,format!("conflict {id} resolved as {resolution}")])?;
conn.execute("DELETE FROM memories_fts WHERE memory_id=?1", [loser])?;
#[cfg(feature = "embeddings")]
crate::ann::on_invalidate(conn, loser);
}
conn.execute(
"UPDATE memory_conflicts SET resolved_at=?2,resolution=?3 WHERE conflict_id=?1",
params![id, ts, resolution],
)?;
Ok(())
}
#[cfg(test)]
mod tests {
use super::*;
use crate::embeddings::{NoopEmbedder, StubEmbedder, encode_embedding};
use kimetsu_core::memory::normalize_memory_text;
use rusqlite::Connection;
fn open_test_brain() -> Connection {
let conn = Connection::open_in_memory().expect("open in-memory");
crate::schema::initialize(&conn).expect("init schema");
conn
}
fn insert_memory(
conn: &Connection,
memory_id: &str,
scope: &str,
kind: &str,
text: &str,
embedder: &dyn Embedder,
) {
let normalized = normalize_memory_text(text);
let vec = embedder.embed(text).expect("embed test row");
let blob = encode_embedding(&vec);
conn.execute(
"
INSERT INTO memories (
memory_id, scope, kind, text, normalized_text, confidence,
source_event_id, provenance_snapshot_json, created_at,
use_count, usefulness_score, embedding, embedding_model
)
VALUES (?1, ?2, ?3, ?4, ?5, 1.0, NULL, '{}',
'2026-01-01T00:00:00Z', 0, 0.0, ?6, ?7)
",
params![
memory_id,
scope,
kind,
text,
normalized,
blob,
embedder.model_id(),
],
)
.expect("insert");
conn.execute(
"INSERT INTO memories_fts (memory_id, text, kind, scope)
VALUES (?1, ?2, ?3, ?4)",
params![memory_id, text, kind, scope],
)
.expect("fts");
}
#[test]
fn noop_embedder_returns_no_conflicts() {
let conn = open_test_brain();
let stub = StubEmbedder::new();
insert_memory(
&conn,
"m_existing",
"global_user",
"fact",
"use thiserror for libraries",
&stub,
);
let hits = find_potential_conflicts(
&conn,
&MemoryScope::GlobalUser,
"use anyhow for libraries",
&NoopEmbedder,
DEFAULT_TOP_K,
DEFAULT_CONFLICT_THRESHOLD,
)
.expect("scan");
assert!(hits.is_empty(), "noop embedder should produce no hits");
}
#[test]
fn cross_model_rows_are_skipped() {
let conn = open_test_brain();
let stub = StubEmbedder::new();
insert_memory(
&conn,
"m_xmodel",
"global_user",
"fact",
"use thiserror",
&stub,
);
conn.execute(
"UPDATE memories SET embedding_model = 'bge-small-en-v1.5' WHERE memory_id = 'm_xmodel'",
[],
)
.expect("force mismatch");
let hits = find_potential_conflicts(
&conn,
&MemoryScope::GlobalUser,
"use thiserror everywhere", &stub,
DEFAULT_TOP_K,
0.0,
)
.expect("scan");
assert!(
hits.is_empty(),
"cross-model rows must be skipped from conflict scan"
);
}
#[test]
fn exact_match_is_not_flagged_as_conflict() {
let conn = open_test_brain();
let stub = StubEmbedder::new();
insert_memory(
&conn,
"m_exact",
"global_user",
"fact",
"Use ripgrep",
&stub,
);
let hits = find_potential_conflicts(
&conn,
&MemoryScope::GlobalUser,
"use ripgrep",
&stub,
DEFAULT_TOP_K,
0.0, )
.expect("scan");
assert!(
hits.is_empty(),
"exact normalized-text match should be dedup, not conflict"
);
}
#[test]
fn similar_but_different_text_is_flagged() {
let conn = open_test_brain();
let stub = StubEmbedder::new();
insert_memory(
&conn,
"m_existing",
"global_user",
"fact",
"alpha beta gamma delta",
&stub,
);
let hits = find_potential_conflicts(
&conn,
&MemoryScope::GlobalUser,
"alpha beta gamma omega", &stub,
DEFAULT_TOP_K,
0.4,
)
.expect("scan");
assert!(
!hits.is_empty(),
"high-cosine + different-normalized text should flag a conflict"
);
assert_eq!(hits[0].existing_memory_id, "m_existing");
assert!(
hits[0].similarity >= 0.4,
"similarity should be >= threshold; got {}",
hits[0].similarity
);
}
#[test]
fn record_conflict_is_idempotent() {
let conn = open_test_brain();
let stub = StubEmbedder::new();
insert_memory(&conn, "m_new", "global_user", "fact", "alpha", &stub);
insert_memory(&conn, "m_old", "global_user", "fact", "beta", &stub);
let hit = ConflictHit {
existing_memory_id: "m_old".to_string(),
existing_kind: "fact".to_string(),
existing_text: "beta".to_string(),
similarity: 0.85,
};
let id1 = record_conflict(&conn, "m_new", &MemoryScope::GlobalUser, "fact", &hit)
.expect("record 1");
let id2 = record_conflict(&conn, "m_new", &MemoryScope::GlobalUser, "fact", &hit)
.expect("record 2");
assert_eq!(id1, id2, "re-recording the same pair must return same id");
let count: i64 = conn
.query_row("SELECT COUNT(*) FROM memory_conflicts", [], |row| {
row.get(0)
})
.unwrap();
assert_eq!(count, 1);
}
#[test]
fn list_unresolved_excludes_resolved_rows() {
let conn = open_test_brain();
let stub = StubEmbedder::new();
insert_memory(
&conn,
"m_new1",
"global_user",
"fact",
"use thiserror",
&stub,
);
insert_memory(&conn, "m_old1", "global_user", "fact", "use anyhow", &stub);
insert_memory(
&conn,
"m_new2",
"global_user",
"fact",
"tabs over spaces",
&stub,
);
insert_memory(
&conn,
"m_old2",
"global_user",
"fact",
"spaces over tabs",
&stub,
);
let hit1 = ConflictHit {
existing_memory_id: "m_old1".to_string(),
existing_kind: "fact".to_string(),
existing_text: "use anyhow".to_string(),
similarity: 0.9,
};
let hit2 = ConflictHit {
existing_memory_id: "m_old2".to_string(),
existing_kind: "fact".to_string(),
existing_text: "spaces over tabs".to_string(),
similarity: 0.85,
};
let cid1 =
record_conflict(&conn, "m_new1", &MemoryScope::GlobalUser, "fact", &hit1).unwrap();
let _cid2 =
record_conflict(&conn, "m_new2", &MemoryScope::GlobalUser, "fact", &hit2).unwrap();
assert!(resolve_conflict(&conn, &cid1, "kept_both").unwrap());
let open = list_unresolved_conflicts(&conn, 50).unwrap();
assert_eq!(open.len(), 1, "only the unresolved conflict should list");
assert_eq!(open[0].new_memory_id, "m_new2");
assert_eq!(open[0].existing_memory_id, "m_old2");
assert_eq!(open[0].new_text, "tabs over spaces");
assert_eq!(open[0].existing_text, "spaces over tabs");
}
#[test]
fn resolve_conflict_invalidates_loser_side() {
let conn = open_test_brain();
let stub = StubEmbedder::new();
for (mid, text) in [
("m_keep_new", "alpha"),
("m_old_loses", "beta"),
("m_new_loses", "gamma"),
("m_keep_existing", "delta"),
("m_both_a", "epsilon"),
("m_both_b", "zeta"),
] {
insert_memory(&conn, mid, "global_user", "fact", text, &stub);
}
let mk_hit = |old: &str| ConflictHit {
existing_memory_id: old.to_string(),
existing_kind: "fact".to_string(),
existing_text: "x".to_string(),
similarity: 0.9,
};
let c_kept_new = record_conflict(
&conn,
"m_keep_new",
&MemoryScope::GlobalUser,
"fact",
&mk_hit("m_old_loses"),
)
.unwrap();
let c_kept_existing = record_conflict(
&conn,
"m_new_loses",
&MemoryScope::GlobalUser,
"fact",
&mk_hit("m_keep_existing"),
)
.unwrap();
let c_both = record_conflict(
&conn,
"m_both_a",
&MemoryScope::GlobalUser,
"fact",
&mk_hit("m_both_b"),
)
.unwrap();
assert!(resolve_conflict(&conn, &c_kept_new, "kept_new").unwrap());
assert!(resolve_conflict(&conn, &c_kept_existing, "kept_existing").unwrap());
assert!(resolve_conflict(&conn, &c_both, "kept_both").unwrap());
let invalidated_at: Vec<(String, Option<String>)> = {
let mut stmt = conn
.prepare("SELECT memory_id, invalidated_at FROM memories ORDER BY memory_id")
.unwrap();
stmt.query_map([], |row| {
Ok((row.get::<_, String>(0)?, row.get::<_, Option<String>>(1)?))
})
.unwrap()
.map(|r| r.unwrap())
.collect()
};
let map: std::collections::HashMap<_, _> = invalidated_at.into_iter().collect();
assert!(map["m_keep_new"].is_none(), "winner should stay active");
assert!(
map["m_old_loses"].is_some(),
"kept_new must invalidate the existing memory"
);
assert!(
map["m_keep_existing"].is_none(),
"winner (existing) should stay active"
);
assert!(
map["m_new_loses"].is_some(),
"kept_existing must invalidate the new memory"
);
assert!(
map["m_both_a"].is_none() && map["m_both_b"].is_none(),
"kept_both should leave both memories active"
);
}
#[test]
fn resolve_conflict_is_idempotent() {
let conn = open_test_brain();
let stub = StubEmbedder::new();
insert_memory(&conn, "m_new", "global_user", "fact", "x", &stub);
insert_memory(&conn, "m_old", "global_user", "fact", "y", &stub);
let hit = ConflictHit {
existing_memory_id: "m_old".to_string(),
existing_kind: "fact".to_string(),
existing_text: "y".to_string(),
similarity: 0.95,
};
let cid = record_conflict(&conn, "m_new", &MemoryScope::GlobalUser, "fact", &hit).unwrap();
assert!(resolve_conflict(&conn, &cid, "kept_new").unwrap());
assert!(
!resolve_conflict(&conn, &cid, "kept_existing").unwrap(),
"second resolve must return false (already resolved)"
);
}
#[test]
fn detect_and_record_noop_writes_nothing() {
let conn = open_test_brain();
let stub = StubEmbedder::new();
insert_memory(
&conn,
"m_existing",
"global_user",
"fact",
"alpha beta",
&stub,
);
insert_memory(&conn, "m_new", "global_user", "fact", "alpha gamma", &stub);
let recorded = detect_and_record(
&conn,
"m_new",
&MemoryScope::GlobalUser,
"fact",
"alpha gamma",
&NoopEmbedder,
);
assert_eq!(recorded, 0);
let count: i64 = conn
.query_row("SELECT COUNT(*) FROM memory_conflicts", [], |row| {
row.get(0)
})
.unwrap();
assert_eq!(count, 0);
}
#[test]
fn resolve_conflict_rejects_invalid_resolution_strings() {
let conn = open_test_brain();
let err = resolve_conflict(&conn, "ignored", "delete_them_all").unwrap_err();
let msg = format!("{err}");
assert!(msg.contains("invalid conflict resolution"), "got: {msg}");
}
#[test]
fn conflict_detection_enabled_env_disable_overrides_config_true() {
let lock = crate::user_brain::test_env_lock()
.lock()
.unwrap_or_else(|p| p.into_inner());
let prev = std::env::var("KIMETSU_DETECT_CONFLICTS").ok();
for v in ["0", "false", "off", "no"] {
unsafe {
std::env::set_var("KIMETSU_DETECT_CONFLICTS", v);
}
assert!(
!conflict_detection_enabled(true),
"env={v:?} must disable even when config=true"
);
}
unsafe {
match prev {
Some(v) => std::env::set_var("KIMETSU_DETECT_CONFLICTS", v),
None => std::env::remove_var("KIMETSU_DETECT_CONFLICTS"),
}
}
drop(lock);
}
#[test]
fn conflict_detection_enabled_config_false_when_env_unset() {
let lock = crate::user_brain::test_env_lock()
.lock()
.unwrap_or_else(|p| p.into_inner());
let prev = std::env::var("KIMETSU_DETECT_CONFLICTS").ok();
unsafe {
std::env::remove_var("KIMETSU_DETECT_CONFLICTS");
}
assert!(
!conflict_detection_enabled(false),
"config=false + env unset must be disabled"
);
assert!(
conflict_detection_enabled(true),
"config=true + env unset must be enabled"
);
unsafe {
match prev {
Some(v) => std::env::set_var("KIMETSU_DETECT_CONFLICTS", v),
None => std::env::remove_var("KIMETSU_DETECT_CONFLICTS"),
}
}
drop(lock);
}
#[test]
fn off_switch_prevents_conflict_detection() {
let conn = open_test_brain();
let stub = StubEmbedder::new();
insert_memory(
&conn,
"m_seed",
"global_user",
"fact",
"alpha beta gamma delta",
&stub,
);
let lock = crate::user_brain::test_env_lock()
.lock()
.unwrap_or_else(|p| p.into_inner());
let prev = std::env::var("KIMETSU_DETECT_CONFLICTS").ok();
unsafe {
std::env::remove_var("KIMETSU_DETECT_CONFLICTS");
}
if conflict_detection_enabled(false) {
panic!("detect_conflicts=false must disable the gate");
}
let count: i64 = conn
.query_row("SELECT COUNT(*) FROM memory_conflicts", [], |row| {
row.get(0)
})
.unwrap();
assert_eq!(count, 0, "off-switch must prevent any conflict writes");
let hits = find_potential_conflicts(
&conn,
&MemoryScope::GlobalUser,
"alpha beta gamma omega",
&stub,
DEFAULT_TOP_K,
0.4,
)
.expect("scan");
assert!(
!hits.is_empty(),
"when enabled, near-dup must be detected (test sanity check)"
);
unsafe {
match prev {
Some(v) => std::env::set_var("KIMETSU_DETECT_CONFLICTS", v),
None => std::env::remove_var("KIMETSU_DETECT_CONFLICTS"),
}
}
drop(lock);
}
#[test]
fn exclude_id_prevents_self_conflict() {
let conn = open_test_brain();
let stub = StubEmbedder::new();
insert_memory(
&conn,
"m_self",
"global_user",
"fact",
"alpha beta gamma delta",
&stub,
);
let hits = find_potential_conflicts_with_vec(
&conn,
&MemoryScope::GlobalUser,
"alpha beta gamma delta",
None,
&stub,
Some("m_self"),
DEFAULT_TOP_K,
0.0, )
.expect("scan");
assert!(
hits.is_empty(),
"excluded memory must not appear as a conflict hit"
);
}
#[allow(clippy::too_many_arguments)]
fn insert_memory_with_meta(
conn: &Connection,
memory_id: &str,
scope: &str,
kind: &str,
text: &str,
confidence: f32,
created_at: &str,
embedder: &dyn Embedder,
) {
let normalized = normalize_memory_text(text);
let vec = embedder.embed(text).expect("embed test row");
let blob = encode_embedding(&vec);
conn.execute(
"INSERT INTO memories (
memory_id, scope, kind, text, normalized_text, confidence,
source_event_id, provenance_snapshot_json, created_at,
use_count, usefulness_score, embedding, embedding_model
)
VALUES (?1, ?2, ?3, ?4, ?5, ?6, NULL, '{}', ?7, 0, 0.0, ?8, ?9)",
rusqlite::params![
memory_id,
scope,
kind,
text,
normalized,
confidence as f64,
created_at,
blob,
embedder.model_id(),
],
)
.expect("insert");
conn.execute(
"INSERT INTO memories_fts (memory_id, text, kind, scope) VALUES (?1, ?2, ?3, ?4)",
rusqlite::params![memory_id, text, kind, scope],
)
.expect("fts");
}
#[test]
fn resolution_score_higher_confidence_wins_all_else_equal() {
let now_str = time::OffsetDateTime::now_utc()
.format(&time::format_description::well_known::Rfc3339)
.unwrap();
let score_high = resolution_score(0.9, &now_str);
let score_low = resolution_score(0.5, &now_str);
assert!(
score_high > score_low,
"higher confidence must produce higher score; got {score_high} vs {score_low}"
);
}
#[test]
fn resolution_score_newer_wins_all_else_equal() {
let now_str = time::OffsetDateTime::now_utc()
.format(&time::format_description::well_known::Rfc3339)
.unwrap();
let old_ts = (time::OffsetDateTime::now_utc() - time::Duration::days(90))
.format(&time::format_description::well_known::Rfc3339)
.unwrap();
let score_new = resolution_score(0.8, &now_str);
let score_old = resolution_score(0.8, &old_ts);
assert!(
score_new > score_old,
"newer memory must score higher; got new={score_new} old={score_old}"
);
}
#[test]
fn auto_resolution_stamps_loser_valid_to_when_new_wins() {
let conn = open_test_brain();
let stub = StubEmbedder::new();
let old_ts = "2020-01-01T00:00:00Z";
insert_memory_with_meta(
&conn,
"m_loser",
"global_user",
"fact",
"alpha beta gamma delta",
0.3, old_ts,
&stub,
);
let now_str = time::OffsetDateTime::now_utc()
.format(&time::format_description::well_known::Rfc3339)
.unwrap();
insert_memory_with_meta(
&conn,
"m_winner",
"global_user",
"fact",
"alpha beta gamma omega",
0.95, &now_str,
&stub,
);
let new_score = resolution_score(0.95, &now_str);
let existing_score = resolution_score(0.3, old_ts);
assert!(
new_score > existing_score,
"new high-confidence must score higher; got new={new_score} existing={existing_score}"
);
let delta = (new_score - existing_score).abs();
assert!(
delta >= NEAR_TIE_BAND,
"gap {delta} must exceed NEAR_TIE_BAND for auto-resolution"
);
crate::projector::mark_memory_temporal(&conn, "m_loser", None, Some(&now_str))
.expect("mark valid_to on loser");
let loser_vt: Option<String> = conn
.query_row(
"SELECT valid_to FROM memories WHERE memory_id = 'm_loser'",
[],
|r| r.get(0),
)
.unwrap();
assert!(loser_vt.is_some(), "loser must have valid_to stamped");
let winner_vt: Option<String> = conn
.query_row(
"SELECT valid_to FROM memories WHERE memory_id = 'm_winner'",
[],
|r| r.get(0),
)
.unwrap();
assert!(winner_vt.is_none(), "winner must NOT have valid_to");
}
#[test]
fn auto_resolution_stamps_new_memory_when_existing_wins() {
let conn = open_test_brain();
let stub = StubEmbedder::new();
let now_str = time::OffsetDateTime::now_utc()
.format(&time::format_description::well_known::Rfc3339)
.unwrap();
insert_memory_with_meta(
&conn,
"m_existing_winner",
"global_user",
"fact",
"alpha beta gamma delta",
0.95, &now_str,
&stub,
);
let old_ts = "2020-01-01T00:00:00Z";
insert_memory_with_meta(
&conn,
"m_new_loser",
"global_user",
"fact",
"alpha beta gamma omega",
0.2, old_ts,
&stub,
);
let existing_score = resolution_score(0.95, &now_str);
let new_score = resolution_score(0.2, old_ts);
assert!(
existing_score > new_score,
"existing high-confidence must score higher; existing={existing_score} new={new_score}"
);
let delta = (existing_score - new_score).abs();
assert!(
delta >= NEAR_TIE_BAND,
"gap {delta} must exceed NEAR_TIE_BAND"
);
crate::projector::mark_memory_temporal(&conn, "m_new_loser", None, Some(&now_str))
.expect("mark valid_to on new loser");
let new_vt: Option<String> = conn
.query_row(
"SELECT valid_to FROM memories WHERE memory_id = 'm_new_loser'",
[],
|r| r.get(0),
)
.unwrap();
assert!(new_vt.is_some(), "new loser must have valid_to stamped");
let existing_vt: Option<String> = conn
.query_row(
"SELECT valid_to FROM memories WHERE memory_id = 'm_existing_winner'",
[],
|r| r.get(0),
)
.unwrap();
assert!(
existing_vt.is_none(),
"existing winner must NOT have valid_to"
);
}
#[test]
fn near_tie_goes_to_queue_not_auto_resolved() {
let conn = open_test_brain();
let stub = StubEmbedder::new();
let now_str = time::OffsetDateTime::now_utc()
.format(&time::format_description::well_known::Rfc3339)
.unwrap();
insert_memory_with_meta(
&conn,
"m_tie_existing",
"global_user",
"fact",
"alpha beta gamma delta",
0.8,
&now_str,
&stub,
);
insert_memory_with_meta(
&conn,
"m_tie_new",
"global_user",
"fact",
"alpha beta gamma omega",
0.8,
&now_str,
&stub,
);
let (auto_resolved, queued) = detect_record_and_resolve_with_vec(
&conn,
"m_tie_new",
&MemoryScope::GlobalUser,
"fact",
"alpha beta gamma omega",
None,
&stub,
0.8,
&now_str,
);
assert_eq!(
auto_resolved, 0,
"near-tie must NOT be auto-resolved (got {auto_resolved} auto-resolved)"
);
let existing_vt: Option<String> = conn
.query_row(
"SELECT valid_to FROM memories WHERE memory_id = 'm_tie_existing'",
[],
|r| r.get(0),
)
.unwrap();
let new_vt: Option<String> = conn
.query_row(
"SELECT valid_to FROM memories WHERE memory_id = 'm_tie_new'",
[],
|r| r.get(0),
)
.unwrap();
assert!(
existing_vt.is_none(),
"near-tie existing memory must NOT be stamped; got {existing_vt:?}"
);
assert!(
new_vt.is_none(),
"near-tie new memory must NOT be stamped; got {new_vt:?}"
);
if queued > 0 {
let count: i64 = conn
.query_row(
"SELECT COUNT(*) FROM memory_conflicts WHERE resolved_at IS NULL",
[],
|r| r.get(0),
)
.unwrap();
assert!(
count > 0,
"near-tie must add unresolved row to memory_conflicts"
);
}
}
#[test]
fn high_similarity_and_score_gap_do_not_prove_contradiction() {
let conn = open_test_brain();
let stub = StubEmbedder::new();
let old = "The development service uses SQLite.";
let new = "The production service uses SQLite.";
let now = OffsetDateTime::now_utc().format(&Rfc3339).unwrap();
insert_memory_with_meta(
&conn,
"old",
"global_user",
"fact",
old,
0.1,
"2020-01-01T00:00:00Z",
&stub,
);
insert_memory_with_meta(&conn, "new", "global_user", "fact", new, 1.0, &now, &stub);
let vector = stub.embed(old).unwrap();
let (resolved, queued) = detect_record_and_resolve_with_vec(
&conn,
"new",
&MemoryScope::GlobalUser,
"fact",
new,
Some(&vector),
&stub,
1.0,
&now,
);
assert_eq!(
resolved, 0,
"no automatic retirement without a proven conflicting claim"
);
assert_eq!(queued, 1, "related claims remain reviewable");
let retired: i64 = conn
.query_row(
"SELECT COUNT(*) FROM memories WHERE valid_to IS NOT NULL",
[],
|r| r.get(0),
)
.unwrap();
assert_eq!(retired, 0);
}
#[test]
fn auto_resolution_survives_rebuild() {
let conn = open_test_brain();
let stub = StubEmbedder::new();
let old_ts = "2020-01-01T00:00:00Z";
insert_memory_with_meta(
&conn,
"m_rebuild_old",
"global_user",
"fact",
"alpha beta gamma delta",
0.2,
old_ts,
&stub,
);
let now_str = time::OffsetDateTime::now_utc()
.format(&time::format_description::well_known::Rfc3339)
.unwrap();
insert_memory_with_meta(
&conn,
"m_rebuild_new",
"global_user",
"fact",
"alpha beta gamma omega",
0.95,
&now_str,
&stub,
);
let (auto_resolved, _queued) = detect_record_and_resolve_with_vec(
&conn,
"m_rebuild_new",
&MemoryScope::GlobalUser,
"fact",
"alpha beta gamma omega",
None,
&stub,
0.95,
&now_str,
);
if auto_resolved == 0 {
return;
}
let vt_before: Option<String> = conn
.query_row(
"SELECT valid_to FROM memories WHERE memory_id = 'm_rebuild_old'",
[],
|r| r.get(0),
)
.unwrap();
assert!(
vt_before.is_some(),
"loser must have valid_to before rebuild"
);
crate::projector::rebuild_in_place(&conn).expect("rebuild_in_place");
let vt_after: Option<String> = conn
.query_row(
"SELECT valid_to FROM memories WHERE memory_id = 'm_rebuild_old'",
[],
|r| r.get(0),
)
.unwrap();
assert!(
vt_after.is_some(),
"loser's valid_to must survive rebuild_in_place"
);
}
#[test]
fn resolve_conflicts_enabled_env_disable_overrides_config_true() {
let lock = crate::user_brain::test_env_lock()
.lock()
.unwrap_or_else(|p| p.into_inner());
let prev = std::env::var("KIMETSU_RESOLVE_CONFLICTS").ok();
for v in ["0", "false", "off", "no"] {
unsafe {
std::env::set_var("KIMETSU_RESOLVE_CONFLICTS", v);
}
assert!(
!resolve_conflicts_enabled(true),
"env={v:?} must disable resolution even when config=true"
);
}
unsafe {
match prev {
Some(v) => std::env::set_var("KIMETSU_RESOLVE_CONFLICTS", v),
None => std::env::remove_var("KIMETSU_RESOLVE_CONFLICTS"),
}
}
drop(lock);
}
}