use rusqlite::Connection;
use crate::error::{Error, Result};
pub const CHECKPOINT_FORMAT: i64 = 7;
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum RunStatus {
Running,
Paused,
Completed,
Failed,
}
impl RunStatus {
fn from_str(s: &str) -> Self {
match s {
"paused" => RunStatus::Paused,
"completed" => RunStatus::Completed,
"failed" => RunStatus::Failed,
_ => RunStatus::Running,
}
}
}
#[derive(Debug, Clone, PartialEq)]
pub struct SpawnRow {
pub child_run_id: i64,
pub goal: String,
pub verify_file: String,
pub needle: String,
pub max_steps: Option<u32>,
pub deny_write: String,
}
pub struct Store {
conn: Connection,
}
#[derive(Debug, Clone, PartialEq)]
pub struct CheckpointEvent {
pub run_id: i64,
pub step: u32,
pub kind: String,
pub detail: Option<String>,
}
impl CheckpointEvent {
pub fn checkpoint(run_id: i64, step: u32) -> Self {
Self {
run_id,
step,
kind: "checkpoint".into(),
detail: None,
}
}
pub fn resume(run_id: i64, step: u32, detail: impl Into<String>) -> Self {
Self {
run_id,
step,
kind: "resume".into(),
detail: Some(detail.into()),
}
}
pub fn skipped(run_id: i64, step: u32) -> Self {
Self {
run_id,
step,
kind: "skipped".into(),
detail: None,
}
}
}
#[derive(Debug, Clone, PartialEq)]
pub struct StepRecord {
pub step: u32,
pub decision: String,
pub result: String,
pub prompt: String,
pub tool_call: String,
pub tokens: u64,
}
impl StepRecord {
pub fn new(step: u32, decision: impl Into<String>, result: impl Into<String>) -> Self {
Self {
step,
decision: decision.into(),
result: result.into(),
prompt: String::new(),
tool_call: String::new(),
tokens: 0,
}
}
pub fn with_trace(
mut self,
prompt: impl Into<String>,
tool_call: impl Into<String>,
tokens: u64,
) -> Self {
self.prompt = prompt.into();
self.tool_call = tool_call.into();
self.tokens = tokens;
self
}
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct PolicyEvent {
pub step: u32,
pub kind: String,
pub act: String,
pub target: String,
pub rule: Option<String>,
pub layer: Option<String>,
pub decision: Option<String>,
pub source: Option<String>,
pub performed: Option<String>,
}
impl PolicyEvent {
pub fn refusal(step: u32, act: impl Into<String>, target: impl Into<String>) -> Self {
Self {
step,
kind: "refusal".into(),
act: act.into(),
target: target.into(),
rule: None,
layer: None,
decision: None,
source: None,
performed: None,
}
}
pub fn decision(
step: u32,
act: impl Into<String>,
target: impl Into<String>,
decision: impl Into<String>,
source: impl Into<String>,
) -> Self {
Self {
kind: "decision".into(),
decision: Some(decision.into()),
source: Some(source.into()),
..Self::refusal(step, act, target)
}
}
pub fn with_rule(mut self, rule: impl Into<String>, layer: impl Into<String>) -> Self {
self.rule = Some(rule.into());
self.layer = Some(layer.into());
self
}
pub fn with_performed(mut self, performed: impl Into<String>) -> Self {
self.performed = Some(performed.into());
self
}
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct Pending {
pub id: i64,
pub run_id: i64,
pub step: u32,
pub act: String,
pub target: String,
pub content: Option<String>,
pub resolved: Option<String>,
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct AgentEvent {
pub run_id: i64,
pub step: u32,
pub kind: String,
pub child_run_id: Option<i64>,
pub detail: Option<String>,
pub tokens: Option<u64>,
pub remaining: Option<u64>,
}
impl AgentEvent {
pub fn spawn(run_id: i64, step: u32, child_run_id: i64, goal: impl Into<String>) -> Self {
Self {
run_id,
step,
kind: "spawn".into(),
child_run_id: Some(child_run_id),
detail: Some(goal.into()),
tokens: None,
remaining: None,
}
}
pub fn spawn_refused(run_id: i64, step: u32, cap: &str) -> Self {
Self {
run_id,
step,
kind: "spawn_refused".into(),
child_run_id: None,
detail: Some(cap.into()),
tokens: None,
remaining: None,
}
}
pub fn budget_draw(run_id: i64, step: u32, tokens: u64, remaining: u64) -> Self {
Self {
run_id,
step,
kind: "budget_draw".into(),
child_run_id: None,
detail: None,
tokens: Some(tokens),
remaining: Some(remaining),
}
}
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct SandboxEvent {
pub run_id: i64,
pub step: u32,
pub kind: String,
pub backend: Option<String>,
pub detail: Option<String>,
}
impl SandboxEvent {
pub fn create(run_id: i64, step: u32, backend: &str) -> Self {
Self {
run_id,
step,
kind: "create".into(),
backend: Some(backend.into()),
detail: None,
}
}
pub fn exec(run_id: i64, step: u32, backend: &str, argv: &str) -> Self {
Self {
run_id,
step,
kind: "exec".into(),
backend: Some(backend.into()),
detail: Some(argv.into()),
}
}
pub fn cap_hit(run_id: i64, step: u32, cap: &str) -> Self {
Self {
run_id,
step,
kind: "cap_hit".into(),
backend: None,
detail: Some(cap.into()),
}
}
pub fn destroy(run_id: i64, step: u32) -> Self {
Self {
run_id,
step,
kind: "destroy".into(),
backend: None,
detail: None,
}
}
pub fn gate_phase_failed(run_id: i64, step: u32, phase: &str) -> Self {
Self {
run_id,
step,
kind: "gate_phase_failed".into(),
backend: None,
detail: Some(phase.into()),
}
}
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct McpEvent {
pub step: u32,
pub kind: String,
pub server: String,
pub tool: Option<String>,
pub ok: Option<bool>,
pub millis: Option<u64>,
pub detail: Option<String>,
}
impl McpEvent {
fn new(kind: &str, server: &str) -> Self {
Self {
step: 0,
kind: kind.into(),
server: server.into(),
tool: None,
ok: None,
millis: None,
detail: None,
}
}
pub fn connected(server: &str, transport: &str) -> Self {
Self::new("connected", server).with_detail(transport)
}
pub fn discovered(server: &str, tool: &str) -> Self {
let mut e = Self::new("discovered", server);
e.tool = Some(tool.into());
e
}
pub fn called(server: &str, tool: &str, ok: bool) -> Self {
let mut e = Self::new("called", server);
e.tool = Some(tool.into());
e.ok = Some(ok);
e
}
pub fn disconnected(server: &str) -> Self {
Self::new("disconnected", server)
}
pub fn at_step(mut self, step: u32) -> Self {
self.step = step;
self
}
pub fn with_millis(mut self, millis: u64) -> Self {
self.millis = Some(millis);
self
}
pub fn with_detail(mut self, detail: impl Into<String>) -> Self {
let detail = detail.into();
self.detail = (!detail.is_empty()).then_some(detail);
self
}
}
#[derive(Debug, Clone, PartialEq)]
pub struct MemoryEntry {
pub key: String,
pub value: String,
pub run_id: i64,
pub step: u32,
pub created_at: String,
}
pub const MEMORY_MAX_ENTRIES: usize = 64;
pub const MEMORY_MAX_CHARS: usize = 16_000;
pub const MEMORY_MAX_ENTRY_CHARS: usize = MEMORY_MAX_CHARS / 8;
const MEMORY_TRUNCATED: &str = "…[truncated]";
fn truncate_memory_value(value: &str) -> String {
if value.chars().count() <= MEMORY_MAX_ENTRY_CHARS {
return value.to_string();
}
let keep = MEMORY_MAX_ENTRY_CHARS - MEMORY_TRUNCATED.chars().count();
let mut out: String = value.chars().take(keep).collect();
out.push_str(MEMORY_TRUNCATED);
out
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct ContextEvent {
pub step: u32,
pub kind: String,
pub detail: Option<String>,
pub est_tokens: Option<u64>,
pub reported_tokens: Option<u64>,
}
impl ContextEvent {
pub fn assembled(step: u32, detail: impl Into<String>, est_tokens: u64) -> Self {
Self {
step,
kind: "assembled".into(),
detail: Some(detail.into()),
est_tokens: Some(est_tokens),
reported_tokens: None,
}
}
pub fn reread(step: u32, detail: impl Into<String>) -> Self {
Self {
step,
kind: "reread".into(),
detail: Some(detail.into()),
est_tokens: None,
reported_tokens: None,
}
}
pub fn reread_refused(step: u32, detail: impl Into<String>) -> Self {
Self {
step,
kind: "reread_refused".into(),
detail: Some(detail.into()),
est_tokens: None,
reported_tokens: None,
}
}
pub fn memory_write(step: u32, detail: impl Into<String>) -> Self {
Self::of("memory_write", step, detail)
}
pub fn memory_evict(step: u32, detail: impl Into<String>) -> Self {
Self::of("memory_evict", step, detail)
}
pub fn memory_recall(step: u32, detail: impl Into<String>) -> Self {
Self::of("memory_recall", step, detail)
}
fn of(kind: &str, step: u32, detail: impl Into<String>) -> Self {
Self {
step,
kind: kind.into(),
detail: Some(detail.into()),
est_tokens: None,
reported_tokens: None,
}
}
}
impl Store {
pub fn open(path: impl AsRef<std::path::Path>) -> Result<Self> {
Self::from_conn(Connection::open(path)?)
}
pub fn memory() -> Result<Self> {
Self::from_conn(Connection::open_in_memory()?)
}
fn from_conn(conn: Connection) -> Result<Self> {
conn.execute_batch(
"CREATE TABLE IF NOT EXISTS runs (
id INTEGER PRIMARY KEY AUTOINCREMENT,
goal TEXT NOT NULL,
file TEXT NOT NULL,
outcome TEXT,
provider TEXT
);
CREATE TABLE IF NOT EXISTS steps (
id INTEGER PRIMARY KEY AUTOINCREMENT,
run_id INTEGER NOT NULL REFERENCES runs(id),
step INTEGER NOT NULL,
decision TEXT NOT NULL,
result TEXT NOT NULL,
prompt TEXT NOT NULL DEFAULT '',
tool_call TEXT NOT NULL DEFAULT '',
tokens INTEGER NOT NULL DEFAULT 0
);",
)?;
for col in [
"prompt TEXT NOT NULL DEFAULT ''",
"tool_call TEXT NOT NULL DEFAULT ''",
"tokens INTEGER NOT NULL DEFAULT 0",
] {
let _ = conn.execute(&format!("ALTER TABLE steps ADD COLUMN {col}"), []);
}
let _ = conn.execute("ALTER TABLE runs ADD COLUMN provider TEXT", []);
conn.execute_batch(
"CREATE TABLE IF NOT EXISTS policy_events (
id INTEGER PRIMARY KEY AUTOINCREMENT,
run_id INTEGER NOT NULL,
step INTEGER NOT NULL,
kind TEXT NOT NULL,
act TEXT NOT NULL,
target TEXT NOT NULL,
rule TEXT,
layer TEXT,
decision TEXT,
source TEXT,
performed TEXT
);
CREATE TABLE IF NOT EXISTS pending_approvals (
id INTEGER PRIMARY KEY AUTOINCREMENT,
run_id INTEGER NOT NULL,
step INTEGER NOT NULL,
act TEXT NOT NULL,
target TEXT NOT NULL,
content TEXT,
resolved TEXT
);",
)?;
let _ = conn.execute("ALTER TABLE runs ADD COLUMN parent_run_id INTEGER", []);
let _ = conn.execute(
"ALTER TABLE runs ADD COLUMN depth INTEGER NOT NULL DEFAULT 0",
[],
);
conn.execute_batch(
"CREATE TABLE IF NOT EXISTS agent_events (
id INTEGER PRIMARY KEY AUTOINCREMENT,
run_id INTEGER NOT NULL,
step INTEGER NOT NULL,
kind TEXT NOT NULL,
child_run_id INTEGER,
detail TEXT,
tokens INTEGER,
remaining INTEGER
);",
)?;
conn.execute_batch(
"CREATE TABLE IF NOT EXISTS sandbox_events (
id INTEGER PRIMARY KEY AUTOINCREMENT,
run_id INTEGER NOT NULL,
step INTEGER NOT NULL,
kind TEXT NOT NULL,
backend TEXT,
detail TEXT
);",
)?;
let _ = conn.execute(
"ALTER TABLE runs ADD COLUMN status TEXT NOT NULL DEFAULT 'running'",
[],
);
let _ = conn.execute("ALTER TABLE runs ADD COLUMN started_at TEXT", []);
conn.execute_batch(
"CREATE TABLE IF NOT EXISTS checkpoint_events (
id INTEGER PRIMARY KEY AUTOINCREMENT,
run_id INTEGER NOT NULL,
step INTEGER NOT NULL,
kind TEXT NOT NULL,
detail TEXT
);
CREATE TABLE IF NOT EXISTS spawns (
id INTEGER PRIMARY KEY AUTOINCREMENT,
parent_run_id INTEGER NOT NULL,
step INTEGER NOT NULL,
child_run_id INTEGER NOT NULL,
goal TEXT NOT NULL,
verify_file TEXT NOT NULL,
needle TEXT NOT NULL,
max_steps INTEGER,
deny_write TEXT NOT NULL DEFAULT '[]'
);",
)?;
conn.execute_batch(
"CREATE TABLE IF NOT EXISTS mcp_events (
id INTEGER PRIMARY KEY AUTOINCREMENT,
run_id INTEGER NOT NULL,
step INTEGER NOT NULL,
kind TEXT NOT NULL,
server TEXT NOT NULL,
tool TEXT,
ok INTEGER,
millis INTEGER,
detail TEXT
);",
)?;
conn.execute_batch(
"CREATE TABLE IF NOT EXISTS memory (
id INTEGER PRIMARY KEY,
workspace TEXT NOT NULL,
key TEXT NOT NULL,
value TEXT NOT NULL,
run_id INTEGER NOT NULL,
step INTEGER NOT NULL,
created_at TEXT NOT NULL,
UNIQUE(workspace, key)
);",
)?;
conn.execute_batch(
"CREATE TABLE IF NOT EXISTS context_events (
id INTEGER PRIMARY KEY,
run_id INTEGER NOT NULL,
step INTEGER NOT NULL,
kind TEXT NOT NULL,
detail TEXT,
est_tokens INTEGER,
reported_tokens INTEGER
);",
)?;
let format: i64 = conn.query_row("PRAGMA user_version", [], |r| r.get(0))?;
if format < CHECKPOINT_FORMAT {
conn.execute_batch(&format!("PRAGMA user_version = {CHECKPOINT_FORMAT}"))?;
}
Ok(Self { conn })
}
pub fn record_event(&self, run_id: i64, e: &PolicyEvent) -> Result<()> {
self.conn.execute(
"INSERT INTO policy_events
(run_id, step, kind, act, target, rule, layer, decision, source, performed)
VALUES (?1, ?2, ?3, ?4, ?5, ?6, ?7, ?8, ?9, ?10)",
(
run_id,
e.step,
&e.kind,
&e.act,
&e.target,
&e.rule,
&e.layer,
&e.decision,
&e.source,
&e.performed,
),
)?;
Ok(())
}
pub fn events(&self, run_id: i64) -> Result<Vec<PolicyEvent>> {
let mut stmt = self.conn.prepare(
"SELECT step, kind, act, target, rule, layer, decision, source, performed
FROM policy_events WHERE run_id = ?1 ORDER BY id ASC",
)?;
let rows = stmt.query_map([run_id], |r| {
Ok(PolicyEvent {
step: r.get::<_, i64>(0)? as u32,
kind: r.get(1)?,
act: r.get(2)?,
target: r.get(3)?,
rule: r.get(4)?,
layer: r.get(5)?,
decision: r.get(6)?,
source: r.get(7)?,
performed: r.get(8)?,
})
})?;
Ok(rows.collect::<std::result::Result<_, _>>()?)
}
pub fn put_pending(
&self,
run_id: i64,
step: u32,
act: &str,
target: &str,
content: Option<&str>,
) -> Result<i64> {
self.conn.execute(
"INSERT INTO pending_approvals (run_id, step, act, target, content)
VALUES (?1, ?2, ?3, ?4, ?5)",
(run_id, step, act, target, content),
)?;
Ok(self.conn.last_insert_rowid())
}
pub fn pending(&self, request_id: i64) -> Result<Option<Pending>> {
let mut stmt = self.conn.prepare(
"SELECT id, run_id, step, act, target, content, resolved
FROM pending_approvals WHERE id = ?1",
)?;
let mut rows = stmt.query_map([request_id], |r| {
Ok(Pending {
id: r.get(0)?,
run_id: r.get(1)?,
step: r.get::<_, i64>(2)? as u32,
act: r.get(3)?,
target: r.get(4)?,
content: r.get(5)?,
resolved: r.get(6)?,
})
})?;
Ok(rows.next().transpose()?)
}
pub fn resolve_pending(&self, request_id: i64, decision: &str) -> Result<()> {
self.conn.execute(
"UPDATE pending_approvals SET resolved = ?1 WHERE id = ?2",
(decision, request_id),
)?;
Ok(())
}
pub fn start_run(&self, goal: &str, file: &str) -> Result<i64> {
self.conn.execute(
"INSERT INTO runs (goal, file, status, started_at)
VALUES (?1, ?2, 'running', strftime('%Y-%m-%dT%H:%M:%fZ','now'))",
(goal, file),
)?;
Ok(self.conn.last_insert_rowid())
}
pub fn start_child_run(
&self,
goal: &str,
file: &str,
parent_run_id: i64,
depth: u32,
) -> Result<i64> {
self.conn.execute(
"INSERT INTO runs (goal, file, parent_run_id, depth, status, started_at)
VALUES (?1, ?2, ?3, ?4, 'running', strftime('%Y-%m-%dT%H:%M:%fZ','now'))",
(goal, file, parent_run_id, depth),
)?;
Ok(self.conn.last_insert_rowid())
}
pub fn record_agent_event(&self, e: &AgentEvent) -> Result<()> {
self.conn.execute(
"INSERT INTO agent_events (run_id, step, kind, child_run_id, detail, tokens, remaining)
VALUES (?1, ?2, ?3, ?4, ?5, ?6, ?7)",
(
e.run_id,
e.step,
&e.kind,
e.child_run_id,
&e.detail,
e.tokens,
e.remaining,
),
)?;
Ok(())
}
pub fn agent_events(&self, run_id: i64) -> Result<Vec<AgentEvent>> {
let mut stmt = self.conn.prepare(
"SELECT run_id, step, kind, child_run_id, detail, tokens, remaining
FROM agent_events WHERE run_id = ?1 ORDER BY id ASC",
)?;
let rows = stmt.query_map([run_id], |r| {
Ok(AgentEvent {
run_id: r.get(0)?,
step: r.get::<_, i64>(1)? as u32,
kind: r.get(2)?,
child_run_id: r.get(3)?,
detail: r.get(4)?,
tokens: r.get::<_, Option<i64>>(5)?.map(|n| n as u64),
remaining: r.get::<_, Option<i64>>(6)?.map(|n| n as u64),
})
})?;
Ok(rows.collect::<std::result::Result<_, _>>()?)
}
pub fn record_sandbox_event(&self, e: &SandboxEvent) -> Result<()> {
self.conn.execute(
"INSERT INTO sandbox_events (run_id, step, kind, backend, detail)
VALUES (?1, ?2, ?3, ?4, ?5)",
(e.run_id, e.step, &e.kind, &e.backend, &e.detail),
)?;
Ok(())
}
pub fn record_mcp(&self, run_id: i64, e: &McpEvent) -> Result<()> {
self.conn.execute(
"INSERT INTO mcp_events (run_id, step, kind, server, tool, ok, millis, detail)
VALUES (?1, ?2, ?3, ?4, ?5, ?6, ?7, ?8)",
(
run_id,
e.step,
&e.kind,
&e.server,
&e.tool,
e.ok,
e.millis.map(|m| m as i64),
&e.detail,
),
)?;
Ok(())
}
pub fn mcp_events(&self, run_id: i64) -> Result<Vec<McpEvent>> {
let mut stmt = self.conn.prepare(
"SELECT step, kind, server, tool, ok, millis, detail
FROM mcp_events WHERE run_id = ?1 ORDER BY id ASC",
)?;
let rows = stmt.query_map([run_id], |r| {
Ok(McpEvent {
step: r.get::<_, i64>(0)? as u32,
kind: r.get(1)?,
server: r.get(2)?,
tool: r.get(3)?,
ok: r.get(4)?,
millis: r.get::<_, Option<i64>>(5)?.map(|m| m as u64),
detail: r.get(6)?,
})
})?;
Ok(rows.collect::<std::result::Result<_, _>>()?)
}
pub fn record_context_event(&self, run_id: i64, e: &ContextEvent) -> Result<()> {
self.conn.execute(
"INSERT INTO context_events (run_id, step, kind, detail, est_tokens, reported_tokens)
VALUES (?1, ?2, ?3, ?4, ?5, ?6)",
(
run_id,
e.step,
&e.kind,
&e.detail,
e.est_tokens.map(|n| n as i64),
e.reported_tokens.map(|n| n as i64),
),
)?;
Ok(())
}
pub fn record_context_reported(&self, run_id: i64, step: u32, reported: u64) -> Result<()> {
self.conn.execute(
"UPDATE context_events SET reported_tokens = ?1
WHERE run_id = ?2 AND step = ?3 AND kind = 'assembled'",
(reported as i64, run_id, step),
)?;
Ok(())
}
pub fn context_events(&self, run_id: i64) -> Result<Vec<ContextEvent>> {
let mut stmt = self.conn.prepare(
"SELECT step, kind, detail, est_tokens, reported_tokens
FROM context_events WHERE run_id = ?1 ORDER BY id ASC",
)?;
let rows = stmt.query_map([run_id], |r| {
Ok(ContextEvent {
step: r.get::<_, i64>(0)? as u32,
kind: r.get(1)?,
detail: r.get(2)?,
est_tokens: r.get::<_, Option<i64>>(3)?.map(|n| n as u64),
reported_tokens: r.get::<_, Option<i64>>(4)?.map(|n| n as u64),
})
})?;
Ok(rows.collect::<std::result::Result<_, _>>()?)
}
pub fn sandbox_events(&self, run_id: i64) -> Result<Vec<SandboxEvent>> {
let mut stmt = self.conn.prepare(
"SELECT run_id, step, kind, backend, detail
FROM sandbox_events WHERE run_id = ?1 ORDER BY id ASC",
)?;
let rows = stmt.query_map([run_id], |r| {
Ok(SandboxEvent {
run_id: r.get(0)?,
step: r.get::<_, i64>(1)? as u32,
kind: r.get(2)?,
backend: r.get(3)?,
detail: r.get(4)?,
})
})?;
Ok(rows.collect::<std::result::Result<_, _>>()?)
}
pub fn children(&self, run_id: i64) -> Result<Vec<i64>> {
let mut stmt = self
.conn
.prepare("SELECT id FROM runs WHERE parent_run_id = ?1 ORDER BY id ASC")?;
let rows = stmt.query_map([run_id], |r| r.get(0))?;
Ok(rows.collect::<std::result::Result<_, _>>()?)
}
pub fn parent(&self, run_id: i64) -> Result<Option<i64>> {
Ok(self.conn.query_row(
"SELECT parent_run_id FROM runs WHERE id = ?1",
[run_id],
|r| r.get(0),
)?)
}
pub fn depth(&self, run_id: i64) -> Result<u32> {
let d: i64 =
self.conn
.query_row("SELECT depth FROM runs WHERE id = ?1", [run_id], |r| {
r.get(0)
})?;
Ok(d as u32)
}
pub fn record(&self, run_id: i64, step: &StepRecord) -> Result<()> {
self.conn.execute(
"INSERT INTO steps (run_id, step, decision, result, prompt, tool_call, tokens)
VALUES (?1, ?2, ?3, ?4, ?5, ?6, ?7)",
(
run_id,
step.step,
&step.decision,
&step.result,
&step.prompt,
&step.tool_call,
step.tokens,
),
)?;
Ok(())
}
pub fn checkpoint_step(&self, run_id: i64, step: &StepRecord) -> Result<()> {
let tx = self.conn.unchecked_transaction()?;
tx.execute(
"INSERT INTO steps (run_id, step, decision, result, prompt, tool_call, tokens)
VALUES (?1, ?2, ?3, ?4, ?5, ?6, ?7)",
(
run_id,
step.step,
&step.decision,
&step.result,
&step.prompt,
&step.tool_call,
step.tokens,
),
)?;
tx.execute(
"INSERT INTO checkpoint_events (run_id, step, kind, detail)
VALUES (?1, ?2, 'checkpoint', NULL)",
(run_id, step.step),
)?;
tx.commit()?;
Ok(())
}
pub fn record_checkpoint_event(&self, e: &CheckpointEvent) -> Result<()> {
self.conn.execute(
"INSERT INTO checkpoint_events (run_id, step, kind, detail) VALUES (?1, ?2, ?3, ?4)",
(e.run_id, e.step, &e.kind, &e.detail),
)?;
Ok(())
}
pub fn checkpoint_events(&self, run_id: i64) -> Result<Vec<CheckpointEvent>> {
let mut stmt = self.conn.prepare(
"SELECT run_id, step, kind, detail
FROM checkpoint_events WHERE run_id = ?1 ORDER BY id ASC",
)?;
let rows = stmt.query_map([run_id], |r| {
Ok(CheckpointEvent {
run_id: r.get(0)?,
step: r.get::<_, i64>(1)? as u32,
kind: r.get(2)?,
detail: r.get(3)?,
})
})?;
Ok(rows.collect::<std::result::Result<_, _>>()?)
}
pub fn set_status(&self, run_id: i64, status: &str) -> Result<()> {
self.conn.execute(
"UPDATE runs SET status = ?1 WHERE id = ?2",
(status, run_id),
)?;
Ok(())
}
pub fn status(&self, run_id: i64) -> Result<Option<String>> {
Ok(self
.conn
.query_row("SELECT status FROM runs WHERE id = ?1", [run_id], |r| {
r.get(0)
})
.ok())
}
pub fn elapsed_secs(&self, run_id: i64) -> Result<f64> {
let secs: Option<f64> = self.conn.query_row(
"SELECT (julianday('now') - julianday(started_at)) * 86400.0
FROM runs WHERE id = ?1",
[run_id],
|r| r.get(0),
)?;
Ok(secs.unwrap_or(0.0).max(0.0))
}
pub fn spent_tokens(&self, run_id: i64) -> Result<u64> {
let n: i64 = self.conn.query_row(
"SELECT COALESCE(SUM(tokens), 0) FROM steps WHERE run_id = ?1",
[run_id],
|r| r.get(0),
)?;
Ok(n as u64)
}
pub fn tree_run_ids(&self, root: i64) -> Result<Vec<i64>> {
let mut stmt = self.conn.prepare(
"WITH RECURSIVE tree(id) AS (
SELECT id FROM runs WHERE id = ?1
UNION ALL
SELECT r.id FROM runs r JOIN tree t ON r.parent_run_id = t.id
)
SELECT id FROM tree ORDER BY id ASC",
)?;
let rows = stmt.query_map([root], |r| r.get(0))?;
Ok(rows.collect::<std::result::Result<_, _>>()?)
}
pub fn spent_tokens_tree(&self, root: i64) -> Result<u64> {
let n: i64 = self.conn.query_row(
"WITH RECURSIVE tree(id) AS (
SELECT id FROM runs WHERE id = ?1
UNION ALL
SELECT r.id FROM runs r JOIN tree t ON r.parent_run_id = t.id
)
SELECT COALESCE(SUM(s.tokens), 0)
FROM steps s JOIN tree ON s.run_id = tree.id",
[root],
|r| r.get(0),
)?;
Ok(n as u64)
}
pub fn agent_count_tree(&self, root: i64) -> Result<u32> {
let n: i64 = self.conn.query_row(
"WITH RECURSIVE tree(id) AS (
SELECT id FROM runs WHERE id = ?1
UNION ALL
SELECT r.id FROM runs r JOIN tree t ON r.parent_run_id = t.id
)
SELECT COUNT(*) FROM tree",
[root],
|r| r.get(0),
)?;
Ok(n as u32)
}
#[allow(clippy::too_many_arguments)]
pub fn record_spawn(
&self,
parent_run_id: i64,
step: u32,
child_run_id: i64,
goal: &str,
verify_file: &str,
needle: &str,
max_steps: Option<u32>,
deny_write_json: &str,
) -> Result<()> {
self.conn.execute(
"INSERT INTO spawns
(parent_run_id, step, child_run_id, goal, verify_file, needle, max_steps, deny_write)
VALUES (?1, ?2, ?3, ?4, ?5, ?6, ?7, ?8)",
(
parent_run_id,
step,
child_run_id,
goal,
verify_file,
needle,
max_steps,
deny_write_json,
),
)?;
Ok(())
}
pub fn find_spawn(
&self,
parent_run_id: i64,
step: u32,
goal: &str,
) -> Result<Option<SpawnRow>> {
Ok(self
.conn
.query_row(
"SELECT child_run_id, goal, verify_file, needle, max_steps, deny_write
FROM spawns WHERE parent_run_id = ?1 AND step = ?2 AND goal = ?3
ORDER BY id ASC LIMIT 1",
(parent_run_id, step, goal),
|r| {
Ok(SpawnRow {
child_run_id: r.get(0)?,
goal: r.get(1)?,
verify_file: r.get(2)?,
needle: r.get(3)?,
max_steps: r.get::<_, Option<i64>>(4)?.map(|n| n as u32),
deny_write: r.get(5)?,
})
},
)
.ok())
}
pub fn check_resumable(&self, run_id: i64) -> Result<()> {
let format: i64 = self
.conn
.query_row("PRAGMA user_version", [], |r| r.get(0))?;
if format > CHECKPOINT_FORMAT {
return Err(Error::Resume {
reason: format!(
"checkpoint format {format} is newer than supported {CHECKPOINT_FORMAT}; \
upgrade io-harness to resume this run"
),
});
}
let exists: bool = self
.conn
.query_row("SELECT 1 FROM runs WHERE id = ?1", [run_id], |_| Ok(true))
.unwrap_or(false);
if !exists {
return Err(Error::Resume {
reason: format!("no run with id {run_id} in the store"),
});
}
Ok(())
}
pub fn set_provider(&self, run_id: i64, provider: &str) -> Result<()> {
self.conn.execute(
"UPDATE runs SET provider = ?1 WHERE id = ?2",
(provider, run_id),
)?;
Ok(())
}
pub fn provider(&self, run_id: i64) -> Result<Option<String>> {
Ok(self
.conn
.query_row("SELECT provider FROM runs WHERE id = ?1", [run_id], |r| {
r.get(0)
})?)
}
pub fn finish_run(&self, run_id: i64, outcome: &str) -> Result<()> {
let status = match outcome {
"awaiting_approval" => "paused",
_ => "completed",
};
self.conn.execute(
"UPDATE runs SET outcome = ?1, status = ?2 WHERE id = ?3",
(outcome, status, run_id),
)?;
Ok(())
}
pub fn outcome(&self, run_id: i64) -> Result<Option<String>> {
Ok(self
.conn
.query_row("SELECT outcome FROM runs WHERE id = ?1", [run_id], |r| {
r.get(0)
})
.ok()
.flatten())
}
pub fn run_status(&self, run_id: i64) -> Result<Option<RunStatus>> {
Ok(self.status(run_id)?.map(|s| RunStatus::from_str(&s)))
}
pub fn last_step(&self, run_id: i64) -> Result<u32> {
let n: i64 = self.conn.query_row(
"SELECT COALESCE(MAX(step), 0) FROM steps WHERE run_id = ?1",
[run_id],
|r| r.get(0),
)?;
Ok(n as u32)
}
pub fn steps(&self, run_id: i64) -> Result<Vec<StepRecord>> {
let mut stmt = self.conn.prepare(
"SELECT step, decision, result, prompt, tool_call, tokens
FROM steps WHERE run_id = ?1 ORDER BY step ASC, id ASC",
)?;
let rows = stmt.query_map([run_id], |r| {
Ok(StepRecord {
step: r.get::<_, i64>(0)? as u32,
decision: r.get(1)?,
result: r.get(2)?,
prompt: r.get(3)?,
tool_call: r.get(4)?,
tokens: r.get::<_, i64>(5)? as u64,
})
})?;
Ok(rows.collect::<std::result::Result<_, _>>()?)
}
pub fn memory_put(
&self,
workspace: &str,
key: &str,
value: &str,
run_id: i64,
step: u32,
) -> Result<Vec<String>> {
let value = truncate_memory_value(value);
self.conn.execute(
"INSERT INTO memory (workspace, key, value, run_id, step, created_at)
VALUES (?1, ?2, ?3, ?4, ?5, strftime('%Y-%m-%dT%H:%M:%fZ','now'))
ON CONFLICT(workspace, key) DO UPDATE SET
value = excluded.value,
run_id = excluded.run_id,
step = excluded.step,
created_at = excluded.created_at",
(workspace, key, &value, run_id, step),
)?;
self.enforce_memory_caps(workspace, key)
}
fn enforce_memory_caps(&self, workspace: &str, keep: &str) -> Result<Vec<String>> {
let rows: Vec<(String, i64)> = {
let mut stmt = self.conn.prepare(
"SELECT key, LENGTH(value) FROM memory WHERE workspace = ?1
ORDER BY created_at ASC, id ASC",
)?;
let rows = stmt.query_map([workspace], |r| Ok((r.get(0)?, r.get(1)?)))?;
rows.collect::<std::result::Result<_, _>>()?
};
let mut count = rows.len();
let mut chars: i64 = rows.iter().map(|(_, n)| *n).sum();
let mut evicted = Vec::new();
for (key, n) in &rows {
if count <= MEMORY_MAX_ENTRIES && chars <= MEMORY_MAX_CHARS as i64 {
break;
}
if key == keep {
continue;
}
self.conn.execute(
"DELETE FROM memory WHERE workspace = ?1 AND key = ?2",
(workspace, key),
)?;
count -= 1;
chars -= n;
evicted.push(key.clone());
}
Ok(evicted)
}
pub fn memory_list(&self, workspace: &str) -> Result<Vec<MemoryEntry>> {
let mut stmt = self.conn.prepare(
"SELECT key, value, run_id, step, created_at FROM memory
WHERE workspace = ?1 ORDER BY created_at ASC, id ASC",
)?;
let rows = stmt.query_map([workspace], |r| {
Ok(MemoryEntry {
key: r.get(0)?,
value: r.get(1)?,
run_id: r.get(2)?,
step: r.get::<_, i64>(3)? as u32,
created_at: r.get(4)?,
})
})?;
Ok(rows.collect::<std::result::Result<_, _>>()?)
}
pub fn memory_get(&self, workspace: &str, key: &str) -> Result<Option<MemoryEntry>> {
Ok(self
.conn
.query_row(
"SELECT key, value, run_id, step, created_at FROM memory
WHERE workspace = ?1 AND key = ?2",
(workspace, key),
|r| {
Ok(MemoryEntry {
key: r.get(0)?,
value: r.get(1)?,
run_id: r.get(2)?,
step: r.get::<_, i64>(3)? as u32,
created_at: r.get(4)?,
})
},
)
.ok())
}
pub fn memory_delete(&self, workspace: &str, key: &str) -> Result<bool> {
let n = self.conn.execute(
"DELETE FROM memory WHERE workspace = ?1 AND key = ?2",
(workspace, key),
)?;
Ok(n > 0)
}
pub fn memory_clear(&self, workspace: &str) -> Result<usize> {
Ok(self
.conn
.execute("DELETE FROM memory WHERE workspace = ?1", [workspace])?)
}
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn refusals_record_action_target_rule_and_layer() {
let store = Store::memory().unwrap();
let run = store.start_run("goal", "root").unwrap();
store
.record_event(
run,
&PolicyEvent::refusal(2, "write", "secrets/key.txt").with_rule("secrets/*", "base"),
)
.unwrap();
let events = store.events(run).unwrap();
assert_eq!(events.len(), 1);
let e = &events[0];
assert_eq!(e.kind, "refusal");
assert_eq!(e.act, "write");
assert_eq!(e.target, "secrets/key.txt");
assert_eq!(e.rule.as_deref(), Some("secrets/*"));
assert_eq!(e.layer.as_deref(), Some("base"));
}
#[test]
fn decisions_record_their_value_source_and_any_altered_target() {
let store = Store::memory().unwrap();
let run = store.start_run("goal", "root").unwrap();
store
.record_event(
run,
&PolicyEvent::decision(1, "write", "src/a.rs", "approve", "stdin")
.with_performed("src/sandbox/a.rs"),
)
.unwrap();
store
.record_event(
run,
&PolicyEvent::decision(2, "write", "src/b.rs", "approve", "remembered"),
)
.unwrap();
let events = store.events(run).unwrap();
assert_eq!(events.len(), 2);
assert_eq!(events[0].decision.as_deref(), Some("approve"));
assert_eq!(events[0].target, "src/a.rs");
assert_eq!(events[0].performed.as_deref(), Some("src/sandbox/a.rs"));
assert_eq!(events[1].source.as_deref(), Some("remembered"));
assert_eq!(events[1].performed, None);
}
#[test]
fn a_pre_0_4_database_migrates_in_place_and_keeps_its_rows() {
let dir = tempfile::tempdir().unwrap();
let path = dir.path().join("runs.db");
{
let conn = rusqlite::Connection::open(&path).unwrap();
conn.execute_batch(
"CREATE TABLE runs (id INTEGER PRIMARY KEY AUTOINCREMENT, goal TEXT NOT NULL,
file TEXT NOT NULL, outcome TEXT, provider TEXT);
CREATE TABLE steps (id INTEGER PRIMARY KEY AUTOINCREMENT, run_id INTEGER NOT NULL,
step INTEGER NOT NULL, decision TEXT NOT NULL, result TEXT NOT NULL,
prompt TEXT NOT NULL DEFAULT '', tool_call TEXT NOT NULL DEFAULT '',
tokens INTEGER NOT NULL DEFAULT 0);
INSERT INTO runs (goal, file) VALUES ('old goal', 'old.txt');",
)
.unwrap();
}
let store = Store::open(&path).unwrap();
assert_eq!(store.last_step(1).unwrap(), 0);
store
.record_event(1, &PolicyEvent::refusal(1, "read", ".env"))
.unwrap();
assert_eq!(store.events(1).unwrap().len(), 1);
}
#[test]
fn a_pending_approval_survives_the_store_being_reopened() {
let dir = tempfile::tempdir().unwrap();
let path = dir.path().join("runs.db");
let request_id = {
let store = Store::open(&path).unwrap();
let run = store.start_run("goal", "root").unwrap();
store
.put_pending(run, 3, "write", "src/a.rs", Some("fn a() {}"))
.unwrap()
};
let store = Store::open(&path).unwrap();
let p = store.pending(request_id).unwrap().expect("still pending");
assert_eq!(p.step, 3);
assert_eq!(p.act, "write");
assert_eq!(p.target, "src/a.rs");
assert_eq!(p.content.as_deref(), Some("fn a() {}"));
assert_eq!(p.resolved, None);
store.resolve_pending(request_id, "approve").unwrap();
let p = store.pending(request_id).unwrap().unwrap();
assert_eq!(p.resolved.as_deref(), Some("approve"));
}
#[test]
fn the_tree_is_reconstructable_from_a_reopened_store() {
let dir = tempfile::tempdir().unwrap();
let path = dir.path().join("runs.db");
let (root, c1, c2, gc) = {
let store = Store::open(&path).unwrap();
let root = store.start_run("root goal", "ws").unwrap();
let c1 = store.start_child_run("child 1", "ws", root, 1).unwrap();
let c2 = store.start_child_run("child 2", "ws", root, 1).unwrap();
let gc = store.start_child_run("grandchild", "ws", c1, 2).unwrap();
store
.record_agent_event(&AgentEvent::spawn(root, 1, c1, "child 1"))
.unwrap();
store
.record_agent_event(&AgentEvent::spawn(root, 1, c2, "child 2"))
.unwrap();
store
.record_agent_event(&AgentEvent::spawn(c1, 1, gc, "grandchild"))
.unwrap();
store
.record_agent_event(&AgentEvent::spawn_refused(root, 2, "agents"))
.unwrap();
store
.record_agent_event(&AgentEvent::budget_draw(c1, 1, 30, 70))
.unwrap();
(root, c1, c2, gc)
};
let store = Store::open(&path).unwrap();
assert_eq!(store.children(root).unwrap(), vec![c1, c2]);
assert_eq!(store.children(c1).unwrap(), vec![gc]);
assert_eq!(store.parent(gc).unwrap(), Some(c1));
assert_eq!(store.parent(root).unwrap(), None);
assert_eq!(store.depth(gc).unwrap(), 2);
let root_events = store.agent_events(root).unwrap();
assert_eq!(root_events.iter().filter(|e| e.kind == "spawn").count(), 2);
assert_eq!(
root_events
.iter()
.filter(|e| e.kind == "spawn_refused")
.count(),
1
);
let draws = store.agent_events(c1).unwrap();
let draw = draws.iter().find(|e| e.kind == "budget_draw").unwrap();
assert_eq!(draw.tokens, Some(30));
assert_eq!(draw.remaining, Some(70));
}
#[test]
fn a_pre_0_5_database_migrates_and_keeps_its_rows() {
let dir = tempfile::tempdir().unwrap();
let path = dir.path().join("runs.db");
{
let conn = rusqlite::Connection::open(&path).unwrap();
conn.execute_batch(
"CREATE TABLE runs (id INTEGER PRIMARY KEY AUTOINCREMENT, goal TEXT NOT NULL,
file TEXT NOT NULL, outcome TEXT, provider TEXT);
CREATE TABLE steps (id INTEGER PRIMARY KEY AUTOINCREMENT, run_id INTEGER NOT NULL,
step INTEGER NOT NULL, decision TEXT NOT NULL, result TEXT NOT NULL,
prompt TEXT NOT NULL DEFAULT '', tool_call TEXT NOT NULL DEFAULT '',
tokens INTEGER NOT NULL DEFAULT 0);
INSERT INTO runs (goal, file) VALUES ('old', 'old.txt');",
)
.unwrap();
}
let store = Store::open(&path).unwrap();
assert_eq!(store.parent(1).unwrap(), None);
assert_eq!(store.depth(1).unwrap(), 0);
let child = store.start_child_run("c", "ws", 1, 1).unwrap();
assert_eq!(store.children(1).unwrap(), vec![child]);
}
#[test]
fn a_pre_0_8_database_migrates_in_place_and_keeps_its_rows() {
let dir = tempfile::tempdir().unwrap();
let path = dir.path().join("runs.db");
{
let store = Store::open(&path).unwrap();
let run = store.start_run("old goal", "old.txt").unwrap();
store
.checkpoint_step(run, &StepRecord::new(1, "wrote", "ok"))
.unwrap();
store
.record_event(run, &PolicyEvent::refusal(1, "write", "secrets/k"))
.unwrap();
store
.conn
.execute("DROP TABLE IF EXISTS mcp_events", [])
.unwrap();
}
let store = Store::open(&path).unwrap();
assert_eq!(store.last_step(1).unwrap(), 1);
assert_eq!(store.events(1).unwrap().len(), 1);
assert!(store.mcp_events(1).unwrap().is_empty());
store
.record_mcp(1, &McpEvent::connected("files", "stdio"))
.unwrap();
let events = store.mcp_events(1).unwrap();
assert_eq!(events.len(), 1);
assert_eq!(events[0].detail.as_deref(), Some("stdio"));
assert_eq!(store.steps(1).unwrap().len(), 1);
assert_eq!(store.run_status(1).unwrap(), Some(RunStatus::Running));
}
#[test]
fn full_trace_persists_and_reads_back() {
let store = Store::memory().unwrap();
let run = store.start_run("goal", "out.txt").unwrap();
store
.record(
run,
&StepRecord::new(1, "wrote file", "content v1").with_trace(
"the prompt",
r#"{"content":"content v1"}"#,
128,
),
)
.unwrap();
store
.record(run, &StepRecord::new(2, "verified", "ok"))
.unwrap();
store.finish_run(run, "success").unwrap();
let steps = store.steps(run).unwrap();
assert_eq!(steps.len(), 2);
assert_eq!(steps[0].decision, "wrote file");
assert_eq!(steps[0].prompt, "the prompt");
assert_eq!(steps[0].tokens, 128);
assert_eq!(steps[1].result, "ok");
assert_eq!(store.last_step(run).unwrap(), 2);
}
#[test]
fn migrates_a_0_1_0_steps_table_in_place() {
let conn = Connection::open_in_memory().unwrap();
conn.execute_batch(
"CREATE TABLE runs (id INTEGER PRIMARY KEY AUTOINCREMENT, goal TEXT NOT NULL, file TEXT NOT NULL, outcome TEXT);
CREATE TABLE steps (id INTEGER PRIMARY KEY AUTOINCREMENT, run_id INTEGER NOT NULL, step INTEGER NOT NULL, decision TEXT NOT NULL, result TEXT NOT NULL);
INSERT INTO runs (goal, file) VALUES ('g', 'f');
INSERT INTO steps (run_id, step, decision, result) VALUES (1, 1, 'wrote file', 'old');",
)
.unwrap();
let store = Store::from_conn(conn).unwrap();
let steps = store.steps(1).unwrap();
assert_eq!(steps.len(), 1);
assert_eq!(steps[0].result, "old");
assert_eq!(steps[0].prompt, "");
assert_eq!(steps[0].tokens, 0);
}
#[test]
fn provider_is_recorded_and_read_back() {
let store = Store::memory().unwrap();
let run = store.start_run("g", "f").unwrap();
assert_eq!(store.provider(run).unwrap(), None);
store.set_provider(run, "anthropic").unwrap();
assert_eq!(store.provider(run).unwrap().as_deref(), Some("anthropic"));
}
#[test]
fn migrates_a_pre_0_3_runs_table_adding_provider() {
let conn = Connection::open_in_memory().unwrap();
conn.execute_batch(
"CREATE TABLE runs (id INTEGER PRIMARY KEY AUTOINCREMENT, goal TEXT NOT NULL, file TEXT NOT NULL, outcome TEXT);
CREATE TABLE steps (id INTEGER PRIMARY KEY AUTOINCREMENT, run_id INTEGER NOT NULL, step INTEGER NOT NULL, decision TEXT NOT NULL, result TEXT NOT NULL);
INSERT INTO runs (goal, file) VALUES ('g', 'f');",
)
.unwrap();
let store = Store::from_conn(conn).unwrap();
assert_eq!(store.provider(1).unwrap(), None);
store.set_provider(1, "openai").unwrap();
assert_eq!(store.provider(1).unwrap().as_deref(), Some("openai"));
}
#[test]
fn checkpoint_step_commits_the_step_and_its_event_together() {
let store = Store::memory().unwrap();
let run = store.start_run("goal", "root").unwrap();
store
.checkpoint_step(run, &StepRecord::new(1, "act", "ok"))
.unwrap();
store
.checkpoint_step(run, &StepRecord::new(2, "act", "ok"))
.unwrap();
assert_eq!(store.last_step(run).unwrap(), 2);
assert_eq!(store.steps(run).unwrap().len(), 2);
let cps: Vec<_> = store
.checkpoint_events(run)
.unwrap()
.into_iter()
.filter(|e| e.kind == "checkpoint")
.collect();
assert_eq!(cps.len(), 2);
assert!(cps.iter().all(|e| e.detail.is_none()));
}
#[test]
fn a_rolled_back_step_leaves_the_prior_checkpoint_intact() {
let store = Store::memory().unwrap();
let run = store.start_run("goal", "root").unwrap();
store
.checkpoint_step(run, &StepRecord::new(1, "act", "ok"))
.unwrap();
{
let tx = store.conn.unchecked_transaction().unwrap();
tx.execute(
"INSERT INTO steps (run_id, step, decision, result) VALUES (?1, 2, 'act', 'ok')",
[run],
)
.unwrap();
tx.execute(
"INSERT INTO checkpoint_events (run_id, step, kind) VALUES (?1, 2, 'checkpoint')",
[run],
)
.unwrap();
}
assert_eq!(
store.last_step(run).unwrap(),
1,
"the torn step must not survive"
);
assert_eq!(store.steps(run).unwrap().len(), 1);
}
#[test]
fn check_resumable_refuses_a_newer_format_and_a_missing_run() {
let store = Store::memory().unwrap();
let run = store.start_run("goal", "root").unwrap();
assert!(store.check_resumable(run).is_ok());
assert!(matches!(
store.check_resumable(9999),
Err(Error::Resume { .. })
));
store
.conn
.execute_batch(&format!("PRAGMA user_version = {}", CHECKPOINT_FORMAT + 1))
.unwrap();
assert!(matches!(
store.check_resumable(run),
Err(Error::Resume { .. })
));
}
#[test]
fn spent_tokens_and_elapsed_are_durable_reads() {
let store = Store::memory().unwrap();
let run = store.start_run("goal", "root").unwrap();
store
.checkpoint_step(run, &StepRecord::new(1, "a", "ok").with_trace("p", "t", 30))
.unwrap();
store
.checkpoint_step(run, &StepRecord::new(2, "a", "ok").with_trace("p", "t", 12))
.unwrap();
assert_eq!(store.spent_tokens(run).unwrap(), 42);
assert!(store.elapsed_secs(run).unwrap() >= 0.0);
}
#[test]
fn tree_aggregate_reads_span_root_and_descendants() {
let store = Store::memory().unwrap();
let root = store.start_run("goal", "root").unwrap();
let child = store.start_child_run("sub", "root", root, 1).unwrap();
let grandchild = store.start_child_run("subsub", "root", child, 2).unwrap();
store
.checkpoint_step(
root,
&StepRecord::new(1, "a", "ok").with_trace("p", "t", 10),
)
.unwrap();
store
.checkpoint_step(
child,
&StepRecord::new(1, "a", "ok").with_trace("p", "t", 20),
)
.unwrap();
store
.checkpoint_step(
grandchild,
&StepRecord::new(1, "a", "ok").with_trace("p", "t", 5),
)
.unwrap();
assert_eq!(
store.tree_run_ids(root).unwrap(),
vec![root, child, grandchild]
);
assert_eq!(store.spent_tokens_tree(root).unwrap(), 35);
assert_eq!(store.agent_count_tree(root).unwrap(), 3);
}
#[test]
fn status_round_trips_and_a_pre_0_7_database_migrates() {
let conn = Connection::open_in_memory().unwrap();
conn.execute_batch(
"CREATE TABLE runs (id INTEGER PRIMARY KEY AUTOINCREMENT, goal TEXT NOT NULL, file TEXT NOT NULL, outcome TEXT, provider TEXT, parent_run_id INTEGER, depth INTEGER NOT NULL DEFAULT 0);
INSERT INTO runs (goal, file) VALUES ('g', 'f');",
)
.unwrap();
let store = Store::from_conn(conn).unwrap();
assert_eq!(store.status(1).unwrap().as_deref(), Some("running"));
store.set_status(1, "completed").unwrap();
assert_eq!(store.status(1).unwrap().as_deref(), Some("completed"));
}
#[test]
fn the_entry_count_cap_evicts_oldest_first_and_never_the_new_entry() {
let store = Store::memory().unwrap();
for i in 0..MEMORY_MAX_ENTRIES {
let evicted = store.memory_put("ws", &format!("k{i}"), "v", 1, 1).unwrap();
assert!(evicted.is_empty(), "no eviction while under the cap");
}
assert_eq!(store.memory_list("ws").unwrap().len(), MEMORY_MAX_ENTRIES);
let mut evicted = Vec::new();
for i in 0..3 {
evicted.extend(
store
.memory_put("ws", &format!("new{i}"), "v", 2, 2)
.unwrap(),
);
}
assert_eq!(evicted, vec!["k0", "k1", "k2"]);
let keys: Vec<String> = store
.memory_list("ws")
.unwrap()
.into_iter()
.map(|e| e.key)
.collect();
assert_eq!(
keys.len(),
MEMORY_MAX_ENTRIES,
"the cap holds after eviction"
);
assert!(!keys.contains(&"k0".to_string()));
for i in 0..3 {
assert!(keys.contains(&format!("new{i}")));
}
}
#[test]
fn the_total_chars_cap_evicts_before_the_count_cap_is_reached() {
let store = Store::memory().unwrap();
let big = "x".repeat(MEMORY_MAX_ENTRY_CHARS);
let mut evicted = Vec::new();
for i in 0..10 {
evicted.extend(
store
.memory_put("ws", &format!("k{i}"), &big, 1, 1)
.unwrap(),
);
}
assert_eq!(
evicted,
vec!["k0", "k1"],
"oldest first, count cap untouched"
);
let entries = store.memory_list("ws").unwrap();
assert!(entries.len() < MEMORY_MAX_ENTRIES);
let total: usize = entries.iter().map(|e| e.value.chars().count()).sum();
assert!(total <= MEMORY_MAX_CHARS, "{total} chars is over the cap");
}
#[test]
fn an_oversized_value_is_truncated_with_a_marker_not_rejected() {
let store = Store::memory().unwrap();
let huge = "é".repeat(MEMORY_MAX_ENTRY_CHARS * 2);
assert!(store.memory_put("ws", "k", &huge, 1, 1).is_ok());
let stored = store.memory_get("ws", "k").unwrap().unwrap().value;
assert_eq!(stored.chars().count(), MEMORY_MAX_ENTRY_CHARS);
assert!(stored.ends_with(MEMORY_TRUNCATED), "the cut is visible");
let kept = MEMORY_MAX_ENTRY_CHARS - MEMORY_TRUNCATED.chars().count();
assert!(stored.chars().take(kept).all(|c| c == 'é'));
}
#[test]
fn a_0_9_1_store_opens_unchanged_and_still_resumes() {
let dir = tempfile::tempdir().unwrap();
let path = dir.path().join("runs.db");
let before_format: i64 = {
let store = Store::open(&path).unwrap();
let run = store.start_run("old goal", "old.txt").unwrap();
store
.checkpoint_step(run, &StepRecord::new(1, "wrote", "ok"))
.unwrap();
store
.record_event(run, &PolicyEvent::refusal(1, "write", "secrets/k"))
.unwrap();
store
.put_pending(run, 1, "write", "src/a.rs", None)
.unwrap();
let child = store.start_child_run("sub", "ws", run, 1).unwrap();
store
.record_agent_event(&AgentEvent::spawn(run, 1, child, "sub"))
.unwrap();
store
.record_sandbox_event(&SandboxEvent::create(run, 1, "proc"))
.unwrap();
store
.record_spawn(run, 1, child, "sub", "out.txt", "ok", None, "[]")
.unwrap();
store
.record_mcp(run, &McpEvent::connected("files", "stdio"))
.unwrap();
store.conn.execute("DROP TABLE memory", []).unwrap();
store
.conn
.query_row("PRAGMA user_version", [], |r| r.get(0))
.unwrap()
};
let store = Store::open(&path).unwrap();
let after_format: i64 = store
.conn
.query_row("PRAGMA user_version", [], |r| r.get(0))
.unwrap();
assert_eq!(
after_format, before_format,
"the checkpoint format must not move — a 0.9.1 checkpoint still resumes"
);
assert_eq!(after_format, CHECKPOINT_FORMAT);
assert!(store.check_resumable(1).is_ok());
assert_eq!(store.steps(1).unwrap().len(), 1);
assert_eq!(store.last_step(1).unwrap(), 1);
assert_eq!(store.events(1).unwrap().len(), 1);
assert_eq!(store.pending(1).unwrap().unwrap().act, "write");
assert_eq!(store.checkpoint_events(1).unwrap().len(), 1);
assert_eq!(store.agent_events(1).unwrap().len(), 1);
assert_eq!(store.sandbox_events(1).unwrap().len(), 1);
assert_eq!(store.mcp_events(1).unwrap().len(), 1);
assert_eq!(store.children(1).unwrap(), vec![2]);
assert!(store.find_spawn(1, 1, "sub").is_ok());
assert_eq!(store.run_status(1).unwrap(), Some(RunStatus::Running));
assert!(store.memory_list("ws").unwrap().is_empty());
store.memory_put("ws", "k", "v", 1, 1).unwrap();
assert_eq!(store.memory_get("ws", "k").unwrap().unwrap().value, "v");
}
}