use std::borrow::Cow;
use std::str::FromStr;
use std::time::Duration;
use kimetsu_core::KimetsuResult;
use kimetsu_core::event::Event;
use kimetsu_core::ids::{EventId, RunId};
use rusqlite::{Connection, OptionalExtension, params};
use time::OffsetDateTime;
use time::format_description::well_known::Rfc3339;
use crate::redact;
use crate::schema;
const WRITE_TXN_MAX_ATTEMPTS: u32 = 5;
fn is_sqlite_busy(err: &(dyn std::error::Error + 'static)) -> bool {
err.downcast_ref::<rusqlite::Error>()
.and_then(|e| e.sqlite_error_code())
.is_some_and(|code| {
matches!(
code,
rusqlite::ErrorCode::DatabaseBusy | rusqlite::ErrorCode::DatabaseLocked
)
})
}
fn with_write_txn<F>(conn: &Connection, mut body: F) -> KimetsuResult<()>
where
F: FnMut(&Connection) -> KimetsuResult<()>,
{
let mut attempt = 0u32;
loop {
attempt += 1;
if let Err(e) = conn.execute_batch("BEGIN IMMEDIATE") {
let boxed: Box<dyn std::error::Error + Send + Sync> = e.into();
if is_sqlite_busy(boxed.as_ref()) && attempt < WRITE_TXN_MAX_ATTEMPTS {
std::thread::sleep(Duration::from_millis(20 * attempt as u64));
continue;
}
return Err(boxed);
}
match body(conn) {
Ok(()) => match conn.execute_batch("COMMIT") {
Ok(()) => return Ok(()),
Err(e) => {
let _ = conn.execute_batch("ROLLBACK");
return Err(e.into());
}
},
Err(e) => {
let _ = conn.execute_batch("ROLLBACK");
return Err(e);
}
}
}
}
fn upcast_event(event: &Event) -> Cow<'_, Event> {
Cow::Borrowed(event)
}
pub fn rebuild(conn: &Connection, events: &[Event]) -> KimetsuResult<()> {
reset_projection(conn)?;
apply_events(conn, events)
}
pub fn rebuild_in_place(conn: &Connection) -> KimetsuResult<usize> {
let events = read_events_ordered(conn)?;
with_write_txn(conn, |c| {
reset_projection(c)?;
for event in &events {
project_event(c, event)?;
}
Ok(())
})?;
Ok(events.len())
}
fn read_events_ordered(conn: &Connection) -> KimetsuResult<Vec<Event>> {
let mut stmt = conn.prepare(
"
SELECT event_id, run_id, ts, kind, schema_version, payload_json, origin, hlc
FROM events
ORDER BY hlc, rowid
",
)?;
let rows = stmt.query_map([], |row| {
let event_id_str: String = row.get(0)?;
let run_id_str: String = row.get(1)?;
let ts_str: String = row.get(2)?;
let kind: String = row.get(3)?;
let schema_version: u32 = row.get(4)?;
let payload_json: String = row.get(5)?;
let origin: Option<String> = row.get(6)?;
let hlc: Option<String> = row.get(7)?;
Ok((
event_id_str,
run_id_str,
ts_str,
kind,
schema_version,
payload_json,
origin,
hlc,
))
})?;
let mut events = Vec::new();
for row in rows {
let (event_id_str, run_id_str, ts_str, kind, schema_version, payload_json, origin, hlc) =
row?;
let event_id = EventId(
ulid::Ulid::from_str(&event_id_str)
.map_err(|e| format!("invalid event_id {event_id_str:?}: {e}"))?,
);
let run_id = RunId(
ulid::Ulid::from_str(&run_id_str)
.map_err(|e| format!("invalid run_id {run_id_str:?}: {e}"))?,
);
let ts = OffsetDateTime::parse(&ts_str, &Rfc3339)
.map_err(|e| format!("invalid ts {ts_str:?}: {e}"))?;
let payload: serde_json::Value = serde_json::from_str(&payload_json)?;
events.push(Event {
event_id,
run_id,
ts,
parent_event_id: None, kind,
schema_version,
payload,
origin, hlc, });
}
Ok(events)
}
pub fn apply_events(conn: &Connection, events: &[Event]) -> KimetsuResult<()> {
with_write_txn(conn, |c| {
for event in events {
apply_event(c, event)?;
}
Ok(())
})
}
fn reset_projection(conn: &Connection) -> KimetsuResult<()> {
conn.execute_batch(
"
DELETE FROM runs;
DELETE FROM sources;
DELETE FROM memories;
DELETE FROM memory_proposals;
DELETE FROM memories_fts;
DELETE FROM memory_citations;
DELETE FROM memory_conflicts;
DELETE FROM sync_conflicts;
DELETE FROM memory_edges;
DELETE FROM work_episodes;
",
)?;
Ok(())
}
fn apply_event(conn: &Connection, event: &Event) -> KimetsuResult<()> {
let event = redact_memory_event(event);
let event = event.as_ref();
insert_event(conn, event)?;
project_event(conn, event)
}
fn project_event(conn: &Connection, event: &Event) -> KimetsuResult<()> {
let upcasted = upcast_event(event);
let redacted = redact_memory_event(upcasted.as_ref());
let event = redacted.as_ref();
match event.kind.as_str() {
"run.started" => apply_run_started(conn, event),
"run.finished" | "run.failed" | "run.aborted" => apply_terminal_run(conn, event),
"memory.accepted" => apply_memory_accepted(conn, event),
"memory.proposed" => apply_memory_proposed(conn, event),
"memory.rejected" => apply_memory_rejected(conn, event),
"memory.invalidated" => apply_memory_invalidated(conn, event),
"memory.cited" => apply_memory_cited(conn, event),
"retrieval.regret" => apply_retrieval_regret(conn, event),
"memory.aged" => apply_memory_aged(conn, event),
"memory.superseded" => apply_memory_superseded(conn, event),
"memory.edge" => apply_memory_edge(conn, event),
"memory.temporal" => apply_memory_temporal(conn, event),
"work.episode" => crate::episode::project_work_episode(conn, event),
_ => Ok(()),
}
}
fn redact_memory_event(event: &Event) -> Cow<'_, Event> {
if !matches!(
event.kind.as_str(),
"memory.accepted" | "memory.proposed" | "memory.cited"
) {
return Cow::Borrowed(event);
}
let (payload, changed) = redact_json_strings(&event.payload);
if changed {
Cow::Owned(Event {
payload,
..event.clone()
})
} else {
Cow::Borrowed(event)
}
}
fn redact_json_strings(value: &serde_json::Value) -> (serde_json::Value, bool) {
match value {
serde_json::Value::String(text) => {
let redaction = redact::redact_secrets(text);
let changed = redaction.was_redacted();
(serde_json::Value::String(redaction.text), changed)
}
serde_json::Value::Array(values) => {
let mut changed = false;
let values = values
.iter()
.map(|value| {
let (value, did_change) = redact_json_strings(value);
changed |= did_change;
value
})
.collect();
(serde_json::Value::Array(values), changed)
}
serde_json::Value::Object(map) => {
let mut changed = false;
let map = map
.iter()
.map(|(key, value)| {
let (value, did_change) = redact_json_strings(value);
changed |= did_change;
(key.clone(), value)
})
.collect();
(serde_json::Value::Object(map), changed)
}
other => (other.clone(), false),
}
}
fn apply_memory_cited(conn: &Connection, event: &Event) -> KimetsuResult<()> {
let Some(memory_id) = event
.payload
.get("memory_id")
.and_then(|value| value.as_str())
else {
return Ok(());
};
let turn = event
.payload
.get("turn")
.and_then(|value| value.as_i64())
.unwrap_or(0);
let rationale = event
.payload
.get("rationale")
.and_then(|value| value.as_str());
let cited_at = ts_text(event)?;
conn.execute(
"
INSERT OR REPLACE INTO memory_citations (
run_id, memory_id, turn, cited_at, rationale
)
VALUES (?1, ?2, ?3, ?4, ?5)
",
params![
event.run_id.to_string(),
memory_id,
turn,
cited_at,
rationale,
],
)?;
if event.run_id.0 == ulid::Ulid::nil() {
apply_cited_outcome(conn, memory_id, 1.0, 1.0, &cited_at, true)?;
}
Ok(())
}
fn apply_retrieval_regret(conn: &Connection, event: &Event) -> KimetsuResult<()> {
let is_manual = event
.payload
.get("source")
.and_then(|v| v.as_str())
.map(|s| s == "manual")
.unwrap_or(false);
if !is_manual {
return Ok(());
}
let Some(memory_id) = event.payload.get("memory_id").and_then(|v| v.as_str()) else {
return Ok(());
};
let ts = ts_text(event)?;
apply_cited_outcome(conn, memory_id, -1.0, 0.0, &ts, false)?;
Ok(())
}
fn apply_memory_aged(conn: &Connection, event: &Event) -> KimetsuResult<()> {
let Some(memory_id) = event.payload.get("memory_id").and_then(|v| v.as_str()) else {
return Ok(());
};
if let Some(created) = event.payload.get("created_at").and_then(|v| v.as_str()) {
conn.execute(
"UPDATE memories SET created_at = ?2 WHERE memory_id = ?1",
params![memory_id, created],
)?;
}
if let Some(last_useful) = event.payload.get("last_useful_at").and_then(|v| v.as_str()) {
conn.execute(
"UPDATE memories SET last_useful_at = ?2 WHERE memory_id = ?1",
params![memory_id, last_useful],
)?;
}
Ok(())
}
const CONF_ALPHA: f64 = 0.05;
fn apply_cited_outcome(
conn: &Connection,
memory_id: &str,
usefulness_delta: f64,
conf_target: f64,
ts: &str,
bump_last_useful: bool,
) -> KimetsuResult<()> {
conn.execute(
"UPDATE memories
SET use_count = use_count + 1,
usefulness_score = usefulness_score + ?2,
last_used_at = ?3
WHERE memory_id = ?1",
params![memory_id, usefulness_delta, ts],
)?;
if bump_last_useful {
conn.execute(
"UPDATE memories SET last_useful_at = ?2 WHERE memory_id = ?1",
params![memory_id, ts],
)?;
}
let old_conf: f64 = conn
.query_row(
"SELECT confidence FROM memories WHERE memory_id = ?1",
params![memory_id],
|row| row.get::<_, f64>(0),
)
.unwrap_or(1.0);
let new_conf = (old_conf + CONF_ALPHA * (conf_target - old_conf)).clamp(0.1, 0.99);
conn.execute(
"UPDATE memories SET confidence = ?2 WHERE memory_id = ?1",
params![memory_id, new_conf],
)?;
Ok(())
}
pub(crate) fn insert_event(conn: &Connection, event: &Event) -> KimetsuResult<()> {
let payload = serde_json::to_string(&event.payload)?;
conn.execute(
"
INSERT OR IGNORE INTO events (
event_id, run_id, ts, kind, schema_version, payload_json, origin, hlc
)
VALUES (?1, ?2, ?3, ?4, ?5, ?6, ?7, ?8)
",
params![
event.event_id.to_string(),
event.run_id.to_string(),
ts_text(event)?,
event.kind,
event.schema_version,
payload,
event.origin,
event.hlc,
],
)?;
Ok(())
}
fn apply_run_started(conn: &Connection, event: &Event) -> KimetsuResult<()> {
let project_id = event
.payload
.get("project_id")
.and_then(|value| value.as_str())
.unwrap_or("unknown");
let task = event
.payload
.get("task")
.and_then(|value| value.as_str())
.unwrap_or("");
let model = event
.payload
.get("model")
.and_then(|value| value.as_str())
.map(str::to_string);
conn.execute(
"
INSERT OR IGNORE INTO runs (
run_id, project_id, task, started_at, model, total_cost_usd
)
VALUES (?1, ?2, ?3, ?4, ?5, 0)
",
params![
event.run_id.to_string(),
project_id,
task,
ts_text(event)?,
model
],
)?;
Ok(())
}
fn apply_terminal_run(conn: &Connection, event: &Event) -> KimetsuResult<()> {
let total_cost = event
.payload
.get("total_cost_usd")
.and_then(|value| value.as_f64())
.unwrap_or(0.0);
conn.execute(
"
UPDATE runs
SET ended_at = ?2,
terminal_kind = ?3,
total_cost_usd = ?4
WHERE run_id = ?1
",
params![
event.run_id.to_string(),
ts_text(event)?,
event.kind,
total_cost
],
)?;
apply_memory_usefulness_for_run(conn, event)?;
Ok(())
}
fn apply_memory_usefulness_for_run(conn: &Connection, event: &Event) -> KimetsuResult<()> {
let (strong, weak): (f64, f64) = match event.kind.as_str() {
"run.finished" => (1.0, 0.1),
"run.failed" => {
let category = event
.payload
.get("category")
.and_then(|value| value.as_str())
.unwrap_or("");
if category == "Gate" {
return Ok(());
}
(-1.0, -0.1)
}
_ => return Ok(()), };
let run_id = event.run_id.to_string();
let retrieved = collect_injected_memory_ids(conn, &run_id)?;
if retrieved.is_empty() {
return Ok(());
}
let cited = collect_cited_memory_ids(conn, &run_id)?;
let ts = ts_text(event)?;
let bump_last_useful = event.kind == "run.finished";
let conf_target: Option<f64> = match event.kind.as_str() {
"run.finished" => Some(1.0),
"run.failed" => Some(0.0),
_ => None,
};
for memory_id in &retrieved {
let is_cited = cited.contains(memory_id);
let delta = if is_cited { strong } else { weak };
conn.execute(
"
UPDATE memories
SET use_count = use_count + 1,
usefulness_score = usefulness_score + ?2,
last_used_at = ?3
WHERE memory_id = ?1
",
params![memory_id, delta, ts],
)?;
if is_cited && bump_last_useful {
conn.execute(
"UPDATE memories SET last_useful_at = ?2 WHERE memory_id = ?1",
params![memory_id, ts],
)?;
}
if is_cited {
if let Some(target) = conf_target {
let old_conf: f64 = conn
.query_row(
"SELECT confidence FROM memories WHERE memory_id = ?1",
params![memory_id],
|row| row.get::<_, f64>(0),
)
.unwrap_or(1.0);
let new_conf = (old_conf + CONF_ALPHA * (target - old_conf)).clamp(0.1, 0.99);
conn.execute(
"UPDATE memories SET confidence = ?2 WHERE memory_id = ?1",
params![memory_id, new_conf],
)?;
}
}
}
Ok(())
}
fn collect_cited_memory_ids(
conn: &Connection,
run_id: &str,
) -> KimetsuResult<std::collections::BTreeSet<String>> {
let mut stmt = conn.prepare(
"
SELECT DISTINCT memory_id
FROM memory_citations
WHERE run_id = ?1
",
)?;
let rows = stmt.query_map(params![run_id], |row| row.get::<_, String>(0))?;
let mut out = std::collections::BTreeSet::new();
for row in rows {
out.insert(row?);
}
Ok(out)
}
fn collect_injected_memory_ids(conn: &Connection, run_id: &str) -> KimetsuResult<Vec<String>> {
let mut stmt = conn.prepare(
"
SELECT payload_json
FROM events
WHERE run_id = ?1 AND kind = 'context.injected'
",
)?;
let rows = stmt.query_map(params![run_id], |row| row.get::<_, String>(0))?;
let mut seen = std::collections::BTreeSet::new();
for row in rows {
let payload_json = row?;
let payload: serde_json::Value = serde_json::from_str(&payload_json)?;
if let Some(ids) = payload.get("memory_ids").and_then(|v| v.as_array()) {
for id in ids {
if let Some(id_str) = id.as_str()
&& !id_str.is_empty()
{
seen.insert(id_str.to_string());
}
}
}
}
Ok(seen.into_iter().collect())
}
fn apply_memory_accepted(conn: &Connection, event: &Event) -> KimetsuResult<()> {
let Some(memory_id) = event
.payload
.get("memory_id")
.and_then(|value| value.as_str())
else {
return Ok(());
};
let scope = event
.payload
.get("scope")
.and_then(|value| value.as_str())
.unwrap_or("global_user");
let kind = event
.payload
.get("kind")
.and_then(|value| value.as_str())
.unwrap_or("fact");
let text = event
.payload
.get("text")
.and_then(|value| value.as_str())
.unwrap_or("");
let normalized_text = event
.payload
.get("normalized_text")
.and_then(|value| value.as_str())
.unwrap_or(text);
let confidence = event
.payload
.get("confidence")
.and_then(|value| value.as_f64())
.unwrap_or(1.0);
let initial_usefulness = event
.payload
.get("initial_usefulness")
.and_then(|value| value.as_f64())
.unwrap_or(0.0) as f32;
let provenance_snapshot = event
.payload
.get("provenance_snapshot")
.cloned()
.unwrap_or_else(
|| serde_json::json!({ "source": "event", "event_id": event.event_id.to_string() }),
);
conn.execute(
"
INSERT OR REPLACE INTO memories (
memory_id, scope, kind, text, normalized_text, confidence,
source_event_id, provenance_snapshot_json, created_at, use_count,
usefulness_score
)
VALUES (?1, ?2, ?3, ?4, ?5, ?6, ?7, ?8, ?9, 0, ?10)
",
params![
memory_id,
scope,
kind,
text,
normalized_text,
confidence,
event.event_id.to_string(),
serde_json::to_string(&provenance_snapshot)?,
ts_text(event)?,
initial_usefulness
],
)?;
conn.execute(
"DELETE FROM memories_fts WHERE memory_id = ?1",
params![memory_id],
)?;
conn.execute(
"INSERT INTO memories_fts (memory_id, text, kind, scope) VALUES (?1, ?2, ?3, ?4)",
params![memory_id, text, kind, scope],
)?;
Ok(())
}
fn apply_memory_proposed(conn: &Connection, event: &Event) -> KimetsuResult<()> {
let Some(proposal_id) = event
.payload
.get("proposal_id")
.and_then(|value| value.as_str())
else {
return Ok(());
};
let scope = event
.payload
.get("scope")
.and_then(|value| value.as_str())
.unwrap_or("run");
let kind = event
.payload
.get("kind")
.and_then(|value| value.as_str())
.unwrap_or("fact");
let text = event
.payload
.get("text")
.and_then(|value| value.as_str())
.unwrap_or("");
let rationale = event
.payload
.get("rationale")
.and_then(|value| value.as_str())
.unwrap_or("");
let confidence = event
.payload
.get("proposed_confidence")
.and_then(|value| value.as_f64())
.unwrap_or(0.5);
let source_event_ids = event
.payload
.get("source_event_ids")
.cloned()
.unwrap_or_else(|| serde_json::json!([]));
conn.execute(
"
INSERT OR REPLACE INTO memory_proposals (
proposal_id, run_id, scope, kind, text, rationale,
proposed_confidence, source_event_ids_json, status
)
VALUES (?1, ?2, ?3, ?4, ?5, ?6, ?7, ?8, 'pending')
",
params![
proposal_id,
event.run_id.to_string(),
scope,
kind,
text,
rationale,
confidence,
serde_json::to_string(&source_event_ids)?
],
)?;
Ok(())
}
fn apply_memory_rejected(conn: &Connection, event: &Event) -> KimetsuResult<()> {
let Some(proposal_id) = event
.payload
.get("proposal_id")
.and_then(|value| value.as_str())
else {
return Ok(());
};
let reason = event
.payload
.get("reason")
.and_then(|value| value.as_str())
.map(|s| s.to_string());
conn.execute(
"
UPDATE memory_proposals
SET status = 'rejected',
decided_at = ?2,
decided_by = 'cli',
decided_reason = ?3
WHERE proposal_id = ?1
",
params![proposal_id, ts_text(event)?, reason],
)?;
Ok(())
}
fn apply_memory_invalidated(conn: &Connection, event: &Event) -> KimetsuResult<()> {
let Some(memory_id) = event
.payload
.get("memory_id")
.and_then(|value| value.as_str())
else {
return Ok(());
};
let reason = event
.payload
.get("reason")
.and_then(|value| value.as_str())
.map(|s| s.to_string());
conn.execute(
"
UPDATE memories
SET invalidated_at = ?2,
invalidated_reason = ?3
WHERE memory_id = ?1
",
params![memory_id, ts_text(event)?, reason],
)?;
#[cfg(feature = "embeddings")]
crate::ann::on_invalidate(conn, memory_id);
Ok(())
}
fn apply_memory_superseded(conn: &Connection, event: &Event) -> KimetsuResult<()> {
let Some(memory_id) = event.payload.get("memory_id").and_then(|v| v.as_str()) else {
return Ok(());
};
let Some(survivor_id) = event.payload.get("survivor_id").and_then(|v| v.as_str()) else {
return Ok(());
};
let use_count_delta = event
.payload
.get("use_count_delta")
.and_then(|v| v.as_i64())
.unwrap_or(0);
let score_delta = event
.payload
.get("score_delta")
.and_then(|v| v.as_f64())
.unwrap_or(0.0);
let prior_survivor: Option<String> = conn
.query_row(
"SELECT superseded_by FROM memories WHERE memory_id = ?1",
params![memory_id],
|r| r.get::<_, Option<String>>(0),
)
.optional()?
.flatten();
if let Some(prev) = prior_survivor {
if prev != survivor_id {
let (a, b) = if prev.as_str() < survivor_id {
(prev.as_str(), survivor_id)
} else {
(survivor_id, prev.as_str())
};
let detected_at = ts_text(event)?;
conn.execute(
"INSERT OR IGNORE INTO sync_conflicts
(member_id, survivor_a, survivor_b, detected_at)
VALUES (?1, ?2, ?3, ?4)",
params![memory_id, a, b, detected_at],
)?;
}
}
conn.execute(
"UPDATE memories SET superseded_by = ?2 WHERE memory_id = ?1",
params![memory_id, survivor_id],
)?;
if use_count_delta != 0 || score_delta != 0.0 {
conn.execute(
"UPDATE memories
SET use_count = use_count + ?2,
usefulness_score = usefulness_score + ?3
WHERE memory_id = ?1",
params![survivor_id, use_count_delta, score_delta],
)?;
}
reassign_citations_projection(conn, memory_id, survivor_id)?;
conn.execute(
"DELETE FROM memories_fts WHERE memory_id = ?1",
params![memory_id],
)?;
#[cfg(feature = "embeddings")]
crate::ann::on_supersede(conn, memory_id);
let edge_ts = ts_text(event)?;
insert_memory_edge(conn, survivor_id, memory_id, "supersedes", &edge_ts)?;
Ok(())
}
fn apply_memory_temporal(conn: &Connection, event: &Event) -> KimetsuResult<()> {
let Some(memory_id) = event.payload.get("memory_id").and_then(|v| v.as_str()) else {
return Ok(());
};
let valid_from = event
.payload
.get("valid_from")
.and_then(|v| v.as_str())
.map(|s| s.to_string());
let valid_to = event
.payload
.get("valid_to")
.and_then(|v| v.as_str())
.map(|s| s.to_string());
match (valid_from, valid_to) {
(Some(vf), Some(vt)) => {
conn.execute(
"UPDATE memories SET valid_from = ?2, valid_to = ?3 WHERE memory_id = ?1",
params![memory_id, vf, vt],
)?;
}
(Some(vf), None) => {
conn.execute(
"UPDATE memories SET valid_from = ?2 WHERE memory_id = ?1",
params![memory_id, vf],
)?;
}
(None, Some(vt)) => {
conn.execute(
"UPDATE memories SET valid_to = ?2 WHERE memory_id = ?1",
params![memory_id, vt],
)?;
}
(None, None) => {} }
Ok(())
}
pub fn mark_memory_temporal(
conn: &Connection,
memory_id: &str,
valid_from: Option<&str>,
valid_to: Option<&str>,
) -> KimetsuResult<()> {
use kimetsu_core::ids::RunId;
let run_id = RunId::new();
let mut payload = serde_json::json!({ "memory_id": memory_id });
if let Some(vf) = valid_from {
payload["valid_from"] = serde_json::Value::String(vf.to_string());
}
if let Some(vt) = valid_to {
payload["valid_to"] = serde_json::Value::String(vt.to_string());
}
let event = kimetsu_core::event::Event::new(run_id, "memory.temporal", payload);
apply_event(conn, &event)
}
fn apply_memory_edge(conn: &Connection, event: &Event) -> KimetsuResult<()> {
let Some(src_id) = event.payload.get("src_id").and_then(|v| v.as_str()) else {
return Ok(());
};
let Some(dst_id) = event.payload.get("dst_id").and_then(|v| v.as_str()) else {
return Ok(());
};
let Some(edge_type) = event.payload.get("edge_type").and_then(|v| v.as_str()) else {
return Ok(());
};
if src_id == dst_id {
return Ok(());
}
let edge_ts = ts_text(event)?;
insert_memory_edge(conn, src_id, dst_id, edge_type, &edge_ts)
}
pub fn add_memory_edges(
conn: &Connection,
edges: &[(String, String, String)],
) -> KimetsuResult<usize> {
use kimetsu_core::ids::RunId;
let run_id = RunId::new();
let mut events = Vec::with_capacity(edges.len());
let mut written = 0usize;
for (src_id, dst_id, edge_type) in edges {
if src_id == dst_id {
continue;
}
let payload = serde_json::json!({
"src_id": src_id,
"dst_id": dst_id,
"edge_type": edge_type,
});
events.push(kimetsu_core::event::Event::new(
run_id,
"memory.edge",
payload,
));
written += 1;
}
apply_events(conn, &events)?;
Ok(written)
}
pub(crate) fn insert_memory_edge(
conn: &Connection,
src_id: &str,
dst_id: &str,
edge_type: &str,
created_at: &str,
) -> KimetsuResult<()> {
conn.execute(
"INSERT OR IGNORE INTO memory_edges (src_id, dst_id, edge_type, created_at)
VALUES (?1, ?2, ?3, ?4)",
params![src_id, dst_id, edge_type, created_at],
)?;
Ok(())
}
pub(crate) fn reassign_citations_projection(
conn: &Connection,
from_id: &str,
to_id: &str,
) -> KimetsuResult<()> {
let rows: Vec<(String, i64, String, Option<String>)> = {
let mut stmt = conn.prepare(
"SELECT run_id, turn, cited_at, rationale
FROM memory_citations WHERE memory_id = ?1",
)?;
stmt.query_map(params![from_id], |r| {
Ok((
r.get::<_, String>(0)?,
r.get::<_, i64>(1)?,
r.get::<_, String>(2)?,
r.get::<_, Option<String>>(3)?,
))
})?
.collect::<Result<_, _>>()?
};
for (run_id, turn, cited_at, rationale) in &rows {
conn.execute(
"INSERT OR IGNORE INTO memory_citations
(run_id, memory_id, turn, cited_at, rationale)
VALUES (?1, ?2, ?3, ?4, ?5)",
params![run_id, to_id, turn, cited_at, rationale],
)?;
}
conn.execute(
"DELETE FROM memory_citations WHERE memory_id = ?1",
params![from_id],
)?;
Ok(())
}
pub fn ensure_schema(conn: &Connection) -> KimetsuResult<()> {
schema::initialize(conn)
}
fn ts_text(event: &Event) -> KimetsuResult<String> {
Ok(event.ts.format(&Rfc3339)?)
}
#[cfg(test)]
mod tests {
use std::borrow::Cow;
use kimetsu_core::event::Event;
use kimetsu_core::ids::RunId;
use rusqlite::{Connection, params};
use serde_json::json;
use super::{apply_events, rebuild_in_place, upcast_event};
use crate::schema;
fn make_conn() -> Connection {
let conn = Connection::open_in_memory().expect("open_in_memory");
schema::initialize(&conn).expect("schema::initialize");
conn
}
fn make_event(run_id: RunId, kind: &str, payload: serde_json::Value) -> Event {
Event::new(run_id, kind, payload)
}
fn sentinel_run() -> RunId {
RunId(ulid::Ulid::nil())
}
#[test]
fn concurrent_cites_lose_no_updates() {
use std::sync::atomic::{AtomicU64, Ordering};
use std::sync::{Arc, Barrier};
static CTR: AtomicU64 = AtomicU64::new(0);
let n = CTR.fetch_add(1, Ordering::Relaxed);
let db_path =
std::env::temp_dir().join(format!("kimetsu-concurrency-{}-{n}.db", std::process::id()));
let _ = std::fs::remove_file(&db_path);
let mem_id = "mem-concurrency";
{
let conn = Connection::open(&db_path).expect("open seed");
schema::initialize(&conn).expect("init seed");
let accepted = Event::new(
sentinel_run(),
"memory.accepted",
json!({
"memory_id": mem_id,
"text": "hammer me",
"scope": "global_user",
"kind": "fact"
}),
);
apply_events(&conn, std::slice::from_ref(&accepted)).expect("seed accepted");
}
const THREADS: usize = 6;
const CITES_PER_THREAD: usize = 25;
let barrier = Arc::new(Barrier::new(THREADS));
let path = Arc::new(db_path.clone());
let mut handles = Vec::new();
for _ in 0..THREADS {
let b = Arc::clone(&barrier);
let p = Arc::clone(&path);
handles.push(std::thread::spawn(move || {
let conn = Connection::open(&*p).expect("open writer");
schema::initialize(&conn).expect("init writer");
b.wait(); for _ in 0..CITES_PER_THREAD {
let cited = Event::new(
sentinel_run(),
"memory.cited",
json!({ "memory_id": mem_id, "turn": 0 }),
);
apply_events(&conn, std::slice::from_ref(&cited))
.expect("concurrent cite must succeed");
}
}));
}
for h in handles {
h.join().expect("thread join");
}
let expected = (THREADS * CITES_PER_THREAD) as i64;
let conn = Connection::open(&db_path).expect("open verify");
schema::initialize(&conn).expect("init verify");
let use_count: i64 = conn
.query_row(
"SELECT use_count FROM memories WHERE memory_id = ?1",
params![mem_id],
|r| r.get(0),
)
.expect("read use_count");
assert_eq!(
use_count, expected,
"lost updates under concurrency: got {use_count}, expected {expected}"
);
let event_count: i64 = conn
.query_row("SELECT COUNT(*) FROM events", [], |r| r.get(0))
.expect("count events");
assert_eq!(event_count, expected + 1, "missing durable events");
rebuild_in_place(&conn).expect("rebuild");
let after: i64 = conn
.query_row(
"SELECT use_count FROM memories WHERE memory_id = ?1",
params![mem_id],
|r| r.get(0),
)
.expect("read use_count after rebuild");
assert_eq!(after, expected, "rebuild changed the projected use_count");
drop(conn);
let _ = std::fs::remove_file(&db_path);
let _ = std::fs::remove_file(db_path.with_extension("db-wal"));
let _ = std::fs::remove_file(db_path.with_extension("db-shm"));
}
#[test]
fn event_carries_and_roundtrips_origin() {
use super::{insert_event, read_events_ordered};
let conn = make_conn();
kimetsu_core::event::set_process_origin("test-machine/unit");
let ev = Event::new(
sentinel_run(),
"memory.accepted",
json!({
"memory_id": "m-origin",
"text": "with origin",
"scope": "global_user",
"kind": "fact"
}),
);
let stamped = ev.origin.clone();
insert_event(&conn, &ev).expect("insert");
let read_back = read_events_ordered(&conn).expect("read");
assert_eq!(read_back.len(), 1);
assert_eq!(read_back[0].origin, stamped, "origin must round-trip");
}
#[test]
fn upcast_is_identity_at_v1() {
let run_id = RunId::new();
let event = make_event(
run_id,
"run.started",
json!({"project_id": "p1", "task": "t"}),
);
assert_eq!(
event.schema_version, 1,
"Event::new must stamp schema_version=1"
);
let cow = upcast_event(&event);
assert!(
matches!(cow, Cow::Borrowed(_)),
"upcast_event must return Cow::Borrowed for current schema_version"
);
let out = cow.as_ref();
assert_eq!(out.kind, event.kind);
assert_eq!(out.schema_version, event.schema_version);
assert_eq!(out.payload, event.payload);
}
fn assert_empty_payload_ok(kind: &str) {
let conn = make_conn();
let run_id = RunId::new();
let event = make_event(run_id, kind, json!({}));
let result = apply_events(&conn, &[event]);
assert!(
result.is_ok(),
"apply_events with empty payload for kind={kind:?} must return Ok(()), got: {result:?}"
);
}
#[test]
fn empty_payload_run_started() {
assert_empty_payload_ok("run.started");
}
#[test]
fn empty_payload_run_finished() {
assert_empty_payload_ok("run.finished");
}
#[test]
fn empty_payload_run_failed() {
assert_empty_payload_ok("run.failed");
}
#[test]
fn empty_payload_run_aborted() {
assert_empty_payload_ok("run.aborted");
}
#[test]
fn empty_payload_memory_accepted() {
assert_empty_payload_ok("memory.accepted");
}
#[test]
fn empty_payload_memory_proposed() {
assert_empty_payload_ok("memory.proposed");
}
#[test]
fn empty_payload_memory_rejected() {
assert_empty_payload_ok("memory.rejected");
}
#[test]
fn empty_payload_memory_invalidated() {
assert_empty_payload_ok("memory.invalidated");
}
#[test]
fn empty_payload_memory_cited() {
assert_empty_payload_ok("memory.cited");
}
#[test]
fn empty_payload_work_episode() {
assert_empty_payload_ok("work.episode");
}
#[test]
fn empty_payload_memory_temporal() {
assert_empty_payload_ok("memory.temporal");
}
#[test]
fn well_formed_run_started_projects_correctly() {
let conn = make_conn();
let run_id = RunId::new();
let event = make_event(
run_id,
"run.started",
json!({
"project_id": "proj-abc",
"task": "fix the bug",
"model": "claude-sonnet-4-6"
}),
);
apply_events(&conn, &[event])
.expect("apply_events must succeed for well-formed run.started");
let row: (String, String, String) = conn
.query_row(
"SELECT run_id, project_id, task FROM runs WHERE run_id = ?1",
[run_id.to_string()],
|r| Ok((r.get(0)?, r.get(1)?, r.get(2)?)),
)
.expect("runs row must exist after apply_events");
assert_eq!(row.0, run_id.to_string());
assert_eq!(row.1, "proj-abc");
assert_eq!(row.2, "fix the bug");
}
#[test]
fn reset_projection_keeps_events() {
use super::reset_projection;
let conn = make_conn();
let run_id = RunId::new();
let mem_id = "mem-reset-test";
let events = vec![
make_event(
run_id,
"run.started",
json!({"project_id": "p", "task": "t"}),
),
make_event(
run_id,
"memory.accepted",
json!({
"memory_id": mem_id,
"text": "hello",
"scope": "global_user",
"kind": "fact"
}),
),
];
apply_events(&conn, &events).expect("apply_events");
let event_count: i64 = conn
.query_row("SELECT COUNT(*) FROM events", [], |r| r.get(0))
.unwrap();
assert!(event_count > 0, "events must be stored before reset");
let mem_count: i64 = conn
.query_row("SELECT COUNT(*) FROM memories", [], |r| r.get(0))
.unwrap();
assert_eq!(mem_count, 1, "memory must be projected before reset");
reset_projection(&conn).expect("reset_projection");
let event_count_after: i64 = conn
.query_row("SELECT COUNT(*) FROM events", [], |r| r.get(0))
.unwrap();
assert_eq!(
event_count_after, event_count,
"reset_projection must NOT delete from events"
);
let memories_after: i64 = conn
.query_row("SELECT COUNT(*) FROM memories", [], |r| r.get(0))
.unwrap();
assert_eq!(
memories_after, 0,
"memories must be cleared by reset_projection"
);
let runs_after: i64 = conn
.query_row("SELECT COUNT(*) FROM runs", [], |r| r.get(0))
.unwrap();
assert_eq!(runs_after, 0, "runs must be cleared by reset_projection");
let citations_after: i64 = conn
.query_row("SELECT COUNT(*) FROM memory_citations", [], |r| r.get(0))
.unwrap();
assert_eq!(
citations_after, 0,
"memory_citations must be cleared by reset_projection"
);
let conflicts_after: i64 = conn
.query_row("SELECT COUNT(*) FROM memory_conflicts", [], |r| r.get(0))
.unwrap();
assert_eq!(
conflicts_after, 0,
"memory_conflicts must be cleared by reset_projection"
);
let episodes_after: i64 = conn
.query_row("SELECT COUNT(*) FROM work_episodes", [], |r| r.get(0))
.unwrap();
assert_eq!(
episodes_after, 0,
"work_episodes must be cleared by reset_projection"
);
}
#[test]
fn rebuild_in_place_no_dup_events() {
use super::rebuild_in_place;
let conn = make_conn();
let run_id = RunId::new();
let mem_id = "mem-dup-test";
let events = vec![
make_event(
run_id,
"run.started",
json!({"project_id": "p", "task": "t"}),
),
make_event(
run_id,
"memory.accepted",
json!({
"memory_id": mem_id,
"text": "no dup",
"scope": "global_user",
"kind": "fact"
}),
),
make_event(run_id, "run.finished", json!({"total_cost_usd": 0.01})),
];
apply_events(&conn, &events).expect("apply_events");
let event_count_before: i64 = conn
.query_row("SELECT COUNT(*) FROM events", [], |r| r.get(0))
.unwrap();
assert_eq!(event_count_before, 3, "expected 3 events seeded");
conn.execute_batch("DELETE FROM memories; DELETE FROM memories_fts;")
.unwrap();
let mem_count_wiped: i64 = conn
.query_row("SELECT COUNT(*) FROM memories", [], |r| r.get(0))
.unwrap();
assert_eq!(mem_count_wiped, 0, "memories wiped before rebuild_in_place");
let replayed = rebuild_in_place(&conn).expect("rebuild_in_place");
assert_eq!(
replayed, 3,
"rebuild_in_place must return the number of events replayed"
);
let mem_exists: i64 = conn
.query_row(
"SELECT COUNT(*) FROM memories WHERE memory_id = ?1",
[mem_id],
|r| r.get(0),
)
.unwrap();
assert_eq!(
mem_exists, 1,
"memory must be re-projected after rebuild_in_place"
);
let event_count_after: i64 = conn
.query_row("SELECT COUNT(*) FROM events", [], |r| r.get(0))
.unwrap();
assert_eq!(
event_count_after, event_count_before,
"rebuild_in_place must NOT insert duplicate events"
);
}
#[test]
fn rebuild_in_place_reconstructs_citations() {
use super::rebuild_in_place;
let conn = make_conn();
let run_id = RunId::new();
let mem_id = "mem-cite-test";
let events = vec![
make_event(
run_id,
"run.started",
json!({"project_id": "p", "task": "t"}),
),
make_event(
run_id,
"memory.accepted",
json!({
"memory_id": mem_id,
"text": "cite me",
"scope": "global_user",
"kind": "fact"
}),
),
make_event(
run_id,
"memory.cited",
json!({
"memory_id": mem_id,
"turn": 2,
"rationale": "relevant context"
}),
),
make_event(run_id, "run.finished", json!({"total_cost_usd": 0.0})),
];
apply_events(&conn, &events).expect("apply_events");
let citations_before: i64 = conn
.query_row("SELECT COUNT(*) FROM memory_citations", [], |r| r.get(0))
.unwrap();
assert_eq!(
citations_before, 1,
"citation must exist after apply_events"
);
let replayed = rebuild_in_place(&conn).expect("rebuild_in_place");
assert_eq!(replayed, 4, "expected 4 events replayed");
let citations_after: i64 = conn
.query_row("SELECT COUNT(*) FROM memory_citations", [], |r| r.get(0))
.unwrap();
assert_eq!(
citations_after, 1,
"memory_citations must be repopulated by rebuild_in_place"
);
}
#[test]
fn add_memory_edges_writes_and_survives_rebuild() {
use super::{add_memory_edges, rebuild_in_place};
let conn = make_conn();
let run_id = RunId::new();
let m1 = "mem-edge-a";
let m2 = "mem-edge-b";
let events = vec![
make_event(
run_id,
"memory.accepted",
json!({"memory_id": m1, "text": "alpha", "scope": "global_user", "kind": "fact"}),
),
make_event(
run_id,
"memory.accepted",
json!({"memory_id": m2, "text": "beta", "scope": "global_user", "kind": "fact"}),
),
];
apply_events(&conn, &events).expect("apply_events");
let written = add_memory_edges(
&conn,
&[
(m1.to_string(), m1.to_string(), "relates_to".to_string()),
(m1.to_string(), m2.to_string(), "relates_to".to_string()),
],
)
.expect("add_memory_edges");
assert_eq!(
written, 1,
"self-loop must be skipped, one real edge written"
);
let edge_count = |c: &Connection| -> i64 {
c.query_row(
"SELECT COUNT(*) FROM memory_edges WHERE src_id=?1 AND dst_id=?2 AND edge_type='relates_to'",
params![m1, m2],
|r| r.get(0),
)
.unwrap()
};
assert_eq!(edge_count(&conn), 1, "edge present after write");
let total_edges_before: i64 = conn
.query_row("SELECT COUNT(*) FROM memory_edges", [], |r| r.get(0))
.unwrap();
rebuild_in_place(&conn).expect("rebuild_in_place");
let total_edges_after: i64 = conn
.query_row("SELECT COUNT(*) FROM memory_edges", [], |r| r.get(0))
.unwrap();
assert_eq!(
total_edges_before, total_edges_after,
"rebuild must reproduce exactly the same edge set"
);
assert_eq!(edge_count(&conn), 1, "edge survives rebuild_in_place");
}
#[test]
fn memory_proposed_redacts_event_and_projection_payloads() {
let conn = make_conn();
let run_id = RunId::new();
let secret = "sk-ant-api03-AbCdEfGhIjKlMnOpQrStUv0123456789AbCdEf";
let event = make_event(
run_id,
"memory.proposed",
json!({
"proposal_id": "prop-redact",
"scope": "project",
"kind": "fact",
"text": format!("lesson uses {secret}"),
"rationale": format!("model repeated {secret}"),
"proposed_confidence": 0.5,
"source_event_ids": [],
}),
);
apply_events(&conn, &[event]).expect("apply_events");
let payload: String = conn
.query_row(
"SELECT payload_json FROM events WHERE kind = 'memory.proposed'",
[],
|r| r.get(0),
)
.expect("event payload");
assert!(!payload.contains(secret), "event leaked secret: {payload}");
assert!(payload.contains("[REDACTED:anthropic_oauth]"));
let row: (String, String) = conn
.query_row(
"SELECT text, rationale FROM memory_proposals WHERE proposal_id = 'prop-redact'",
[],
|r| Ok((r.get(0)?, r.get(1)?)),
)
.expect("proposal row");
assert!(!row.0.contains(secret), "proposal text leaked: {}", row.0);
assert!(
!row.1.contains(secret),
"proposal rationale leaked: {}",
row.1
);
}
#[test]
fn memory_cited_redacts_event_and_projection_rationale() {
let conn = make_conn();
let run_id = RunId::new();
let secret = "sk-ant-api03-AbCdEfGhIjKlMnOpQrStUv0123456789AbCdEf";
let event = make_event(
run_id,
"memory.cited",
json!({
"memory_id": "mem-redact",
"turn": 1,
"rationale": format!("used because output showed {secret}"),
}),
);
apply_events(&conn, &[event]).expect("apply_events");
let payload: String = conn
.query_row(
"SELECT payload_json FROM events WHERE kind = 'memory.cited'",
[],
|r| r.get(0),
)
.expect("event payload");
assert!(!payload.contains(secret), "event leaked secret: {payload}");
assert!(payload.contains("[REDACTED:anthropic_oauth]"));
let rationale: String = conn
.query_row(
"SELECT rationale FROM memory_citations WHERE memory_id = 'mem-redact'",
[],
|r| r.get(0),
)
.expect("citation rationale");
assert!(
!rationale.contains(secret),
"citation rationale leaked: {rationale}"
);
assert!(rationale.contains("[REDACTED:anthropic_oauth]"));
}
#[test]
fn memory_temporal_stamps_validity_and_survives_rebuild() {
use super::{mark_memory_temporal, rebuild_in_place};
let conn = make_conn();
let run_id = RunId::new();
let mem_id = "mem-temporal-test";
let events = vec![make_event(
run_id,
"memory.accepted",
json!({
"memory_id": mem_id,
"text": "old fact that expired",
"scope": "project",
"kind": "fact",
"confidence": 0.9
}),
)];
apply_events(&conn, &events).expect("apply_events");
mark_memory_temporal(
&conn,
mem_id,
Some("2020-01-01T00:00:00Z"),
Some("2025-01-01T00:00:00Z"),
)
.expect("mark_memory_temporal");
let (vf, vt): (Option<String>, Option<String>) = conn
.query_row(
"SELECT valid_from, valid_to FROM memories WHERE memory_id = ?1",
[mem_id],
|r| Ok((r.get(0)?, r.get(1)?)),
)
.expect("query valid_from/valid_to");
assert_eq!(
vf.as_deref(),
Some("2020-01-01T00:00:00Z"),
"valid_from must be set"
);
assert_eq!(
vt.as_deref(),
Some("2025-01-01T00:00:00Z"),
"valid_to must be set"
);
rebuild_in_place(&conn).expect("rebuild_in_place");
let (vf2, vt2): (Option<String>, Option<String>) = conn
.query_row(
"SELECT valid_from, valid_to FROM memories WHERE memory_id = ?1",
[mem_id],
|r| Ok((r.get(0)?, r.get(1)?)),
)
.expect("query valid_from/valid_to after rebuild");
assert_eq!(
vf2.as_deref(),
Some("2020-01-01T00:00:00Z"),
"valid_from must survive rebuild_in_place"
);
assert_eq!(
vt2.as_deref(),
Some("2025-01-01T00:00:00Z"),
"valid_to must survive rebuild_in_place"
);
}
#[test]
fn rebuild_in_place_payload_fidelity() {
use super::rebuild_in_place;
let conn = make_conn();
let run_id = RunId::new();
let mem_id = "mem-fidelity-test";
let expected_text = "Rust edition 2024 requires explicit use of `use` for trait impls";
let expected_scope = "project";
let expected_kind = "guideline";
let events = vec![
make_event(
run_id,
"run.started",
json!({"project_id": "p", "task": "t"}),
),
make_event(
run_id,
"memory.accepted",
json!({
"memory_id": mem_id,
"text": expected_text,
"scope": expected_scope,
"kind": expected_kind,
"confidence": 0.9
}),
),
];
apply_events(&conn, &events).expect("apply_events");
conn.execute_batch("DELETE FROM memories; DELETE FROM memories_fts; DELETE FROM runs;")
.unwrap();
rebuild_in_place(&conn).expect("rebuild_in_place");
let row: (String, String, String) = conn
.query_row(
"SELECT text, scope, kind FROM memories WHERE memory_id = ?1",
[mem_id],
|r| Ok((r.get(0)?, r.get(1)?, r.get(2)?)),
)
.expect("memory must exist after rebuild_in_place");
assert_eq!(row.0, expected_text, "text must round-trip through rebuild");
assert_eq!(
row.1, expected_scope,
"scope must round-trip through rebuild"
);
assert_eq!(row.2, expected_kind, "kind must round-trip through rebuild");
}
#[test]
fn initial_usefulness_seeds_score_and_survives_rebuild() {
use super::rebuild_in_place;
let conn = make_conn();
let run_id = RunId::new();
let events = vec![
make_event(
run_id,
"memory.accepted",
json!({
"memory_id": "salient",
"text": "rm -rf node_modules then reinstall fixes the EBUSY lock",
"scope": "project",
"kind": "failure_pattern",
"confidence": 1.0,
"initial_usefulness": 0.3
}),
),
make_event(
run_id,
"memory.accepted",
json!({
"memory_id": "neutral",
"text": "the readme mentions a port number",
"scope": "project",
"kind": "fact",
"confidence": 1.0
}),
),
];
apply_events(&conn, &events).expect("apply_events");
let read = |id: &str| -> f64 {
conn.query_row(
"SELECT usefulness_score FROM memories WHERE memory_id = ?1",
[id],
|r| r.get(0),
)
.unwrap()
};
assert!(
(read("salient") - 0.3).abs() < 1e-6,
"salient memory must be seeded to 0.3"
);
assert!(
read("neutral").abs() < 1e-6,
"memory without initial_usefulness must default to 0.0"
);
assert!(
read("salient") > read("neutral"),
"salient new memory must outrank a neutral one from day one"
);
conn.execute_batch("DELETE FROM memories; DELETE FROM memories_fts;")
.unwrap();
rebuild_in_place(&conn).expect("rebuild_in_place");
assert!(
(read("salient") - 0.3).abs() < 1e-6,
"initial_usefulness seed must survive rebuild"
);
}
fn cite_and_terminate_confidence(terminal_kind: &str) -> (Connection, f64) {
let conn = make_conn();
let run_id = RunId::new();
let mem_id = "cal-mem";
let events = vec![
make_event(
run_id,
"run.started",
json!({"project_id": "p", "task": "t"}),
),
make_event(
run_id,
"memory.accepted",
json!({
"memory_id": mem_id,
"text": "use lld linker on windows",
"scope": "project",
"kind": "convention",
"confidence": 0.7
}),
),
make_event(
run_id,
"context.injected",
json!({"stage": "loc", "memory_ids": [mem_id], "used_tokens": 100}),
),
make_event(
run_id,
"memory.cited",
json!({"memory_id": mem_id, "turn": 1}),
),
make_event(run_id, terminal_kind, json!({"total_cost_usd": 0.0})),
];
apply_events(&conn, &events).expect("apply_events");
let conf: f64 = conn
.query_row(
"SELECT confidence FROM memories WHERE memory_id = ?1",
[mem_id],
|r| r.get(0),
)
.unwrap();
(conn, conf)
}
#[test]
fn confidence_calibration_rewards_success_and_survives_rebuild() {
use super::rebuild_in_place;
let (success_conn, success_conf) = cite_and_terminate_confidence("run.finished");
let (_fail_conn, fail_conf) = cite_and_terminate_confidence("run.failed");
assert!(
success_conf > 0.7,
"successful citation must raise confidence above 0.7, got {success_conf}"
);
assert!(
fail_conf < 0.7,
"failed citation must lower confidence below 0.7, got {fail_conf}"
);
assert!(
success_conf > fail_conf,
"cited-in-success must beat cited-in-failure: {success_conf} vs {fail_conf}"
);
let mem_id = "cal-mem";
success_conn
.execute_batch("DELETE FROM memories; DELETE FROM memories_fts;")
.unwrap();
rebuild_in_place(&success_conn).expect("rebuild_in_place");
let post: f64 = success_conn
.query_row(
"SELECT confidence FROM memories WHERE memory_id = ?1",
[mem_id],
|r| r.get(0),
)
.unwrap();
assert!(
(post - success_conf).abs() < 1e-9,
"calibrated confidence must reproduce after rebuild: {post} vs {success_conf}"
);
}
}