use std::path::Path;
use kimetsu_core::KimetsuResult;
use kimetsu_core::event::Event;
use kimetsu_core::ids::RunId;
use rusqlite::{Connection, OptionalExtension};
use crate::lock::ProjectLock;
use crate::project::*;
use crate::projector;
use crate::schema;
use crate::trace::TraceWriter;
pub fn abort_run(start: &Path, run_id_str: &str) -> KimetsuResult<()> {
{
let (_paths, _config, ro_conn) = load_project_readonly(start)?;
let row: Option<Option<String>> = ro_conn
.query_row(
"SELECT terminal_kind FROM runs WHERE run_id = ?1",
rusqlite::params![run_id_str],
|row| row.get::<_, Option<String>>(0),
)
.optional()?;
match row {
None => {
return Err(format!("run abort: unknown run_id `{run_id_str}`").into());
}
Some(Some(terminal_kind)) => {
return Err(format!(
"run abort: run `{run_id_str}` is already terminal ({})",
terminal_kind
)
.into());
}
Some(None) => {} }
}
let run_id: RunId = run_id_str
.parse::<ulid::Ulid>()
.map(RunId)
.map_err(|_| format!("run abort: `{run_id_str}` is not a valid ULID run id"))?;
let (paths, _config, conn) = load_project(start)?;
let lock = ProjectLock::acquire(&paths, "run abort", Some(run_id))?;
let (mut writer, _run_paths) = TraceWriter::create(&paths, run_id)?;
let aborted_event = Event::new(
run_id,
"run.aborted",
serde_json::json!({
"reason": "manual_abort_via_cli",
}),
);
writer.append(&aborted_event, true)?;
projector::apply_events(&conn, &[aborted_event])?;
lock.release()?;
crate::lock::clear_force(&paths)?;
Ok(())
}
pub fn log_telemetry_event(
start: &Path,
kind: &str,
payload: serde_json::Value,
) -> KimetsuResult<()> {
let paths = kimetsu_core::paths::ProjectPaths::discover(start)?;
let conn = Connection::open(&paths.brain_db)?;
schema::initialize(&conn)?;
let event = Event::new(RunId::new(), kind, payload);
projector::insert_event(&conn, &event)?;
Ok(())
}
pub fn record_context_exposure(start: &Path, event: &Event) -> KimetsuResult<()> {
if event.kind != "context.injected"
|| event.payload.get("memory_revisions").is_none()
|| event.run_id.0 == ulid::Ulid::nil()
{
return Err("exposure requires final delivered revision map".into());
}
let (paths, _, conn) = load_project(start)?;
let _lock = ProjectLock::acquire(&paths, "context exposure", Some(event.run_id))?;
projector::apply_events(&conn, std::slice::from_ref(event))
}
fn load_exposure(
conn: &Connection,
exposure_id: &str,
) -> KimetsuResult<(RunId, serde_json::Value)> {
let row: Option<(String, String)> = conn
.query_row(
"SELECT run_id,payload_json FROM events WHERE event_id=?1 AND kind='context.injected'",
[exposure_id],
|r| Ok((r.get(0)?, r.get(1)?)),
)
.optional()?;
let (run, payload) = row.ok_or("unknown context exposure")?;
Ok((
RunId(run.parse::<ulid::Ulid>()?),
serde_json::from_str(&payload)?,
))
}
pub fn record_exposure_citation(
start: &Path,
exposure_id: &str,
memory_id: &str,
note: Option<&str>,
) -> KimetsuResult<()> {
let (paths, _, conn) = load_project(start)?;
let _lock = ProjectLock::acquire(&paths, "exposure citation", None)?;
projector::with_write_txn(&conn, |conn| {
let (run, payload) = load_exposure(conn, exposure_id)?;
if !payload["memory_ids"]
.as_array()
.is_some_and(|ids| ids.iter().any(|id| id.as_str() == Some(memory_id)))
{
return Err("memory was not delivered in this exposure".into());
}
let revision = payload["memory_revisions"][memory_id]
.as_str()
.filter(|s| !s.is_empty())
.ok_or("delivered claim is unbound")?;
let exists:bool=conn.query_row("SELECT EXISTS(SELECT 1 FROM events WHERE kind='memory.cited' AND json_extract(payload_json,'$.exposure_id')=?1 AND json_extract(payload_json,'$.memory_id')=?2)",rusqlite::params![exposure_id,memory_id],|r|r.get(0))?;
if exists {
return Ok(());
}
let event = Event::new(
run,
"memory.cited",
serde_json::json!({"memory_id":memory_id,"revision_event_id":revision,"exposure_id":exposure_id,"rationale":note,"evidence_kind":"reliance"}),
);
projector::apply_event(conn, &event)
})
}
pub fn record_exposure_outcome(
start: &Path,
exposure_id: &str,
passed: Option<bool>,
) -> KimetsuResult<usize> {
let Some(passed) = passed else { return Ok(0) };
let (paths, _, conn) = load_project(start)?;
let _lock = ProjectLock::acquire(&paths, "exposure outcome", None)?;
let mut count = 0;
projector::with_write_txn(&conn, |conn| {
let (run, payload) = load_exposure(conn, exposure_id)?;
let exists:bool=conn.query_row("SELECT EXISTS(SELECT 1 FROM events WHERE run_id=?1 AND kind IN ('run.finished','run.failed','run.aborted'))",[run.to_string()],|r|r.get(0))?;
if exists {
return Ok(());
}
count = payload["memory_ids"]
.as_array()
.map(|ids| {
ids.iter()
.filter(|id| {
id.as_str()
.is_some_and(|id| payload["memory_revisions"][id].as_str().is_some())
})
.count()
})
.unwrap_or(0);
let event = Event::new(
run,
if passed { "run.finished" } else { "run.failed" },
serde_json::json!({"exposure_id":exposure_id,"evidence_kind":"outcome_association"}),
);
projector::apply_event(conn, &event)
})?;
Ok(count)
}
pub fn emit_regret_for_cited_memories(start: &Path, events: &[kimetsu_core::event::Event]) {
use crate::dropped_capsule;
use kimetsu_core::paths::{ProjectPaths, user_cache_dir_for};
let cache_dir = match ProjectPaths::discover(start) {
Ok(paths) => user_cache_dir_for(&paths.repo_root),
Err(_) => return,
};
let cited_at = dropped_capsule::now_secs();
for event in events {
if event.kind != "memory.cited" {
continue;
}
let Some(memory_id) = event.payload.get("memory_id").and_then(|v| v.as_str()) else {
continue;
};
let Some(dropped_entry) = dropped_capsule::take_if_dropped(&cache_dir, memory_id, cited_at)
else {
continue;
};
let _ = log_telemetry_event(
start,
"retrieval.regret",
serde_json::json!({
"memory_id": memory_id,
"dropped_at": dropped_entry.dropped_at,
"cited_at": cited_at,
}),
);
}
}
pub fn record_mcp_citation(start: &Path, memory_id: &str, note: Option<&str>) -> KimetsuResult<()> {
record_citations(start, &[memory_id.to_string()], note, None)
}
pub fn record_citations(
start: &Path,
memory_ids: &[String],
note: Option<&str>,
query: Option<&str>,
) -> KimetsuResult<()> {
if memory_ids.is_empty() {
return Ok(());
}
let paths = kimetsu_core::paths::ProjectPaths::discover(start)?;
let conn = Connection::open(&paths.brain_db)?;
schema::initialize(&conn)?;
let store_queries = crate::project::load_config(&paths)
.map(|cfg| cfg.learning.store_queries)
.unwrap_or(false);
let group_run_id = RunId::new();
let mut events = Vec::with_capacity(memory_ids.len());
for (turn, memory_id) in memory_ids.iter().enumerate() {
let mut payload = serde_json::json!({
"memory_id": memory_id,
"turn": turn as i64,
"standalone": true,
});
if let Some(n) = note {
payload["rationale"] = serde_json::json!(n);
}
if let Some(q) = query.filter(|_| store_queries) {
payload["query"] = serde_json::json!(q);
}
events.push(kimetsu_core::event::Event::new(
group_run_id,
"memory.cited",
payload,
));
}
projector::apply_events(&conn, &events)?;
emit_regret_for_cited_memories(start, &events);
Ok(())
}
pub fn record_regret(start: &Path, memory_id: &str) -> KimetsuResult<()> {
let paths = kimetsu_core::paths::ProjectPaths::discover(start)?;
let conn = Connection::open(&paths.brain_db)?;
schema::initialize(&conn)?;
let sentinel_run_id = RunId(ulid::Ulid::nil());
let event = kimetsu_core::event::Event::new(
sentinel_run_id,
"retrieval.regret",
serde_json::json!({ "memory_id": memory_id, "source": "manual" }),
);
projector::apply_events(&conn, std::slice::from_ref(&event))?;
Ok(())
}
pub fn record_set_age(start: &Path, memory_id: &str, days_ago: u32) -> KimetsuResult<()> {
use time::format_description::well_known::Rfc3339;
let paths = kimetsu_core::paths::ProjectPaths::discover(start)?;
let conn = Connection::open(&paths.brain_db)?;
schema::initialize(&conn)?;
let target = time::OffsetDateTime::now_utc() - time::Duration::days(days_ago as i64);
let ts = target.format(&Rfc3339).unwrap_or_default();
let sentinel_run_id = RunId(ulid::Ulid::nil());
let event = kimetsu_core::event::Event::new(
sentinel_run_id,
"memory.aged",
serde_json::json!({ "memory_id": memory_id, "created_at": ts, "last_useful_at": ts }),
);
projector::apply_events(&conn, std::slice::from_ref(&event))?;
Ok(())
}