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
)
})
}
pub(crate) 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<()> {
with_write_txn(conn, |c| {
for event in events {
insert_event(c, redact_memory_event(event).as_ref())?;
}
replay_locked(c).map(|_| ())
})
}
pub fn rebuild_in_place(conn: &Connection) -> KimetsuResult<usize> {
let mut count = 0;
with_write_txn(conn, |c| {
count = replay_locked(c)?;
Ok(())
})?;
Ok(count)
}
pub(crate) fn replay_locked(conn: &Connection) -> KimetsuResult<usize> {
let events = read_events_ordered(conn)?;
let existing = {
let mut stmt = conn.prepare("SELECT memory_id FROM memories")?;
stmt.query_map([], |r| r.get::<_, String>(0))?
.collect::<Result<std::collections::BTreeSet<_>, _>>()?
};
reset_projection(conn)?;
for event in &events {
let bound = bind_injected_revisions(conn, event)?;
if matches!(&bound, Cow::Owned(_)) {
conn.execute(
"UPDATE events SET payload_json=?2 WHERE event_id=?1",
params![
event.event_id.to_string(),
serde_json::to_string(&bound.payload)?
],
)?;
}
project_event(conn, bound.as_ref())?;
}
let mut stmt = conn.prepare("SELECT memory_id FROM memories")?;
let restored = stmt
.query_map([], |r| r.get::<_, String>(0))?
.collect::<Result<std::collections::BTreeSet<_>, _>>()?;
let missing = existing.difference(&restored).count();
if missing > 0 {
return Err(format!("rebuild refused: {missing} existing memories absent from replay; transaction rolled back. Back up the brain and recover missing events before rebuilding; legacy unlogged rows require migration.").into());
}
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<()> {
apply_events_checked(conn, events, |_| Ok(()))
}
pub(crate) fn apply_events_checked<F>(
conn: &Connection,
events: &[Event],
mut validate: F,
) -> KimetsuResult<()>
where
F: FnMut(&Connection) -> KimetsuResult<()>,
{
with_write_txn(conn, |c| {
validate(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_revisions;
DELETE FROM memory_facts;
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 memory_entities;
DELETE FROM work_episodes;
",
)?;
Ok(())
}
pub(crate) fn apply_event(conn: &Connection, event: &Event) -> KimetsuResult<()> {
let event = redact_memory_event(event);
let event = bind_injected_revisions(conn, event.as_ref())?;
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.restored" => apply_memory_restored(conn, event),
"conflict.resolved" => crate::conflict::project_resolution(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.corrected" => apply_memory_corrected(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" | "memory.corrected"
) {
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 current_revision = claim_revision_at(conn, memory_id, None)?;
let explicit = event
.payload
.get("revision_event_id")
.and_then(|v| v.as_str());
let exposed =
if let Some(exposure_id) = event.payload.get("exposure_id").and_then(|v| v.as_str()) {
exact_claim_exposure(conn, exposure_id, &event.run_id.to_string(), memory_id)?
} else {
run_claim_revision(conn, &event.run_id.to_string(), memory_id)?
};
let exposed = match exposed {
ClaimExposure::Unbound => return Ok(()),
ClaimExposure::Absent => None,
ClaimExposure::Bound(revision) => Some(revision),
};
if explicit
.zip(exposed.as_deref())
.is_some_and(|(explicit, delivered)| explicit != delivered)
{
return Ok(());
}
let evidence_revision = explicit.map(str::to_owned).or(exposed);
if evidence_revision
.as_ref()
.is_some_and(|r| r != ¤t_revision)
|| (evidence_revision.is_none() && !current_revision.starts_with("baseline:"))
{
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 let Some(query) = event.payload.get("query").and_then(|v| v.as_str()) {
conn.execute(
"UPDATE memory_citations SET query = ?4
WHERE run_id = ?1 AND memory_id = ?2 AND turn = ?3",
params![event.run_id.to_string(), memory_id, turn, query],
)?;
}
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(())
}
use crate::scoring::{CITED_DELTA, CONF_ALPHA, FAILURE_PENALTY_CITES_DIVISOR, PASSENGER_DELTA};
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" => (CITED_DELTA, PASSENGER_DELTA),
"run.failed" => {
let category = event
.payload
.get("category")
.and_then(|value| value.as_str())
.unwrap_or("");
if category == "Gate" {
return Ok(());
}
(-CITED_DELTA, -PASSENGER_DELTA)
}
_ => 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 exposure = run_claim_revision(conn, &run_id, memory_id)?;
if matches!(exposure, ClaimExposure::Unbound) {
continue;
}
if let ClaimExposure::Bound(exposed_revision) = exposure {
if exposed_revision != claim_revision_at(conn, memory_id, None)? {
let was_cited: bool = conn.query_row(
"SELECT EXISTS(SELECT 1 FROM events WHERE run_id=?1 AND kind='memory.cited' AND json_extract(payload_json,'$.memory_id')=?2
AND (json_extract(payload_json,'$.revision_event_id') IS NULL OR json_extract(payload_json,'$.revision_event_id')=?3)
AND rowid <= (SELECT rowid FROM events WHERE event_id=?4))",
params![run_id,memory_id,exposed_revision,event.event_id.to_string()], |r| r.get(0))?;
let delta = if was_cited { strong } else { weak };
conn.execute("UPDATE memory_revisions SET use_count=use_count+1,usefulness_score=usefulness_score+?2,
confidence=CASE WHEN ?3 THEN MAX(0.1,MIN(0.99,confidence+?4*(?5-confidence))) ELSE confidence END
WHERE event_id=?1 AND memory_id=?6", params![exposed_revision,delta,was_cited,CONF_ALPHA,conf_target.unwrap_or(1.0),memory_id])?;
continue;
}
}
let is_cited = cited.contains(memory_id);
let delta = if is_cited {
if strong < 0.0 {
let prior_cites: i64 = conn
.query_row(
"SELECT COUNT(*) FROM memory_citations
WHERE memory_id = ?1 AND run_id != ?2",
params![memory_id, run_id],
|row| row.get(0),
)
.unwrap_or(0);
strong / (1.0 + prior_cites as f64 / FAILURE_PENALTY_CITES_DIVISOR)
} else {
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 event_validity<'a>(
conn: &Connection,
event: &'a Event,
) -> KimetsuResult<(Option<&'a str>, Option<&'a str>)> {
let endpoint = |name| -> KimetsuResult<Option<&'a str>> {
match event.payload.get(name) {
None | Some(serde_json::Value::Null) => Ok(None),
Some(serde_json::Value::String(value)) => Ok(Some(value.as_str())),
_ => Err(format!("{name} must be a timestamp string or null").into()),
}
};
let from = endpoint("valid_from")?;
let to = endpoint("valid_to")?;
let (start, end): (Option<f64>, Option<f64>) = conn.query_row(
"SELECT julianday(?1),julianday(?2)",
params![from, to],
|r| Ok((r.get(0)?, r.get(1)?)),
)?;
if (from.is_some() && start.is_none())
|| (to.is_some() && end.is_none())
|| matches!((start,end), (Some(a),Some(b)) if a >= b)
{
return Err("invalid or empty temporal validity interval".into());
}
Ok((from, to))
}
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 (valid_from, valid_to) = event_validity(conn, event)?;
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, valid_from, valid_to
)
VALUES (?1, ?2, ?3, ?4, ?5, ?6, ?7, ?8, ?9, 0, ?10, ?11, ?12)
",
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,
valid_from,
valid_to
],
)?;
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],
)?;
let _ = crate::graph::project_entities(conn, memory_id, text);
if let Some(proposal_id) = event.payload.get("proposal_id").and_then(|v| v.as_str()) {
conn.execute("UPDATE memory_proposals SET status='accepted', decided_at=?2, decided_by='cli' WHERE proposal_id=?1",
params![proposal_id,ts_text(event)?])?;
}
crate::fact_store::refresh(conn, memory_id)?;
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 (valid_from, valid_to) = event_validity(conn, event)?;
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, valid_from, valid_to
)
VALUES (?1, ?2, ?3, ?4, ?5, ?6, ?7, ?8, 'pending', ?9, ?10)
",
params![
proposal_id,
event.run_id.to_string(),
scope,
kind,
text,
rationale,
confidence,
serde_json::to_string(&source_event_ids)?,
valid_from,
valid_to
],
)?;
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());
if reason
.as_deref()
.is_some_and(|r| matches!(r, "forgotten" | "forgotten/archived" | "forgotten_archived"))
{
let active:bool=conn.query_row("SELECT EXISTS(SELECT 1 FROM memories WHERE memory_id=?1 AND invalidated_at IS NULL AND superseded_by IS NULL)",[memory_id],|r|r.get(0))?;
if !active {
return Ok(());
}
}
conn.execute(
"
UPDATE memories
SET invalidated_at = ?2,
invalidated_reason = ?3
WHERE memory_id = ?1
",
params![memory_id, ts_text(event)?, reason],
)?;
conn.execute("DELETE FROM memories_fts WHERE memory_id=?1", [memory_id])?;
#[cfg(feature = "embeddings")]
crate::ann::on_invalidate(conn, memory_id);
let _ = crate::graph::forget_entities(conn, memory_id);
Ok(())
}
fn apply_memory_restored(conn: &Connection, event: &Event) -> KimetsuResult<()> {
let Some(id) = event.payload.get("memory_id").and_then(|v| v.as_str()) else {
return Ok(());
};
let changed = conn.execute(
"UPDATE memories SET invalidated_at=NULL, invalidated_reason=NULL
WHERE memory_id=?1 AND superseded_by IS NULL AND invalidated_at IS NOT NULL
AND invalidated_reason IN ('forgotten','forgotten/archived','forgotten_archived')",
[id],
)?;
if changed > 0 {
conn.execute("DELETE FROM memories_fts WHERE memory_id=?1", [id])?;
conn.execute("INSERT INTO memories_fts(memory_id,text,kind,scope) SELECT memory_id,text,kind,scope FROM memories WHERE memory_id=?1", [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 _ = crate::graph::forget_entities(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) => {} }
crate::fact_store::refresh(conn, memory_id)?;
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
}
#[test]
fn rebuild_import_failure_preserves_existing_projection() {
let conn = make_conn();
let accepted = Event::new(
RunId::new(),
"memory.accepted",
json!({
"memory_id":"kept", "text":"keep my evidence", "scope":"project", "kind":"fact"
}),
);
apply_events(&conn, &[accepted]).unwrap();
let malformed = Event::new(
RunId::new(),
"memory.accepted",
json!({
"memory_id":"broken", "text":"bad validity", "scope":"project", "kind":"fact", "valid_from":"nonsense"
}),
);
assert!(super::rebuild(&conn, &[malformed]).is_err());
let text: String = conn
.query_row(
"SELECT text FROM memories WHERE memory_id='kept'",
[],
|r| r.get(0),
)
.unwrap();
assert_eq!(text, "keep my evidence");
}
#[test]
fn rebuild_import_keeps_durable_events_missing_from_trace() {
let conn = make_conn();
let accepted = Event::new(
RunId::new(),
"memory.accepted",
json!({
"memory_id":"kept", "text":"durable but not in trace", "scope":"project", "kind":"fact"
}),
);
apply_events(&conn, &[accepted]).unwrap();
super::rebuild(&conn, &[]).unwrap();
assert_eq!(
conn.query_row(
"SELECT count(*) FROM memories WHERE memory_id='kept'",
[],
|r| r.get::<_, i64>(0)
)
.unwrap(),
1
);
}
#[test]
fn trace_import_binds_historical_exposure_before_later_correction() {
let conn = make_conn();
let accepted = Event::new(
RunId::new(),
"memory.accepted",
json!({
"memory_id":"m", "text":"old claim", "scope":"project", "kind":"fact"
}),
);
apply_events(&conn, std::slice::from_ref(&accepted)).unwrap();
let original_revision = super::claim_revision_at(&conn, "m", None).unwrap();
let exposure = Event::new(
RunId::new(),
"context.injected",
json!({"memory_ids":["m"]}),
);
let corrected = Event::new(
RunId::new(),
"memory.corrected",
json!({"memory_id":"m", "text":"new claim"}),
);
apply_events(&conn, &[corrected]).unwrap();
super::rebuild(&conn, std::slice::from_ref(&exposure)).unwrap();
for _ in 0..2 {
let revision: String = conn.query_row("SELECT json_extract(payload_json,'$.memory_revisions.m') FROM events WHERE event_id=?1", [exposure.event_id.to_string()], |r|r.get(0)).unwrap();
assert_eq!(revision, original_revision);
rebuild_in_place(&conn).unwrap();
}
}
#[test]
fn trace_import_replays_missing_correction_before_later_invalidation() {
let conn = make_conn();
let accepted = Event::new(
RunId::new(),
"memory.accepted",
json!({
"memory_id":"m", "text":"old claim", "scope":"project", "kind":"fact"
}),
);
apply_events(&conn, &[accepted]).unwrap();
let correction = Event::new(
RunId::new(),
"memory.corrected",
json!({"memory_id":"m", "text":"historically corrected"}),
);
let invalidated = Event::new(
RunId::new(),
"memory.invalidated",
json!({"memory_id":"m", "reason":"retired"}),
);
apply_events(&conn, &[invalidated]).unwrap();
super::rebuild(&conn, &[correction]).unwrap();
for _ in 0..2 {
let row: (String, bool) = conn
.query_row(
"SELECT text,invalidated_at IS NOT NULL FROM memories WHERE memory_id='m'",
[],
|r| Ok((r.get(0)?, r.get(1)?)),
)
.unwrap();
assert_eq!(row, ("historically corrected".into(), true));
rebuild_in_place(&conn).unwrap();
}
}
#[test]
fn rebuild_refuses_to_erase_unlogged_legacy_user_memory() {
let conn = make_conn();
conn.execute("INSERT INTO memories(memory_id,scope,kind,text,normalized_text,confidence,provenance_snapshot_json,created_at) VALUES ('legacy','global_user','fact','original','original',0.7,'{\"source\":\"user_brain\"}','2020-01-01T00:00:00Z')", []).unwrap();
let error = rebuild_in_place(&conn).unwrap_err();
assert!(error.to_string().contains("absent from replay"));
assert_eq!(
conn.query_row(
"SELECT text FROM memories WHERE memory_id='legacy'",
[],
|r| r.get::<_, String>(0)
)
.unwrap(),
"original"
);
assert!(super::rebuild(&conn, &[]).is_err());
}
#[test]
fn rebuild_reads_events_after_acquiring_writer_lock() {
use std::sync::atomic::{AtomicBool, Ordering};
static WAITING: AtomicBool = AtomicBool::new(false);
fn busy(_: i32) -> bool {
WAITING.store(true, Ordering::SeqCst);
std::thread::sleep(std::time::Duration::from_millis(1));
true
}
WAITING.store(false, Ordering::SeqCst);
let dir = tempfile::tempdir().unwrap();
let path = dir.path().join("rebuild.db");
let writer = Connection::open(&path).unwrap();
schema::initialize(&writer).unwrap();
writer.execute_batch("PRAGMA journal_mode=WAL").unwrap();
let rebuilding = Connection::open(&path).unwrap();
rebuilding.busy_handler(Some(busy)).unwrap();
writer.execute_batch("BEGIN IMMEDIATE").unwrap();
let event = Event::new(
RunId::new(),
"memory.accepted",
json!({
"memory_id":"concurrent", "text":"committed while rebuild waits", "scope":"project", "kind":"fact"
}),
);
super::apply_event(&writer, &event).unwrap();
let worker =
std::thread::spawn(move || rebuild_in_place(&rebuilding).map_err(|e| e.to_string()));
let deadline = std::time::Instant::now() + std::time::Duration::from_secs(5);
while !WAITING.load(Ordering::SeqCst) && std::time::Instant::now() < deadline {
std::thread::sleep(std::time::Duration::from_millis(1));
}
let was_waiting = WAITING.load(Ordering::SeqCst);
writer.execute_batch("COMMIT").unwrap();
assert!(
was_waiting,
"rebuild did not reach the contested write lock"
);
assert_eq!(worker.join().unwrap().unwrap(), 1);
assert_eq!(
writer
.query_row(
"SELECT count(*) FROM memories WHERE memory_id='concurrent'",
[],
|r| r.get::<_, i64>(0)
)
.unwrap(),
1
);
}
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_manual_regrets_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(),
"retrieval.regret",
json!({ "memory_id": mem_id, "source": "manual" }),
);
apply_events(&conn, std::slice::from_ref(&cited))
.expect("concurrent explicit regret 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 failure_penalty_scales_with_prior_citations() {
let conn = make_conn();
for (id, text) in [
("proven", "use lld on windows"),
("rookie", "try the new flag"),
] {
apply_events(
&conn,
&[make_event(
RunId::new(),
"memory.accepted",
json!({"memory_id": id, "text": text, "scope": "project", "kind": "convention", "confidence": 0.7}),
)],
)
.unwrap();
}
for _ in 0..3 {
let run = RunId::new();
apply_events(
&conn,
&[
make_event(run, "run.started", json!({"project_id": "p", "task": "t"})),
make_event(
run,
"context.injected",
json!({"stage": "loc", "memory_ids": ["proven"], "used_tokens": 50}),
),
make_event(
run,
"memory.cited",
json!({"memory_id": "proven", "turn": 1}),
),
make_event(run, "run.finished", json!({"outcome": "ok"})),
],
)
.unwrap();
}
let score = |id: &str| -> f64 {
conn.query_row(
"SELECT usefulness_score FROM memories WHERE memory_id = ?1",
[id],
|r| r.get(0),
)
.unwrap()
};
let proven_before = score("proven");
let rookie_before = score("rookie");
let fail_run = RunId::new();
apply_events(
&conn,
&[
make_event(
fail_run,
"run.started",
json!({"project_id": "p", "task": "t"}),
),
make_event(
fail_run,
"context.injected",
json!({"stage": "loc", "memory_ids": ["proven", "rookie"], "used_tokens": 80}),
),
make_event(
fail_run,
"memory.cited",
json!({"memory_id": "proven", "turn": 1}),
),
make_event(
fail_run,
"memory.cited",
json!({"memory_id": "rookie", "turn": 1}),
),
make_event(
fail_run,
"run.failed",
json!({"category": "Verification", "message": "flaky test"}),
),
],
)
.unwrap();
let proven_drop = proven_before - score("proven");
let rookie_drop = rookie_before - score("rookie");
assert!(
(rookie_drop - 1.0).abs() < 1e-6,
"unproven memory takes the full -1.0, got -{rookie_drop}"
);
assert!(
proven_drop < rookie_drop,
"3 prior citations must shrink the penalty: proven -{proven_drop} vs rookie -{rookie_drop}"
);
assert!(
(proven_drop - 0.5).abs() < 1e-6,
"penalty with 3 priors must be -0.5, got -{proven_drop}"
);
}
#[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}"
);
}
}
fn apply_memory_corrected(conn: &Connection, event: &Event) -> KimetsuResult<()> {
let id = event
.payload
.get("memory_id")
.and_then(|v| v.as_str())
.ok_or("correction requires memory_id")?;
if conn.query_row(
"SELECT EXISTS(SELECT 1 FROM memory_revisions WHERE event_id=?1)",
params![event.event_id.to_string()],
|r| r.get::<_, bool>(0),
)? {
return Ok(());
}
let old: Option<(String, String, Option<String>)> = conn
.query_row(
"SELECT text, kind, invalidated_at FROM memories WHERE memory_id=?1",
params![id],
|r| Ok((r.get(0)?, r.get(1)?, r.get(2)?)),
)
.optional()?;
let (old_text, old_kind, invalidated) = old.ok_or_else(|| format!("memory not found: {id}"))?;
if invalidated.is_some() {
return Err(format!("memory {id} is already invalidated").into());
}
let text = event.payload.get("text").and_then(|v| v.as_str());
let kind = event.payload.get("kind").and_then(|v| v.as_str());
if text.is_none() && kind.is_none() {
return Err("correction requires text or kind".into());
}
if text.is_some_and(|t| t.trim().is_empty()) {
return Err("correction text cannot be empty".into());
}
if let Some(k) = kind {
k.parse::<kimetsu_core::memory::MemoryKind>()?;
}
let now = ts_text(event)?;
let effective = event
.payload
.get("effective_at")
.and_then(|v| v.as_str())
.unwrap_or(&now);
OffsetDateTime::parse(effective, &Rfc3339)?;
conn.execute("INSERT OR IGNORE INTO memory_revisions (memory_id,event_id,text,kind,known_at,effective_at,confidence,use_count,usefulness_score)
SELECT memory_id, 'baseline:' || memory_id, text,kind,created_at,COALESCE(valid_from,'0001-01-01T00:00:00Z'),confidence,use_count,usefulness_score FROM memories WHERE memory_id=?1", params![id])?;
conn.execute(
"UPDATE memory_revisions SET
confidence=(SELECT confidence FROM memories WHERE memory_id=?1),
use_count=(SELECT use_count FROM memories WHERE memory_id=?1),
usefulness_score=(SELECT usefulness_score FROM memories WHERE memory_id=?1)
WHERE revision_id=(SELECT MAX(revision_id) FROM memory_revisions WHERE memory_id=?1)",
params![id],
)?;
let changed = text.is_some_and(|t| t != old_text);
if changed {
conn.execute("UPDATE memories SET confidence=1.0,use_count=0,usefulness_score=0,last_used_at=NULL,last_useful_at=NULL,embedding=NULL,embedding_model=NULL WHERE memory_id=?1", params![id])?;
conn.execute(
"DELETE FROM memory_citations WHERE memory_id=?1",
params![id],
)?;
conn.execute("DELETE FROM query_routes WHERE memory_id=?1", params![id])?;
}
let text = text.unwrap_or(&old_text);
let kind = kind.unwrap_or(&old_kind);
conn.execute(
"UPDATE memories SET text=?2,normalized_text=?3,kind=?4 WHERE memory_id=?1",
params![
id,
text,
kimetsu_core::memory::normalize_memory_text(text),
kind
],
)?;
conn.execute("DELETE FROM memories_fts WHERE memory_id=?1", params![id])?;
conn.execute("INSERT INTO memories_fts(memory_id,text,kind,scope) SELECT memory_id,text,kind,scope FROM memories WHERE memory_id=?1", params![id])?;
crate::graph::project_entities(conn, id, text)?;
conn.execute("INSERT INTO memory_revisions (memory_id,event_id,text,kind,known_at,effective_at,confidence,use_count,usefulness_score)
SELECT memory_id,?2,text,kind,?3,?4,confidence,use_count,usefulness_score FROM memories WHERE memory_id=?1", params![id,event.event_id.to_string(),now,effective])?;
crate::fact_store::refresh(conn, id)?;
Ok(())
}
#[cfg(test)]
mod correction_regressions {
use super::*;
fn event(kind: &str, payload: serde_json::Value, at: &str) -> Event {
let mut e = Event::new(RunId::new(), kind, payload);
e.ts = OffsetDateTime::parse(at, &Rfc3339).unwrap();
e
}
fn seed(c: &Connection) {
schema::initialize(c).unwrap();
apply_events(c, &[event("memory.accepted", serde_json::json!({"memory_id":"m", "scope":"project", "kind":"fact", "text":"original quokka"}), "2026-01-01T00:00:00Z")]).unwrap();
}
#[test]
fn explicit_unbound_exposure_never_credits_a_claim() {
for corrected in [false, true] {
for bindings in [
serde_json::json!({}),
serde_json::json!({"other":"baseline:other"}),
] {
let c = Connection::open_in_memory().unwrap();
seed(&c);
if corrected {
apply_events(
&c,
&[event(
"memory.corrected",
serde_json::json!({"memory_id":"m","text":"new claim"}),
"2026-01-02T00:00:00Z",
)],
)
.unwrap();
}
let run = RunId::new();
let mut events = vec![
event(
"run.started",
serde_json::json!({"project_id":"p","task":"t"}),
"2026-01-03T00:00:00Z",
),
event(
"context.injected",
serde_json::json!({"memory_ids":["m"],"memory_revisions":bindings}),
"2026-01-04T00:00:00Z",
),
event(
"memory.cited",
serde_json::json!({"memory_id":"m","turn":1}),
"2026-01-05T00:00:00Z",
),
event(
"run.finished",
serde_json::json!({"total_cost_usd":0}),
"2026-01-06T00:00:00Z",
),
];
for e in &mut events {
e.run_id = run;
}
apply_events(&c, &events).unwrap();
for _ in 0..2 {
let current: (i64, f64) = c
.query_row(
"SELECT use_count,usefulness_score FROM memories WHERE memory_id='m'",
[],
|r| Ok((r.get(0)?, r.get(1)?)),
)
.unwrap();
assert_eq!(
current,
(0, 0.0),
"explicit unbound exposure must not fall back to current claim"
);
let citations: i64 = c
.query_row("SELECT count(*) FROM memory_citations", [], |r| r.get(0))
.unwrap();
assert_eq!(citations, 0);
let revision_uses: i64 = c
.query_row(
"SELECT COALESCE(SUM(use_count),0) FROM memory_revisions",
[],
|r| r.get(0),
)
.unwrap();
assert_eq!(revision_uses, 0);
rebuild_in_place(&c).unwrap();
}
}
}
}
#[test]
fn delayed_run_evidence_stays_on_the_retiring_claim() {
let c = Connection::open_in_memory().unwrap();
seed(&c);
let run = RunId::new();
let mut events = vec![
event(
"run.started",
serde_json::json!({"project_id":"p","task":"t"}),
"2026-01-02T00:00:00Z",
),
event(
"context.injected",
serde_json::json!({"memory_ids":["m"]}),
"2026-01-03T00:00:00Z",
),
event(
"memory.cited",
serde_json::json!({"memory_id":"m","turn":1}),
"2026-01-04T00:00:00Z",
),
event(
"memory.corrected",
serde_json::json!({"memory_id":"m","text":"new claim"}),
"2026-01-05T00:00:00Z",
),
event(
"memory.cited",
serde_json::json!({"memory_id":"m","turn":2}),
"2026-01-06T00:00:00Z",
),
event(
"run.finished",
serde_json::json!({"total_cost_usd":0}),
"2026-01-07T00:00:00Z",
),
];
for e in &mut events {
e.run_id = run;
e.ts = OffsetDateTime::parse("2026-01-02T00:00:00Z", &Rfc3339).unwrap();
}
apply_events(&c, &events).unwrap();
for _ in 0..2 {
let current: (i64, f64) = c
.query_row(
"SELECT use_count,usefulness_score FROM memories WHERE memory_id='m'",
[],
|r| Ok((r.get(0)?, r.get(1)?)),
)
.unwrap();
assert_eq!(current, (0, 0.0));
let retired:(i64,f64)=c.query_row("SELECT use_count,usefulness_score FROM memory_revisions WHERE event_id='baseline:m'",[],|r|Ok((r.get(0)?,r.get(1)?))).unwrap();
assert_eq!(retired, (1, 1.0));
rebuild_in_place(&c).unwrap();
}
c.execute("UPDATE events SET payload_json=json_remove(payload_json,'$.memory_revisions') WHERE kind='context.injected'",[]).unwrap();
rebuild_in_place(&c).unwrap();
assert_eq!(
c.query_row(
"SELECT use_count FROM memories WHERE memory_id='m'",
[],
|r| r.get::<_, i64>(0)
)
.unwrap(),
0
);
}
#[test]
fn correction_validation_rolls_back_events_text_and_fts() {
let c = Connection::open_in_memory().unwrap();
seed(&c);
let before: i64 = c
.query_row("SELECT COUNT(*) FROM events", [], |r| r.get(0))
.unwrap();
let bad = event(
"memory.corrected",
serde_json::json!({"memory_id":"m", "text":"changed narwhal", "kind":"invalid-kind"}),
"2026-03-01T00:00:00Z",
);
assert!(apply_events(&c, &[bad]).is_err());
assert_eq!(
c.query_row("SELECT COUNT(*) FROM events", [], |r| r.get::<_, i64>(0))
.unwrap(),
before
);
assert_eq!(
c.query_row("SELECT text FROM memories", [], |r| r.get::<_, String>(0))
.unwrap(),
"original quokka"
);
assert_eq!(
c.query_row(
"SELECT COUNT(*) FROM memories_fts WHERE memories_fts MATCH 'quokka'",
[],
|r| r.get::<_, i64>(0)
)
.unwrap(),
1
);
c.execute_batch("CREATE TRIGGER fail_correction BEFORE INSERT ON memory_revisions WHEN NEW.event_id NOT LIKE 'baseline:%' BEGIN SELECT RAISE(ABORT,'injected failure'); END;").unwrap();
let valid = event(
"memory.corrected",
serde_json::json!({"memory_id":"m", "text":"changed narwhal"}),
"2026-03-01T00:00:00Z",
);
assert!(apply_events(&c, &[valid]).is_err());
assert_eq!(
c.query_row("SELECT text FROM memories", [], |r| r.get::<_, String>(0))
.unwrap(),
"original quokka"
);
assert_eq!(
c.query_row(
"SELECT COUNT(*) FROM memories_fts WHERE memories_fts MATCH 'quokka'",
[],
|r| r.get::<_, i64>(0)
)
.unwrap(),
1
);
}
#[test]
fn correction_history_separates_known_and_effective_time_and_replays() {
let c = Connection::open_in_memory().unwrap();
seed(&c);
apply_events(&c, &[event("memory.corrected", serde_json::json!({"memory_id":"m", "text":"corrected narwhal", "effective_at":"2026-02-01T00:00:00Z"}), "2026-03-01T00:00:00Z")]).unwrap();
for _ in 0..2 {
assert_eq!(
crate::bitemporal::memories_at(
&c,
"2026-02-15T00:00:00Z",
"2026-02-15T00:00:00Z",
0
)
.unwrap()[0]
.text,
"original quokka"
);
assert_eq!(
crate::bitemporal::memories_at(
&c,
"2026-02-15T00:00:00Z",
"2026-04-01T00:00:00Z",
0
)
.unwrap()[0]
.text,
"corrected narwhal"
);
assert_eq!(
crate::bitemporal::memories_at(
&c,
"2026-01-15T00:00:00Z",
"2026-04-01T00:00:00Z",
0
)
.unwrap()[0]
.text,
"original quokka"
);
rebuild_in_place(&c).unwrap();
}
let (_, jsonl) = crate::sync::export_events(&c, 0, None, false).unwrap();
let imported = Connection::open_in_memory().unwrap();
schema::initialize(&imported).unwrap();
crate::sync::import_events(&imported, &jsonl.unwrap(), false).unwrap();
rebuild_in_place(&imported).unwrap();
assert_eq!(
imported
.query_row("SELECT text FROM memories WHERE memory_id='m'", [], |r| {
r.get::<_, String>(0)
})
.unwrap(),
"corrected narwhal"
);
}
#[test]
fn corpus_revision_observes_existing_embedding_updates_from_another_connection() {
let dir = tempfile::tempdir().unwrap();
let db = dir.path().join("brain.db");
let reader = Connection::open(&db).unwrap();
seed(&reader);
let writer = Connection::open(&db).unwrap();
let revision = || {
reader
.query_row("SELECT revision FROM corpus_revision", [], |r| {
r.get::<_, i64>(0)
})
.unwrap()
};
let before = revision();
writer
.execute(
"UPDATE memories SET embedding=?1,embedding_model='stub' WHERE memory_id='m'",
params![vec![0u8; 8]],
)
.unwrap();
assert!(revision() > before);
let before = revision();
writer
.execute(
"UPDATE memories SET embedding=?1 WHERE memory_id='m'",
params![vec![1u8; 8]],
)
.unwrap();
assert!(revision() > before);
}
}
pub(crate) fn claim_revision_at(
conn: &Connection,
memory_id: &str,
known_at: Option<&str>,
) -> KimetsuResult<String> {
let revision = conn
.query_row(
"SELECT event_id FROM (
SELECT event_id, known_at, revision_id, text,
LAG(text) OVER (ORDER BY revision_id) AS previous_text
FROM memory_revisions WHERE memory_id=?1)
WHERE (previous_text IS NULL OR text != previous_text)
AND (?2 IS NULL OR julianday(known_at)<julianday(?2))
ORDER BY revision_id DESC LIMIT 1",
params![memory_id, known_at],
|r| r.get::<_, String>(0),
)
.optional()?;
Ok(revision.unwrap_or_else(|| format!("baseline:{memory_id}")))
}
enum ClaimExposure {
Absent,
Unbound,
Bound(String),
}
fn exact_claim_exposure(
conn: &Connection,
exposure_id: &str,
run_id: &str,
memory_id: &str,
) -> KimetsuResult<ClaimExposure> {
let payload:Option<String>=conn.query_row("SELECT payload_json FROM events e WHERE event_id=?1 AND run_id=?2 AND kind='context.injected' AND EXISTS(SELECT 1 FROM json_each(e.payload_json,'$.memory_ids') WHERE value=?3)",params![exposure_id,run_id,memory_id],|r|r.get(0)).optional()?;
let Some(payload) = payload else {
return Ok(ClaimExposure::Unbound);
};
let payload: serde_json::Value = serde_json::from_str(&payload)?;
Ok(
match payload["memory_revisions"][memory_id]
.as_str()
.filter(|r| !r.is_empty())
{
Some(r) => ClaimExposure::Bound(r.to_string()),
None => ClaimExposure::Unbound,
},
)
}
fn run_claim_revision(
conn: &Connection,
run_id: &str,
memory_id: &str,
) -> KimetsuResult<ClaimExposure> {
let mut stmt=conn.prepare("SELECT ts,payload_json FROM events e WHERE run_id=?1 AND kind='context.injected' AND EXISTS(SELECT 1 FROM json_each(e.payload_json,'$.memory_ids') WHERE value=?2) ORDER BY julianday(ts),rowid")?;
let rows = stmt
.query_map(params![run_id, memory_id], |r| {
Ok((r.get::<_, String>(0)?, r.get::<_, String>(1)?))
})?
.collect::<Result<Vec<_>, _>>()?;
let mut bound: Option<String> = None;
for (at, payload) in rows {
let payload: serde_json::Value = serde_json::from_str(&payload)?;
let revision = if let Some(map) = payload.get("memory_revisions") {
let Some(r) = map
.get(memory_id)
.and_then(|v| v.as_str())
.filter(|r| !r.is_empty())
else {
return Ok(ClaimExposure::Unbound);
};
r.to_string()
} else {
claim_revision_at(conn, memory_id, Some(&at))?
};
if bound.as_ref().is_some_and(|r| r != &revision) {
return Ok(ClaimExposure::Unbound);
};
bound = Some(revision);
}
Ok(bound
.map(ClaimExposure::Bound)
.unwrap_or(ClaimExposure::Absent))
}
fn bind_injected_revisions<'a>(
conn: &Connection,
event: &'a Event,
) -> KimetsuResult<Cow<'a, Event>> {
if event.kind != "context.injected" || event.payload.get("memory_revisions").is_some() {
return Ok(Cow::Borrowed(event));
}
let mut bound = event.clone();
let mut revisions = serde_json::Map::new();
if let Some(ids) = event.payload.get("memory_ids").and_then(|v| v.as_array()) {
for id in ids.iter().filter_map(|v| v.as_str()) {
revisions.insert(
id.to_string(),
serde_json::Value::String(claim_revision_at(conn, id, None)?),
);
}
}
bound.payload["memory_revisions"] = serde_json::Value::Object(revisions);
Ok(Cow::Owned(bound))
}