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 sentinel_run_id = RunId(ulid::Ulid::nil());
let event = Event::new(sentinel_run_id, kind, payload);
projector::insert_event(&conn, &event)?;
Ok(())
}
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 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 {
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(())
}