use std::borrow::Cow;
use std::str::FromStr;
use kimetsu_core::KimetsuResult;
use kimetsu_core::event::Event;
use kimetsu_core::ids::{EventId, RunId};
use rusqlite::{Connection, params};
use time::OffsetDateTime;
use time::format_description::well_known::Rfc3339;
use crate::redact;
use crate::schema;
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)?;
let tx = conn.unchecked_transaction()?;
reset_projection(&tx)?;
for event in &events {
project_event(&tx, event)?;
}
tx.commit()?;
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
FROM events
ORDER BY ts, event_id
",
)?;
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)?;
Ok((
event_id_str,
run_id_str,
ts_str,
kind,
schema_version,
payload_json,
))
})?;
let mut events = Vec::new();
for row in rows {
let (event_id_str, run_id_str, ts_str, kind, schema_version, payload_json) = 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,
});
}
Ok(events)
}
pub fn apply_events(conn: &Connection, events: &[Event]) -> KimetsuResult<()> {
let tx = conn.unchecked_transaction()?;
for event in events {
apply_event(&tx, event)?;
}
tx.commit()?;
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;
",
)?;
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),
"memory.superseded" => apply_memory_superseded(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,
],
)?;
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
)
VALUES (?1, ?2, ?3, ?4, ?5, ?6)
",
params![
event.event_id.to_string(),
event.run_id.to_string(),
ts_text(event)?,
event.kind,
event.schema_version,
payload
],
)?;
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";
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],
)?;
}
}
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 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
)
VALUES (?1, ?2, ?3, ?4, ?5, ?6, ?7, ?8, ?9, 0)
",
params![
memory_id,
scope,
kind,
text,
normalized_text,
confidence,
event.event_id.to_string(),
serde_json::to_string(&provenance_snapshot)?,
ts_text(event)?
],
)?;
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);
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);
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;
use serde_json::json;
use super::{apply_events, 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)
}
#[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 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"
);
}
#[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 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 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");
}
}