Skip to main content

agent_top_core/harness/
kodelet.rs

1//! Kodelet: cumulative conversation snapshots in `$KODELET_BASE_PATH/storage.db`
2//! (default `~/.kodelet/storage.db`), opened read-only.
3//!
4//! Format verified on 2026-09-19 against Kodelet v0.6.17-beta (`6dcef0a4`).
5//! Source: `pkg/conversations/sqlite`, `pkg/types/llm/usage.go`, and provider
6//! persistence code in `pkg/llm`.
7//!
8//! Non-obvious storage contracts:
9//! * JSON columns may be BLOBs despite TEXT affinity. Usage and USD costs are
10//!   cumulative snapshots; `conversation_summaries` duplicates the accounting.
11//! * Responses input includes cache reads; Anthropic/Chat input excludes them.
12//!   Cache writes have no TTL split, so the 5m slot holds the unsplit aggregate.
13//! * Forks copy messages/results but reset usage. Only explicit parent metadata
14//!   establishes hierarchy; copied results are not new child work.
15//! * Compaction discards calls/results but retains usage, with no lifetime tool
16//!   counter. Retained/observed calls therefore provide only a lower bound.
17//! * Tool-result timestamps mark completion; `metadata.executionTime` is in
18//!   nanoseconds. Per-response usage and message timestamps are not persisted.
19//! * An MCP call is an ordinary tool named `mcp__<server>_<tool>`
20//!   (`sdk/src/extensions/mcp/register.ts`, `extensionToolName`, read on
21//!   2026-09-20). Taken from that registration rather than a live session, as
22//!   no Kodelet install was to hand; the split is unit-tested, not golden.
23
24mod subagent;
25
26use super::{AttributeContext, HarnessAdapter, RegistryHints, SessionSummary, SessionTracker, SpanRetention, parse_rfc3339_utc};
27use crate::model::{Activity, Attribution, CostBreakdown, Harness, PriceSource, ProcNode, SubagentInfo, TokenUsage};
28use crate::process::{RawProc, kodelet_command};
29use rusqlite::{Connection, OpenFlags, OptionalExtension};
30use serde_json::Value;
31use std::cell::RefCell;
32use std::collections::{HashMap, HashSet};
33use std::path::{Path, PathBuf};
34use std::rc::Rc;
35use std::time::{Duration, SystemTime, UNIX_EPOCH};
36use subagent::SubagentIndex;
37
38type SharedSubagents = Rc<RefCell<SubagentIndex>>;
39
40pub fn kodelet_dir() -> Option<PathBuf> {
41    std::env::var_os("KODELET_BASE_PATH")
42        .filter(|v| !v.is_empty())
43        .map(PathBuf::from)
44        .or_else(|| std::env::var_os("HOME").map(|h| PathBuf::from(h).join(".kodelet")))
45}
46
47pub fn db_path() -> Option<PathBuf> {
48    let path = kodelet_dir()?.join("storage.db");
49    path.is_file().then_some(path)
50}
51
52pub fn session_path(db: &Path, id: &str) -> PathBuf {
53    db.join(id)
54}
55
56pub fn session_id_of(path: &Path) -> Option<String> {
57    path.file_name()?.to_str().map(str::to_owned)
58}
59
60fn open_ro(path: &Path) -> rusqlite::Result<Connection> {
61    let conn = Connection::open_with_flags(path, OpenFlags::SQLITE_OPEN_READ_ONLY)?;
62    conn.busy_timeout(Duration::ZERO)?;
63    Ok(conn)
64}
65
66fn valid_id(id: &str) -> bool {
67    !id.is_empty() && id.trim() == id && id != "." && id != ".." && !id.contains(['/', '\\'])
68}
69
70fn table_exists(conn: &Connection, name: &str) -> bool {
71    conn.query_row("SELECT EXISTS(SELECT 1 FROM sqlite_master WHERE type = 'table' AND name = ?1)", [name], |r| r.get(0)).unwrap_or(false)
72}
73
74/// JSON uses RFC3339Nano; modernc SQLite also writes Go time strings with a
75/// space separator, numeric offset and optional zone name. CURRENT_TIMESTAMP
76/// (older event tables) is UTC without a zone. Preserve fractional seconds.
77fn timestamp(text: &str) -> Option<SystemTime> {
78    if text.as_bytes().get(10) == Some(&b'T') {
79        return parse_rfc3339_utc(text);
80    }
81    let mut parts = text.split_whitespace();
82    let date = parts.next()?;
83    let clock = parts.next()?;
84    let zone = parts.next().unwrap_or("Z");
85    let zone = if zone == "UTC" { "Z" } else { zone };
86    let combined = format!("{date}T{clock}");
87    parse_rfc3339_utc(&combined).or_else(|| parse_rfc3339_utc(&format!("{combined}{zone}")))
88}
89
90fn string(value: &Value, key: &str) -> Option<String> {
91    value.get(key)?.as_str().map(str::trim).filter(|s| !s.is_empty()).map(str::to_owned)
92}
93
94fn read_json(path: &Path) -> Option<Value> {
95    serde_json::from_slice(&std::fs::read(path).ok()?).ok()
96}
97
98#[derive(Debug)]
99struct Session {
100    id: String,
101    cwd: Option<PathBuf>,
102    created: Option<SystemTime>,
103    updated: Option<SystemTime>,
104    child: bool,
105}
106
107fn recent_sessions(db: &Path, since: SystemTime) -> Vec<Session> {
108    let Ok(conn) = open_ro(db) else { return Vec::new() };
109    // Avoid both timestamp string ordering (different offsets) and SQLite's
110    // date parser (which cannot read all of Go's time.String representations).
111    let Ok(mut stmt) =
112        conn.prepare("SELECT id, cwd, CAST(created_at AS TEXT), CAST(updated_at AS TEXT), CAST(metadata AS TEXT) FROM conversations")
113    else {
114        return Vec::new();
115    };
116    let rows = stmt.query_map([], |r| {
117        let metadata = r.get::<_, Option<String>>(4)?.and_then(|s| serde_json::from_str::<Value>(&s).ok()).unwrap_or_default();
118        Ok(Session {
119            id: r.get(0)?,
120            cwd: r.get::<_, Option<String>>(1)?.filter(|s| !s.is_empty()).map(PathBuf::from),
121            created: timestamp(&r.get::<_, String>(2)?),
122            updated: timestamp(&r.get::<_, String>(3)?),
123            child: string(&metadata, "parent_conversation_id").is_some(),
124        })
125    });
126    let mut sessions: Vec<_> = rows
127        .map(|rows| rows.flatten().filter(|s| s.updated.is_some_and(|at| at >= since) || since == UNIX_EPOCH).collect())
128        .unwrap_or_default();
129    // A virtual path must stay inside the database, even for a malformed ID.
130    sessions.retain(|s| valid_id(&s.id));
131    sessions.sort_by(|a, b| b.updated.cmp(&a.updated).then_with(|| a.id.cmp(&b.id)));
132    sessions
133}
134
135/// Each conversation is a separate row, including children. Nothing is folded
136/// into a parent, so snapshots and exports can sum their accounting once.
137#[derive(Default)]
138pub struct KodeletAdapter {
139    db: Option<PathBuf>,
140    recent: Vec<Session>,
141    exact: HashMap<u32, (Vec<PathBuf>, Attribution)>,
142    reserved: HashSet<PathBuf>,
143    hints: HashMap<u32, RegistryHints>,
144    modern: bool,
145    subagents: RefCell<HashMap<PathBuf, SharedSubagents>>,
146}
147
148impl KodeletAdapter {
149    fn db(&self) -> Option<PathBuf> {
150        self.db.clone().or_else(db_path)
151    }
152
153    fn reserve(&mut self, db: &Path, pid: u32, id: &str, attribution: Attribution) {
154        if !valid_id(id) {
155            return;
156        }
157        let path = session_path(db, id);
158        if self.reserved.insert(path.clone()) {
159            self.exact.entry(pid).or_insert_with(|| (Vec::new(), attribution)).0.push(path);
160        }
161    }
162
163    fn prepare_registry(&mut self, db: &Path, roots: &[&ProcNode], by_pid: &HashMap<u32, &RawProc>) {
164        let Ok(conn) = open_ro(db) else { return };
165        self.modern = table_exists(&conn, "runner_registrations") || table_exists(&conn, "chat_turns");
166        let Some(base) = db.parent() else { return };
167        let host = read_json(&base.join("runners/host.json")).and_then(|v| string(&v, "instanceId"));
168        let now = SystemTime::now();
169        let root_pids: HashSet<_> = roots.iter().map(|r| r.pid).collect();
170
171        // Never use host_pid alone: a remote runner may have exactly the same
172        // PID as a local client. A stale registration also cannot own a reused
173        // PID: its connection must postdate this process's start.
174        if let Some(host) = host {
175            let sql = "SELECT DISTINCT rr.conversation_id, r.host_pid, CAST(r.connected_at AS TEXT), \
176                       CAST(r.last_heartbeat_at AS TEXT), r.kodelet_version FROM runner_runs rr \
177                       JOIN runner_registrations r ON r.id = rr.runner_id JOIN conversations c ON c.id = rr.conversation_id \
178                       WHERE rr.status IN ('opening', 'running') AND r.status IN ('idle', 'busy') AND r.host_instance_id = ?1";
179            if let Ok(mut stmt) = conn.prepare(sql)
180                && let Ok(rows) = stmt.query_map([host], |r| {
181                    Ok((
182                        r.get::<_, String>(0)?,
183                        r.get::<_, u32>(1)?,
184                        r.get::<_, Option<String>>(2)?,
185                        r.get::<_, Option<String>>(3)?,
186                        r.get::<_, Option<String>>(4)?,
187                    ))
188                })
189            {
190                for (id, pid, connected, heartbeat, version) in rows.flatten() {
191                    let Some(raw) = by_pid.get(&pid).filter(|_| root_pids.contains(&pid)) else { continue };
192                    let start = UNIX_EPOCH + Duration::from_secs(raw.start_time);
193                    if !execution_host(raw)
194                        || connected.as_deref().and_then(timestamp).is_none_or(|at| at < start)
195                        || !heartbeat
196                            .as_deref()
197                            .and_then(timestamp)
198                            .is_some_and(|at| at >= start && now.duration_since(at).is_ok_and(|age| age <= Duration::from_secs(45)))
199                    {
200                        continue;
201                    }
202                    self.reserve(db, pid, &id, Attribution::HarnessRegistry);
203                    let hints = self.hints.entry(pid).or_default();
204                    hints.version = version;
205                    hints.status = Some("busy".into());
206                }
207            }
208        }
209
210        // The local daemon owns model execution even when its workspace runs
211        // remotely. Prefer a proven local runner above; otherwise active turn
212        // receipts can be assigned to this verified daemon, never to a client.
213        let connection_path = base.join("server/connection.json");
214        if let Some(connection) = read_json(&connection_path)
215            && connection.get("schemaVersion").and_then(Value::as_u64) == Some(1)
216            && string(&connection, "instanceId").is_some()
217            && let Some(pid) = connection.get("pid").and_then(Value::as_u64).and_then(|p| u32::try_from(p).ok())
218            && root_pids.contains(&pid)
219            && let Some(raw) = by_pid.get(&pid)
220            && kodelet_command(raw.cmd.get(1..).unwrap_or_default()).first().is_some_and(|a| a == "serve")
221            && std::fs::metadata(&connection_path)
222                .and_then(|m| m.modified())
223                .ok()
224                .is_some_and(|at| at >= UNIX_EPOCH + Duration::from_secs(raw.start_time))
225        {
226            self.hints.entry(pid).or_default().version = string(&connection, "version");
227            if let Ok(mut stmt) = conn.prepare("SELECT DISTINCT t.conversation_id FROM chat_turns t JOIN conversations c ON c.id = t.conversation_id WHERE t.status IN ('accepted', 'running')")
228                && let Ok(rows) = stmt.query_map([], |r| r.get::<_, String>(0))
229            {
230                for id in rows.flatten() {
231                    self.reserve(db, pid, &id, Attribution::HarnessRegistry);
232                    self.hints.entry(pid).or_default().status = Some("busy".into());
233                }
234            }
235        }
236    }
237}
238
239fn execution_host(raw: &RawProc) -> bool {
240    let args = kodelet_command(raw.cmd.get(1..).unwrap_or_default());
241    args.first().is_some_and(|a| a == "serve") || (args.first().is_some_and(|a| a == "runner") && args.get(1).is_some_and(|a| a == "start"))
242}
243
244impl HarnessAdapter for KodeletAdapter {
245    fn harness(&self) -> Harness {
246        Harness::Kodelet
247    }
248
249    fn rescan(&mut self, since: SystemTime) {
250        self.db = db_path();
251        self.recent = self.db.as_deref().map(|db| recent_sessions(db, since)).unwrap_or_default();
252    }
253
254    fn prepare(&mut self, roots: &[&ProcNode], by_pid: &HashMap<u32, &RawProc>) {
255        self.exact.clear();
256        self.reserved.clear();
257        self.hints.clear();
258        self.modern = false;
259        let Some(db) = self.db() else { return };
260        self.prepare_registry(&db, roots, by_pid);
261    }
262
263    fn hints(&self, pid: u32) -> Option<RegistryHints> {
264        self.hints.get(&pid).cloned()
265    }
266
267    fn attribute(&self, root: &ProcNode, raw: Option<&RawProc>, ctx: &AttributeContext) -> (Vec<PathBuf>, Attribution) {
268        if let Some((paths, attribution)) = self.exact.get(&root.pid) {
269            let paths: Vec<_> = paths.iter().filter(|p| !ctx.attached.contains(*p)).cloned().collect();
270            return if paths.is_empty() { (paths, Attribution::None) } else { (paths, *attribution) };
271        }
272        let (Some(cwd), Some(db), Some(raw)) = (ctx.cwd, self.db(), raw) else { return (Vec::new(), Attribution::None) };
273        if self.modern || !execution_host(raw) || raw.cmd.iter().any(|a| a == "--server" || a.starts_with("--server=")) {
274            return (Vec::new(), Attribution::None);
275        }
276        let matches: Vec<_> = self
277            .recent
278            .iter()
279            .filter(|s| !s.child && s.cwd.as_deref() == Some(cwd))
280            .filter(|s| s.created.is_some_and(|at| at + Duration::from_secs(60) >= ctx.proc_start))
281            .filter(|s| s.updated.is_some_and(|at| ctx.now.duration_since(at).is_ok_and(|age| age <= ctx.activity_timeout)))
282            .map(|s| session_path(&db, &s.id))
283            .filter(|p| !self.reserved.contains(p) && !ctx.attached.contains(p))
284            .collect();
285        if matches.len() == 1 { (matches, Attribution::CwdHeuristic) } else { (Vec::new(), Attribution::None) }
286    }
287
288    fn unowned(&self, attached: &HashSet<PathBuf>) -> Vec<PathBuf> {
289        let Some(db) = self.db() else { return Vec::new() };
290        self.recent.iter().map(|s| session_path(&db, &s.id)).filter(|p| !attached.contains(p)).collect()
291    }
292
293    fn open(&self, path: &Path, spans: SpanRetention) -> Box<dyn SessionTracker> {
294        let mut tracker = KodeletTranscript::new(path, spans);
295        if let Some(base) = path.parent().and_then(Path::parent) {
296            tracker.subagents = Some(
297                self.subagents
298                    .borrow_mut()
299                    .entry(base.to_path_buf())
300                    .or_insert_with(|| Rc::new(RefCell::new(SubagentIndex::new(base.to_path_buf()))))
301                    .clone(),
302            );
303        }
304        Box::new(tracker)
305    }
306
307    fn detect(&self, _path: &Path) -> bool {
308        false
309    }
310
311    fn transcripts(&self) -> Vec<(String, PathBuf)> {
312        let Some(db) = self.db() else { return Vec::new() };
313        recent_sessions(&db, UNIX_EPOCH).into_iter().map(|s| (s.id.clone(), session_path(&db, &s.id))).collect()
314    }
315}
316
317pub struct KodeletTranscript {
318    path: PathBuf,
319    retention: SpanRetention,
320    conn: Option<Connection>,
321    data_version: Option<i64>,
322    stamp: Option<SessionStamp>,
323    subagents: Option<SharedSubagents>,
324    sidecar: Option<SubagentInfo>,
325    /// Call ID -> whether it is a native web search. Only metadata is retained,
326    /// so a compaction cannot make an already observed invocation disappear.
327    observed_calls: HashMap<String, bool>,
328    summary: SessionSummary,
329}
330
331#[derive(Debug, Default, PartialEq, Eq)]
332struct SessionStamp {
333    updated: Option<String>,
334    turns: u64,
335    turn_updated: Option<String>,
336    status: Option<String>,
337}
338
339impl SessionStamp {
340    fn read(conn: &Connection, id: &str) -> rusqlite::Result<Self> {
341        let mut stamp = Self {
342            updated: conn.query_row("SELECT CAST(updated_at AS TEXT) FROM conversations WHERE id = ?1", [id], |r| r.get(0)).optional()?,
343            ..Default::default()
344        };
345        if table_exists(conn, "chat_turns") {
346            (stamp.turns, stamp.turn_updated) =
347                conn.query_row("SELECT count(*), CAST(max(updated_at) AS TEXT) FROM chat_turns WHERE conversation_id = ?1", [id], |r| {
348                    Ok((r.get(0)?, r.get(1)?))
349                })?;
350            stamp.status = conn
351                .query_row(
352                    "SELECT status FROM chat_turns WHERE conversation_id = ?1 ORDER BY created_at DESC, turn_id DESC LIMIT 1",
353                    [id],
354                    |r| r.get(0),
355                )
356                .optional()?;
357        }
358        Ok(stamp)
359    }
360}
361
362impl KodeletTranscript {
363    pub fn new(path: &Path, retention: SpanRetention) -> Self {
364        Self {
365            path: path.to_path_buf(),
366            retention,
367            conn: None,
368            data_version: None,
369            stamp: None,
370            subagents: None,
371            sidecar: None,
372            observed_calls: HashMap::new(),
373            summary: SessionSummary { harness: Some(Harness::Kodelet), spans: retention.log(), ..Default::default() },
374        }
375    }
376}
377
378impl SessionTracker for KodeletTranscript {
379    fn refresh(&mut self) -> anyhow::Result<bool> {
380        let Some(db) = self.path.parent() else { return Ok(false) };
381        let Some(id) = session_id_of(&self.path) else { return Ok(false) };
382        let sidecar = self.subagents.as_ref().and_then(|index| {
383            let mut index = index.borrow_mut();
384            index.refresh();
385            index.get(&id).cloned()
386        });
387        if self.conn.is_none() {
388            self.conn = Some(open_ro(db)?);
389        }
390        let conn = self.conn.as_ref().expect("opened above");
391        let version = conn.query_row("PRAGMA data_version", [], |r| r.get::<_, i64>(0))?;
392        if self.data_version == Some(version) && self.sidecar == sidecar {
393            return Ok(false);
394        }
395        // One read snapshot, including usage, tools and turn status. Only mark
396        // it cached after success, so a failed/busy read is retried next tick.
397        let tx = conn.unchecked_transaction()?;
398        let stamp = SessionStamp::read(&tx, &id)?;
399        if self.stamp.as_ref() == Some(&stamp) && self.sidecar == sidecar {
400            // data_version covers the whole database: another conversation or
401            // a runner heartbeat must not make us scan this history again.
402            tx.commit()?;
403            self.data_version = Some(version);
404            return Ok(false);
405        }
406        let (mut summary, calls) = read_summary(&tx, &id, self.retention, &stamp, sidecar.as_ref())?;
407        tx.commit()?;
408        if summary.session_id.is_none() || summary.started_at != self.summary.started_at {
409            self.observed_calls.clear();
410        }
411        self.observed_calls.extend(calls);
412        summary.tool_calls = self.observed_calls.len() as u64;
413        summary.web_searches = self.observed_calls.values().filter(|search| **search).count() as u64;
414        self.summary = summary;
415        self.stamp = Some(stamp);
416        self.data_version = Some(version);
417        self.sidecar = sidecar;
418        Ok(false)
419    }
420
421    fn summary(&self) -> &SessionSummary {
422        &self.summary
423    }
424    fn path(&self) -> &Path {
425        &self.path
426    }
427}
428
429fn read_summary(
430    conn: &Connection,
431    id: &str,
432    retention: SpanRetention,
433    stamp: &SessionStamp,
434    sidecar: Option<&SubagentInfo>,
435) -> anyhow::Result<(SessionSummary, HashMap<String, bool>)> {
436    let mut summary = SessionSummary { harness: Some(Harness::Kodelet), spans: retention.log(), ..Default::default() };
437    let record = conn.query_row(
438        "SELECT cwd, provider, CAST(usage AS TEXT), CAST(metadata AS TEXT), CAST(created_at AS TEXT), CAST(updated_at AS TEXT) FROM conversations WHERE id = ?1", [id],
439        |r| Ok((r.get::<_, Option<String>>(0)?, r.get::<_, String>(1)?, r.get::<_, String>(2)?, r.get::<_, Option<String>>(3)?, r.get::<_, String>(4)?, r.get::<_, String>(5)?)),
440    ).optional()?;
441    let Some((cwd, provider, usage, metadata, created, updated)) = record else { return Ok((summary, HashMap::new())) };
442    let usage: Value = serde_json::from_str(&usage)?;
443    let metadata: Value = serde_json::from_str(metadata.as_deref().unwrap_or("null"))?;
444    summary.session_id = Some(id.to_owned());
445    // Compaction can precede our first read, and not every provider leaves a
446    // structural marker. Never imply a full lifetime count from a snapshot.
447    summary.tool_calls_lower_bound = true;
448    summary.cwd = cwd.filter(|s| !s.is_empty()).map(PathBuf::from);
449    summary.model = string(&metadata, "model").or_else(|| metadata.get("config_snapshot").and_then(|v| string(v, "model")));
450    summary.started_at = timestamp(&created);
451    summary.last_activity = timestamp(&updated);
452    let parent = string(&metadata, "parent_conversation_id");
453    if let Some(parent) = parent.as_ref().filter(|parent| valid_id(parent) && *parent != id) {
454        summary.subagent = Some(SubagentInfo {
455            parent_session_id: parent.clone(),
456            nickname: string(&metadata, "conversation_name"),
457            role: subagent::role(&metadata).map(str::to_owned),
458        });
459    }
460    if let Some(info) = sidecar.filter(|info| valid_id(&info.parent_session_id) && info.parent_session_id != id) {
461        match summary.subagent.as_mut() {
462            Some(central) if central.parent_session_id == info.parent_session_id => {
463                central.nickname = central.nickname.take().or_else(|| info.nickname.clone());
464                central.role = central.role.take().or_else(|| info.role.clone());
465            }
466            None if parent.is_none() => summary.subagent = Some(info.clone()),
467            _ => {}
468        }
469    }
470    let fork = metadata.get("conversation_fork").is_some_and(Value::is_object);
471    let responses = provider == "openai" && responses_mode(conn, id, &metadata)?;
472    accounting(&usage, responses, &mut summary);
473
474    // Durable turn receipts are local to this conversation, unlike the copied
475    // provider history of a fork. They also survive context compaction.
476    summary.turns = stamp.turns;
477    if let Some(status) = stamp.status.as_deref() {
478        summary.activity = if matches!(status, "accepted" | "running") { Activity::Working } else { Activity::Waiting };
479    }
480    if !fork {
481        summary.health.billable_messages = conn.query_row("SELECT count(*) FROM conversations c, json_each(CAST(c.raw_messages AS TEXT)) m WHERE c.id = ?1 AND json_extract(m.value, '$.role') = 'assistant' AND COALESCE(json_extract(m.value, '$.type'), 'message') = 'message'", [id], |r| r.get(0))?;
482        if summary.turns == 0 {
483            summary.turns = summary.health.billable_messages;
484        }
485    }
486    if summary.activity == Activity::Unknown {
487        let tail: Option<String> = conn.query_row("SELECT CASE WHEN json_extract(m.value, '$.role') = 'assistant' AND COALESCE(json_extract(m.value, '$.type'), 'message') = 'message' THEN 'waiting' ELSE 'working' END FROM conversations c, json_each(CAST(c.raw_messages AS TEXT)) m WHERE c.id = ?1 ORDER BY CAST(m.key AS INTEGER) DESC LIMIT 1", [id], |r| r.get(0)).optional()?;
488        // A pristine fork's copied tail says nothing about new work.
489        if !fork || summary.usage.total() > 0 {
490            summary.activity = match tail.as_deref() {
491                Some("waiting") => Activity::Waiting,
492                Some(_) => Activity::Working,
493                None => Activity::Unknown,
494            };
495        }
496    }
497    if usage.is_object() {
498        summary.health.usage_records = summary.health.billable_messages;
499    }
500    if summary.usage.total() == 0 {
501        summary.health.empty_usage_records = summary.health.usage_records;
502    }
503    let calls = read_tools(conn, id, fork, &mut summary)?;
504    Ok((summary, calls))
505}
506
507fn responses_mode(conn: &Connection, id: &str, metadata: &Value) -> rusqlite::Result<bool> {
508    if let Some(mode) =
509        string(metadata, "api_mode").or_else(|| metadata.pointer("/config_snapshot/openai").and_then(|v| string(v, "api_mode")))
510    {
511        return Ok(mode == "responses");
512    }
513    let first: Option<String> =
514        conn.query_row("SELECT json_extract(CAST(raw_messages AS TEXT), '$[0].type') FROM conversations WHERE id = ?1", [id], |r| {
515            r.get(0)
516        })?;
517    Ok(matches!(
518        first.as_deref(),
519        Some("message" | "reasoning" | "function_call" | "function_call_output" | "compaction" | "compaction_summary")
520    ))
521}
522
523fn accounting(usage: &Value, responses: bool, summary: &mut SessionSummary) {
524    let count = |key| usage.get(key).and_then(Value::as_u64).unwrap_or(0);
525    let cached = count("cacheReadInputTokens");
526    summary.usage = TokenUsage {
527        input: if responses { count("inputTokens").saturating_sub(cached) } else { count("inputTokens") },
528        output: count("outputTokens"),
529        cache_read: cached,
530        cache_write_5m: 0,
531        cache_write_1h: 0,
532        // One cumulative counter with no TTL breakdown saved, so it cannot go
533        // in either priced bucket.
534        cache_write_unsplit: count("cacheCreationInputTokens"),
535    };
536    let mut costs = [0.0; 4];
537    for (i, (key, tokens)) in [
538        ("inputCost", summary.usage.input),
539        ("outputCost", summary.usage.output),
540        ("cacheReadCost", cached),
541        ("cacheCreationCost", summary.usage.cache_write()),
542    ]
543    .into_iter()
544    .enumerate()
545    {
546        if let Some(cost) = usage.get(key).and_then(Value::as_f64).filter(|cost| cost.is_finite() && *cost >= 0.0) {
547            costs[i] = cost;
548        } else {
549            summary.unpriced_tokens += tokens;
550        }
551    }
552    summary.cost_breakdown =
553        CostBreakdown { input: costs[0], output: costs[1], cache_read: costs[2], cache_write_unsplit: costs[3], ..Default::default() };
554    summary.cost_usd = summary.cost_breakdown.total();
555    // Kodelet's own figures are the only prices this row can have, whether or
556    // not it recorded one. Leaving the source empty would let the collector
557    // fill it in from the price table, which never priced anything here;
558    // `unpriced_tokens` is what says a figure is missing.
559    summary.price_source = Some(PriceSource::Harness);
560}
561
562/// The server behind a Kodelet MCP tool name.
563///
564/// Kodelet registers an MCP tool as `mcp__<server>_<tool>`: a double
565/// underscore after the prefix, a single one before the tool
566/// (`sdk/src/extensions/mcp/register.ts`, `extensionToolName`). So the
567/// boundary between server and tool is the first underscore after the prefix,
568/// and a server whose own name contains one is split the same way Kodelet
569/// splits it, which keeps the grouping consistent with what the harness shows.
570pub fn mcp_server_of(tool_name: &str) -> Option<&str> {
571    let rest = tool_name.strip_prefix("mcp__")?;
572    let server = rest.split_once('_').map(|(a, _)| a).unwrap_or(rest);
573    if server.is_empty() { None } else { Some(server) }
574}
575
576fn read_tools(conn: &Connection, id: &str, fork: bool, summary: &mut SessionSummary) -> anyhow::Result<HashMap<String, bool>> {
577    let mut seen = HashMap::new();
578    let mut spans = Vec::new();
579    // (name, failed, completed at), by call id, so the fold below counts a
580    // call once however many places recorded it.
581    let mut calls: HashMap<String, (String, bool, Option<SystemTime>)> = HashMap::new();
582    let mut stmt = conn.prepare("SELECT t.key, json_extract(t.value, '$.toolName'), json_extract(t.value, '$.timestamp'), json_extract(t.value, '$.success'), json_extract(t.value, '$.metadata.executionTime') FROM conversations c, json_each(CAST(c.tool_results AS TEXT)) t WHERE c.id = ?1 AND t.type = 'object'")?;
583    let rows = stmt.query_map([id], |r| {
584        Ok((
585            r.get::<_, String>(0)?,
586            r.get::<_, Option<String>>(1)?,
587            r.get::<_, Option<String>>(2)?,
588            r.get::<_, Option<bool>>(3)?,
589            r.get::<_, Option<i64>>(4)?,
590        ))
591    })?;
592    for (call_id, name, at, success, nanos) in rows.flatten() {
593        let ended = at.as_deref().and_then(timestamp);
594        if fork && ended.zip(summary.started_at).is_none_or(|(end, created)| end < created) {
595            continue;
596        }
597        let Some(name) = name.filter(|s| !s.is_empty()) else { continue };
598        if call_id.is_empty() {
599            continue;
600        }
601        seen.insert(call_id.clone(), name == "openai_web_search");
602        calls.insert(call_id.clone(), (name.clone(), success == Some(false), ended));
603        if let Some(end) = ended {
604            let duration = Duration::from_nanos(nanos.unwrap_or(0).max(0) as u64);
605            let start = end.checked_sub(duration).unwrap_or(end);
606            spans.push((start, end, call_id, name, success == Some(false)));
607        }
608    }
609    // Keep timestamps in order before applying the live retention cap. JSON
610    // map order is call-id order, not completion or invocation order.
611    spans.sort_by(|a, b| a.0.cmp(&b.0).then_with(|| a.2.cmp(&b.2)));
612    for (start, end, call_id, name, error) in spans {
613        summary.spans.open(call_id.clone(), name, start, summary.subagent.is_some());
614        summary.spans.close(&call_id, end, error);
615    }
616    if !fork {
617        // Count older/missing structured results too, deduplicated by call ID.
618        // SQL projects only structural fields, never content/input/output.
619        let sql = "WITH messages AS (SELECT m.value FROM conversations c, json_each(CAST(c.raw_messages AS TEXT)) m WHERE c.id = ?1) \
620            SELECT json_extract(b.value, '$.id'), json_extract(b.value, '$.name') FROM messages m, \
621            json_each(CASE WHEN json_type(m.value, '$.content') = 'array' THEN json_extract(m.value, '$.content') ELSE '[]' END) b \
622            WHERE json_extract(b.value, '$.type') = 'tool_use' \
623            UNION SELECT json_extract(t.value, '$.id'), json_extract(t.value, '$.function.name') FROM messages m, json_each(json_extract(m.value, '$.tool_calls')) t \
624            UNION SELECT json_extract(value, '$.call_id'), CASE WHEN json_extract(value, '$.type') = 'web_search_call' THEN 'openai_web_search' ELSE json_extract(value, '$.name') END FROM messages WHERE json_extract(value, '$.type') IN ('function_call', 'web_search_call')";
625        let mut stmt = conn.prepare(sql)?;
626        let rows = stmt.query_map([id], |r| Ok((r.get::<_, Option<String>>(0)?, r.get::<_, Option<String>>(1)?)))?;
627        for (call_id, name) in rows.flatten() {
628            if let Some(call_id) = call_id.filter(|s| !s.is_empty()) {
629                seen.entry(call_id.clone()).or_insert(name.as_deref() == Some("openai_web_search"));
630                // No result was recorded for these, so nothing is known about
631                // whether they failed or when they finished.
632                if let Some(name) = name.filter(|s| !s.is_empty()) {
633                    calls.entry(call_id).or_insert((name, false, None));
634                }
635            }
636        }
637    }
638    for (name, failed, ended) in calls.into_values() {
639        let Some(server) = mcp_server_of(&name) else { continue };
640        let u = summary.mcp.entry(server.to_string()).or_default();
641        u.calls += 1;
642        u.errors += u64::from(failed);
643        u.last_call = u.last_call.max(ended);
644    }
645    Ok(seen)
646}
647
648#[cfg(test)]
649mod tests {
650    use super::*;
651    use crate::model::{ProcKind, SpanKind};
652    use rusqlite::params;
653    use serde_json::json;
654    use std::sync::atomic::{AtomicU64, Ordering};
655
656    const CREATED: &str = "2026-09-19T10:00:00Z";
657    const UPDATED: &str = "2026-09-19T10:05:00Z";
658
659    struct Fixture {
660        dir: PathBuf,
661        conn: Connection,
662    }
663
664    impl Fixture {
665        fn new() -> Self {
666            static NEXT: AtomicU64 = AtomicU64::new(0);
667            let dir =
668                std::env::temp_dir().join(format!("agent-top-kodelet-{}-{}", std::process::id(), NEXT.fetch_add(1, Ordering::Relaxed)));
669            std::fs::create_dir_all(&dir).unwrap();
670            let conn = Connection::open(dir.join("storage.db")).unwrap();
671            conn.execute_batch(
672                "PRAGMA journal_mode=WAL; CREATE TABLE conversations (
673                id TEXT PRIMARY KEY, cwd TEXT, provider TEXT NOT NULL, raw_messages TEXT NOT NULL,
674                usage TEXT NOT NULL, metadata TEXT, tool_results TEXT, created_at DATETIME NOT NULL, updated_at DATETIME NOT NULL);",
675            )
676            .unwrap();
677            Self { dir, conn }
678        }
679
680        fn db(&self) -> PathBuf {
681            self.dir.join("storage.db")
682        }
683
684        fn insert(&self, id: &str, provider: &str, metadata: Value, usage: Value, messages: Value, tools: Value) {
685            // Match the Go driver's []byte bindings, not just JSON TEXT fixtures.
686            self.conn
687                .execute(
688                    "INSERT INTO conversations VALUES (?1, '/work', ?2, ?3, ?4, ?5, ?6, ?7, ?8)",
689                    params![
690                        id,
691                        provider,
692                        serde_json::to_vec(&messages).unwrap(),
693                        serde_json::to_vec(&usage).unwrap(),
694                        serde_json::to_vec(&metadata).unwrap(),
695                        serde_json::to_vec(&tools).unwrap(),
696                        CREATED,
697                        UPDATED,
698                    ],
699                )
700                .unwrap();
701        }
702
703        fn basic(&self, id: &str) {
704            self.insert(id, "anthropic", json!({"model":"custom-model"}), usage(), json!([]), json!({}));
705        }
706
707        fn tracker(&self, id: &str) -> KodeletTranscript {
708            KodeletTranscript::new(&session_path(&self.db(), id), SpanRetention::All)
709        }
710
711        fn adapter(&self) -> KodeletAdapter {
712            KodeletAdapter { db: Some(self.db()), recent: recent_sessions(&self.db(), UNIX_EPOCH), ..Default::default() }
713        }
714
715        fn runtime(&self) {
716            self.conn
717                .execute_batch(
718                    "CREATE TABLE runner_registrations (id TEXT PRIMARY KEY, host_pid INTEGER, host_instance_id TEXT, status TEXT,
719                connected_at DATETIME, last_heartbeat_at DATETIME, kodelet_version TEXT);
720                CREATE TABLE runner_runs (id TEXT PRIMARY KEY, conversation_id TEXT, runner_id TEXT, status TEXT);
721                CREATE TABLE chat_turns (conversation_id TEXT, turn_id TEXT, status TEXT, created_at DATETIME, updated_at DATETIME);",
722                )
723                .unwrap();
724        }
725
726        fn receipt(&self, id: &str, status: &str) {
727            self.conn.execute("INSERT INTO chat_turns VALUES (?1, 'turn-1', ?2, ?3, ?4)", params![id, status, CREATED, UPDATED]).unwrap();
728        }
729
730        fn json_file(&self, path: &str, data: Value) {
731            let path = self.dir.join(path);
732            std::fs::create_dir_all(path.parent().unwrap()).unwrap();
733            std::fs::write(path, serde_json::to_vec(&data).unwrap()).unwrap();
734        }
735    }
736
737    impl Drop for Fixture {
738        fn drop(&mut self) {
739            let _ = std::fs::remove_dir_all(&self.dir);
740        }
741    }
742
743    fn usage() -> Value {
744        json!({"inputTokens":100, "outputTokens":20, "cacheReadInputTokens":60, "cacheCreationInputTokens":10,
745            "inputCost":0.1, "outputCost":0.2, "cacheReadCost":0.03, "cacheCreationCost":0.04,
746            "currentContextWindow":45, "maxContextWindow":200000})
747    }
748
749    fn process(pid: u32, args: &[&str]) -> (RawProc, ProcNode) {
750        let raw = RawProc {
751            pid,
752            ppid: None,
753            name: "kodelet".into(),
754            exe: None,
755            cmd: args.iter().map(|s| (*s).into()).collect(),
756            cwd: Some("/work".into()),
757            cpu_percent: 0.0,
758            rss_bytes: 0,
759            start_time: 1,
760            run_time: 1,
761        };
762        let node = ProcNode {
763            pid,
764            ppid: None,
765            name: raw.name.clone(),
766            cmdline: raw.cmdline(),
767            kind: ProcKind::Agent,
768            harness: Some(Harness::Kodelet),
769            cpu_percent: 0.0,
770            rss_bytes: 0,
771            age_secs: 1,
772            cwd: raw.cwd.clone(),
773            children: vec![],
774        };
775        (raw, node)
776    }
777
778    fn context(attached: &HashSet<PathBuf>) -> AttributeContext<'_> {
779        AttributeContext {
780            cwd: Some(Path::new("/work")),
781            proc_start: timestamp(CREATED).unwrap(),
782            now: timestamp(UPDATED).unwrap(),
783            attached,
784            activity_timeout: Duration::from_secs(900),
785        }
786    }
787
788    #[test]
789    fn accepts_go_sqlite_and_json_timestamps_without_losing_nanoseconds() {
790        let expected = timestamp("2026-09-19T10:00:00.123456789Z").unwrap();
791        for time in [
792            "2026-09-19 10:00:00.123456789 +0000 UTC",
793            "2026-09-19 11:00:00.123456789 +0100 BST",
794            "2026-09-19 10:00:00.123456789+00:00",
795            "2026-09-19 10:00:00.123456789",
796            "2026-09-19T11:00:00.123456789+01:00",
797        ] {
798            assert_eq!(timestamp(time), Some(expected), "{time}");
799        }
800        assert_eq!(timestamp("0001-01-01T00:00:00Z"), None);
801        assert_eq!(timestamp("not a time"), None);
802        assert_eq!(expected.duration_since(UNIX_EPOCH).unwrap().subsec_nanos(), 123456789);
803    }
804
805    #[test]
806    fn accounts_cumulative_provider_usage_without_repricing_or_cache_double_counting() {
807        let fixture = Fixture::new();
808        for (id, provider, metadata, messages, expected_input) in [
809            ("anthropic", "anthropic", json!({"model":"custom-model"}), json!([]), 100),
810            ("chat", "openai", json!({"api_mode":"chat_completions"}), json!([]), 100),
811            ("responses", "openai", json!({"api_mode":"responses"}), json!([]), 40),
812            ("responses-old", "openai", json!({}), json!([{"type":"message","role":"user"}]), 40),
813            ("config", "openai", json!({"config_snapshot":{"model":"another","openai":{"api_mode":"responses"}}}), json!([]), 40),
814        ] {
815            fixture.insert(id, provider, metadata, usage(), messages, json!({}));
816            let mut tracker = fixture.tracker(id);
817            tracker.refresh_all().unwrap();
818            let summary = tracker.summary();
819            assert_eq!(
820                summary.usage,
821                TokenUsage { input: expected_input, output: 20, cache_read: 60, cache_write_unsplit: 10, ..Default::default() }
822            );
823            assert!((summary.cost_usd - 0.37).abs() < 1e-12);
824            assert_eq!(summary.price_source, Some(PriceSource::Harness));
825            assert_eq!(summary.unpriced_tokens, 0);
826            assert_eq!(summary.started_at, timestamp(CREATED));
827            assert_eq!(summary.last_activity, timestamp(UPDATED));
828            assert!(summary.context.is_empty(), "cumulative usage cannot reconstruct per-response context");
829        }
830    }
831
832    #[test]
833    fn names_the_server_behind_a_kodelet_mcp_tool() {
834        assert_eq!(mcp_server_of("mcp__filesystem_read_text_file"), Some("filesystem"));
835        assert_eq!(mcp_server_of("mcp__chrome-devtools_take_screenshot"), Some("chrome-devtools"));
836        // One underscore separates server from tool, so a server whose name
837        // has one is split Kodelet's way, not ours.
838        assert_eq!(mcp_server_of("mcp__google_workspace_search"), Some("google"));
839        assert_eq!(mcp_server_of("bash"), None);
840        assert_eq!(mcp_server_of("mcp__"), None);
841        // Claude's shape is not Kodelet's: the tool half keeps its underscores.
842        assert_eq!(mcp_server_of("mcp__scratchfs__read_text_file"), Some("scratchfs"));
843    }
844
845    #[test]
846    fn counts_mcp_calls_per_server_including_the_ones_that_failed() {
847        let fixture = Fixture::new();
848        fixture.insert(
849            "mcp",
850            "anthropic",
851            json!({"model":"claude-sonnet-5"}),
852            usage(),
853            json!([]),
854            json!({
855                "c1": {"toolName":"mcp__scratchfs_read_text_file","success":true,"timestamp":"2026-09-19T10:00:01Z"},
856                "c2": {"toolName":"mcp__scratchfs_write_file","success":false,"timestamp":"2026-09-19T10:00:09Z"},
857                "c3": {"toolName":"mcp__node-repl_eval","success":true,"timestamp":"2026-09-19T10:00:05Z"},
858                "c4": {"toolName":"bash","success":true,"timestamp":"2026-09-19T10:00:07Z"}
859            }),
860        );
861        let mut tracker = fixture.tracker("mcp");
862        tracker.refresh_all().unwrap();
863        let s = tracker.summary();
864
865        let fs = s.mcp.get("scratchfs").expect("the server is named by the tool");
866        assert_eq!((fs.calls, fs.errors), (2, 1));
867        assert_eq!(fs.last_call, timestamp("2026-09-19T10:00:09Z"), "the latest call, not the latest success");
868        let repl = s.mcp.get("node-repl").unwrap();
869        assert_eq!((repl.calls, repl.errors), (1, 0));
870        assert!(!s.mcp.contains_key("bash"), "an ordinary tool is not a server");
871        assert_eq!(s.tool_calls, 4, "MCP calls are tool calls too, counted once");
872    }
873
874    #[test]
875    fn missing_or_invalid_prices_are_unpriced_not_guessed() {
876        let fixture = Fixture::new();
877        fixture.insert(
878            "missing",
879            "anthropic",
880            json!({"model":"claude-sonnet-5"}),
881            json!({"inputTokens":100,"outputTokens":20}),
882            json!([]),
883            json!({}),
884        );
885        let mut tracker = fixture.tracker("missing");
886        tracker.refresh_all().unwrap();
887        assert_eq!(tracker.summary().unpriced_tokens, 120);
888        assert_eq!(tracker.summary().cost_usd, 0.0);
889        // The model is in the built-in table, so an empty source here would let
890        // the collector name that table as this row's price. Kodelet says the
891        // row is its own either way, and `unpriced_tokens` says what is missing.
892        assert_eq!(tracker.summary().price_source, Some(PriceSource::Harness));
893
894        fixture
895            .conn
896            .execute(
897                "UPDATE conversations SET usage = ?1, updated_at = datetime(updated_at, '+1 second') WHERE id = 'missing'",
898                [json!({"inputTokens":100,"outputTokens":20,"inputCost":-1,"outputCost":0.5}).to_string()],
899            )
900            .unwrap();
901        tracker.refresh_all().unwrap();
902        assert_eq!(tracker.summary().unpriced_tokens, 100);
903        assert_eq!(tracker.summary().cost_usd, 0.5);
904        assert_eq!(tracker.summary().price_source, Some(PriceSource::Harness));
905    }
906
907    #[test]
908    fn refresh_replaces_snapshots_observes_wal_and_retries_failed_reads() {
909        let fixture = Fixture::new();
910        fixture.basic("session");
911        let mut tracker = fixture.tracker("session");
912        tracker.refresh_all().unwrap();
913        tracker.refresh_all().unwrap();
914        assert_eq!(tracker.summary().usage.input, 100);
915        assert!(tracker.conn.as_ref().unwrap().execute("DELETE FROM conversations", []).is_err());
916        assert_eq!(tracker.conn.as_ref().unwrap().query_row("PRAGMA busy_timeout", [], |r| r.get::<_, u64>(0)).unwrap(), 0);
917        fixture
918            .conn
919            .execute(
920                "UPDATE conversations SET usage = ?1, updated_at = datetime(updated_at, '+1 second')",
921                [json!({"inputTokens":150}).to_string()],
922            )
923            .unwrap();
924        tracker.refresh_all().unwrap();
925        assert_eq!(tracker.summary().usage.input, 150);
926        fixture.conn.execute("UPDATE conversations SET usage = 'not json', updated_at = datetime(updated_at, '+1 second')", []).unwrap();
927        assert!(tracker.refresh().is_err());
928        assert!(tracker.refresh().is_err(), "failed read was not cached");
929        assert_eq!(tracker.summary().usage.input, 150);
930        fixture
931            .conn
932            .execute("UPDATE conversations SET usage = ?1, updated_at = datetime(updated_at, '+1 second')", [usage().to_string()])
933            .unwrap();
934        tracker.refresh_all().unwrap();
935        assert_eq!(tracker.summary().usage.input, 100, "a corrected cumulative snapshot may shrink");
936        fixture.conn.execute("DELETE FROM conversations", []).unwrap();
937        tracker.refresh_all().unwrap();
938        assert_eq!(tracker.summary().usage.total(), 0);
939        let missing = fixture.dir.join("absent.db");
940        assert!(open_ro(&missing).is_err());
941        assert!(!missing.exists());
942    }
943
944    #[test]
945    fn counts_calls_across_provider_layouts_once_and_only_times_structured_results() {
946        let fixture = Fixture::new();
947        let tools = json!({
948            "b": {"toolName":"bash","success":false,"timestamp":"2026-09-19T10:00:05.500Z","metadata":{"executionTime":2500000000_i64}},
949            "a": {"toolName":"file_read","success":true,"timestamp":"2026-09-19T10:00:02Z"}
950        });
951        for (id, provider, metadata, messages) in [
952            (
953                "anthropic",
954                "anthropic",
955                json!({}),
956                json!([
957                    {"role":"assistant","content":[{"type":"tool_use","id":"a","name":"file_read"},{"type":"tool_use","id":"b","name":"bash"}]},
958                    {"role":"assistant","content":[{"type":"tool_use","id":"c","name":"skill"}]}
959                ]),
960            ),
961            (
962                "chat",
963                "openai",
964                json!({"api_mode":"chat_completions"}),
965                json!([
966                    {"role":"assistant","tool_calls":[{"id":"a","function":{"name":"file_read"}},{"id":"b","function":{"name":"bash"}},{"id":"c","function":{"name":"skill"}}]}
967                ]),
968            ),
969            (
970                "responses",
971                "openai",
972                json!({"api_mode":"responses"}),
973                json!([
974                    {"type":"function_call","call_id":"a","name":"file_read"},{"type":"function_call","call_id":"b","name":"bash"},
975                    {"type":"function_call_output","call_id":"b"},{"type":"function_call","call_id":"c","name":"skill"}
976                ]),
977            ),
978        ] {
979            fixture.insert(id, provider, metadata, usage(), messages, tools.clone());
980            let mut tracker = fixture.tracker(id);
981            tracker.refresh_all().unwrap();
982            let summary = tracker.summary();
983            assert_eq!(summary.tool_calls, 3, "{id}");
984            let spans = summary.spans.to_vec();
985            assert_eq!(spans.len(), 2);
986            assert_eq!(spans[0].id, "a");
987            assert_eq!(spans[0].duration_ms, Some(0));
988            assert_eq!(spans[1].started_at, timestamp("2026-09-19T10:00:03Z").unwrap());
989            assert_eq!(spans[1].duration_ms, Some(2500));
990            assert!(spans[1].error);
991            assert!(spans.iter().all(|s| s.kind == SpanKind::Tool && !s.is_open()));
992        }
993    }
994
995    #[test]
996    fn children_keep_own_usage_and_forks_do_not_recount_inherited_tools_or_turns() {
997        let fixture = Fixture::new();
998        let inherited =
999            json!({"toolName":"bash","success":true,"timestamp":"2026-09-19T09:00:00Z","metadata":{"executionTime":1000000000}});
1000        let fresh =
1001            json!({"toolName":"code_search","success":true,"timestamp":"2026-09-19T10:04:00Z","metadata":{"executionTime":2000000000_i64}});
1002        let history = json!([{"role":"assistant","content":[{"type":"tool_use","id":"old","name":"bash"}]}]);
1003        fixture.insert("parent", "anthropic", json!({}), usage(), history.clone(), json!({"old":inherited.clone()}));
1004        fixture.insert(
1005            "child",
1006            "anthropic",
1007            json!({"parent_conversation_id":"parent","conversation_name":"Scout",
1008            "conversation_fork":{"source_conversation_id":"parent","mode":"live_snapshot"}}),
1009            usage(),
1010            history.clone(),
1011            json!({"old":inherited.clone(),"new":fresh,"undated":{"toolName":"bash","success":true}}),
1012        );
1013        fixture.insert(
1014            "sibling",
1015            "anthropic",
1016            json!({"conversation_fork":{"source_conversation_id":"parent"}}),
1017            json!({}),
1018            history,
1019            json!({"old":inherited}),
1020        );
1021        let mut parent = fixture.tracker("parent");
1022        let mut child = fixture.tracker("child");
1023        let mut sibling = fixture.tracker("sibling");
1024        parent.refresh_all().unwrap();
1025        child.refresh_all().unwrap();
1026        sibling.refresh_all().unwrap();
1027        assert_eq!(parent.summary().usage.input + child.summary().usage.input, 200);
1028        assert_eq!(parent.summary().tool_calls, 1);
1029        assert_eq!(child.summary().tool_calls, 1);
1030        assert_eq!(child.summary().turns, 0, "untimestamped copied messages are not fresh turns");
1031        assert_eq!(
1032            child.summary().subagent,
1033            Some(SubagentInfo { parent_session_id: "parent".into(), nickname: Some("Scout".into()), role: None })
1034        );
1035        assert!(child.summary().spans.iter().all(|s| s.sidechain && s.id == "new"));
1036        assert!(sibling.summary().subagent.is_none(), "fork provenance is not a hierarchy parent");
1037        assert_eq!(sibling.summary().tool_calls, 0);
1038        assert_eq!(sibling.summary().activity, Activity::Unknown);
1039        assert_eq!(fixture.adapter().transcripts().len(), 3, "exports include children and siblings separately");
1040    }
1041
1042    #[test]
1043    fn adapter_shares_sidecar_index_and_preserves_central_lineage_and_accounting() {
1044        let fixture = Fixture::new();
1045        let dir = fixture.dir.join("extensions/data/subagent");
1046        std::fs::create_dir_all(&dir).unwrap();
1047        let sidecar = Connection::open(dir.join("subagents.sqlite")).unwrap();
1048        sidecar
1049            .execute_batch(
1050                "PRAGMA journal_mode=WAL;
1051            CREATE TABLE alembic_version(version_num TEXT);
1052            INSERT INTO alembic_version VALUES ('0002_canceling_state');
1053            CREATE TABLE agents(child_conversation_id TEXT, owner_conversation_id TEXT, name TEXT);",
1054            )
1055            .unwrap();
1056        let adapter = fixture.adapter();
1057        for (id, metadata, owner, expected) in [
1058            ("child", json!({}), "parent", Some(("parent", "sidecar", "subagent"))),
1059            ("matching", json!({"parent_conversation_id":"parent"}), "parent", Some(("parent", "sidecar", "subagent"))),
1060            (
1061                "search",
1062                json!({"parent_conversation_id":"parent","conversation_name":"central","profile":"code-search"}),
1063                "parent",
1064                Some(("parent", "central", "code-search")),
1065            ),
1066            (
1067                "conflict",
1068                json!({"parent_conversation_id":"other","conversation_name":"central","profile":"code-search"}),
1069                "parent",
1070                Some(("other", "central", "code-search")),
1071            ),
1072            ("self", json!({}), "self", None),
1073            ("invalid", json!({}), "../parent", None),
1074            ("central-self", json!({"parent_conversation_id":"central-self"}), "parent", None),
1075        ] {
1076            fixture.insert(id, "anthropic", metadata, usage(), json!([]), json!({}));
1077            sidecar.execute("INSERT INTO agents VALUES (?1, ?2, 'sidecar')", params![id, owner]).unwrap();
1078            // Reset the shared index to exercise a rescan without a five-second
1079            // sleep; production refreshes are throttled by the helper itself.
1080            if let Some(index) = adapter.subagents.borrow().get(&fixture.dir) {
1081                *index.borrow_mut() = SubagentIndex::new(fixture.dir.clone());
1082            }
1083            let mut tracker = adapter.open(&session_path(&fixture.db(), id), SpanRetention::All);
1084            tracker.refresh_all().unwrap();
1085            let summary = tracker.summary();
1086            assert_eq!(
1087                summary.subagent.as_ref().map(|info| (
1088                    info.parent_session_id.as_str(),
1089                    info.nickname.as_deref().unwrap(),
1090                    info.role.as_deref().unwrap()
1091                )),
1092                expected,
1093                "{id}"
1094            );
1095            assert_eq!(summary.usage.input, 100);
1096            assert!((summary.cost_usd - 0.37).abs() < 1e-12);
1097            assert_eq!(summary.activity, Activity::Unknown, "sidecars never imply activity");
1098        }
1099
1100        let path = session_path(&fixture.db(), "child");
1101        let mut first = adapter.open(&path, SpanRetention::All);
1102        let mut second = adapter.open(&path, SpanRetention::All);
1103        let index = adapter.subagents.borrow().get(&fixture.dir).unwrap().clone();
1104        assert_eq!(adapter.subagents.borrow().len(), 1);
1105        assert_eq!(Rc::strong_count(&index), 4, "adapter and both trackers share one index");
1106        first.refresh_all().unwrap();
1107        second.refresh_all().unwrap();
1108        sidecar.execute("UPDATE agents SET name = 'renamed' WHERE child_conversation_id = 'child'", []).unwrap();
1109        *index.borrow_mut() = SubagentIndex::new(fixture.dir.clone());
1110        first.refresh_all().unwrap();
1111        second.refresh_all().unwrap();
1112        assert_eq!(first.summary().subagent.as_ref().unwrap().nickname.as_deref(), Some("renamed"));
1113        assert_eq!(second.summary().subagent, first.summary().subagent, "sidecar changes refresh unchanged central snapshots");
1114    }
1115
1116    #[test]
1117    fn turn_receipts_override_copied_history_and_update_without_conversation_save() {
1118        let fixture = Fixture::new();
1119        fixture.runtime();
1120        fixture.insert(
1121            "child",
1122            "openai",
1123            json!({"conversation_fork":{"source_conversation_id":"parent"}}),
1124            usage(),
1125            json!([{"role":"assistant"},{"role":"assistant"},{"role":"assistant"}]),
1126            json!({}),
1127        );
1128        fixture.receipt("child", "running");
1129        let mut tracker = fixture.tracker("child");
1130        tracker.refresh_all().unwrap();
1131        assert_eq!(tracker.summary().turns, 1);
1132        assert_eq!(tracker.summary().activity, Activity::Working);
1133        fixture.conn.execute("UPDATE chat_turns SET status = 'succeeded'", []).unwrap();
1134        tracker.refresh_all().unwrap();
1135        assert_eq!(tracker.summary().activity, Activity::Waiting);
1136        assert_eq!(tracker.summary().turns, 1);
1137    }
1138
1139    #[test]
1140    fn unrelated_database_writes_do_not_rebuild_unchanged_sessions() {
1141        let fixture = Fixture::new();
1142        fixture.runtime();
1143        fixture.basic("session");
1144        fixture.basic("other");
1145        fixture.receipt("session", "succeeded");
1146        let mut tracker = fixture.tracker("session");
1147        tracker.refresh_all().unwrap();
1148        // The old String remains allocated until a replacement summary has
1149        // been built, so its address changes if we unnecessarily rebuild it.
1150        let model_address = tracker.summary().model.as_ref().unwrap().as_ptr();
1151        let version = tracker.data_version;
1152        fixture
1153            .conn
1154            .execute(
1155                "UPDATE conversations SET usage = ?1, updated_at = datetime(updated_at, '+1 second') WHERE id = 'other'",
1156                [usage().to_string()],
1157            )
1158            .unwrap();
1159        tracker.refresh_all().unwrap();
1160        assert_ne!(tracker.data_version, version, "the shared database changed");
1161        assert_eq!(tracker.summary().model.as_ref().unwrap().as_ptr(), model_address, "only the cheap session signature was read");
1162        assert_eq!(tracker.summary().usage.input, 100);
1163    }
1164
1165    #[test]
1166    fn observed_tool_counts_survive_compaction_but_cold_starts_report_only_retained_calls() {
1167        for fork in [false, true] {
1168            let fixture = Fixture::new();
1169            let metadata = if fork {
1170                json!({"api_mode":"responses","conversation_fork":{"source_conversation_id":"parent"}})
1171            } else {
1172                json!({"api_mode":"responses"})
1173            };
1174            let result = |name| json!({"toolName":name,"success":true,"timestamp":"2026-09-19T10:00:10Z"});
1175            fixture.insert(
1176                "session",
1177                "openai",
1178                metadata,
1179                usage(),
1180                json!([]),
1181                json!({
1182                    "call-1": result("bash"), "call-2": result("openai_web_search")
1183                }),
1184            );
1185            let mut live = fixture.tracker("session");
1186            live.refresh_all().unwrap();
1187            assert_eq!(live.summary().tool_calls, 2);
1188            assert!(live.summary().tool_calls_lower_bound);
1189
1190            // Matches ResetContextStateLocked: old results disappear but
1191            // cumulative usage remains; new calls then fill the fresh context.
1192            fixture
1193                .conn
1194                .execute(
1195                    "UPDATE conversations SET raw_messages = ?1, tool_results = ?2, updated_at = datetime(updated_at, '+1 second')",
1196                    params![
1197                        json!([{"type":"compaction"},{"type":"function_call","call_id":"call-3","name":"bash"}]).to_string(),
1198                        json!({"call-3":result("bash")}).to_string()
1199                    ],
1200                )
1201                .unwrap();
1202            live.refresh_all().unwrap();
1203            assert_eq!(live.summary().tool_calls, 3, "previously observed calls must not disappear");
1204            assert_eq!(live.summary().web_searches, 1);
1205            live.refresh_all().unwrap();
1206            assert_eq!(live.summary().tool_calls, 3, "refresh must not count a call twice");
1207
1208            let mut cold = fixture.tracker("session");
1209            cold.refresh_all().unwrap();
1210            assert_eq!(cold.summary().tool_calls, 1, "deleted history is not recoverable on a cold start");
1211            assert!(cold.summary().tool_calls_lower_bound, "one must be displayed as ≥1, not an exact lifetime total");
1212            assert_eq!(cold.summary().usage, live.summary().usage, "tokens still come from cumulative usage");
1213
1214            fixture.conn.execute("DELETE FROM conversations", []).unwrap();
1215            live.refresh_all().unwrap();
1216            assert_eq!(live.summary().tool_calls, 0, "a deleted session does not retain ghost counters");
1217        }
1218    }
1219
1220    #[test]
1221    fn failed_admissions_and_inherited_messages_are_not_parser_drift_evidence() {
1222        let fixture = Fixture::new();
1223        fixture.runtime();
1224        fixture.insert("failed", "anthropic", json!({}), json!({}), json!([{"role":"user"}]), json!({}));
1225        for i in 0..3 {
1226            fixture
1227                .conn
1228                .execute("INSERT INTO chat_turns VALUES ('failed', ?1, 'failed', ?2, ?2)", params![format!("turn-{i}"), CREATED])
1229                .unwrap();
1230        }
1231        let mut tracker = fixture.tracker("failed");
1232        tracker.refresh_all().unwrap();
1233        assert_eq!(tracker.summary().turns, 3);
1234        assert_eq!(tracker.summary().health.billable_messages, 0);
1235        assert!(!tracker.summary().health.fields_unrecognised());
1236        let replies = json!([{"role":"assistant"},{"role":"assistant"},{"role":"assistant"}]);
1237        fixture.insert("drift", "anthropic", json!({}), json!({"renamedTokens":100}), replies.clone(), json!({}));
1238        fixture.insert("fork", "anthropic", json!({"conversation_fork":{"source_conversation_id":"drift"}}), json!({}), replies, json!({}));
1239        let mut drift = fixture.tracker("drift");
1240        let mut fork = fixture.tracker("fork");
1241        drift.refresh_all().unwrap();
1242        fork.refresh_all().unwrap();
1243        assert!(drift.summary().health.fields_unrecognised());
1244        assert!(!fork.summary().health.fields_unrecognised());
1245    }
1246
1247    #[test]
1248    fn live_retention_is_bounded_but_exports_keep_all_spans() {
1249        let fixture = Fixture::new();
1250        let count = super::super::MAX_SPANS + 20;
1251        let tools: serde_json::Map<String, Value> = (0..count)
1252            .map(|i| {
1253                (
1254                    format!("call-{i:04}"),
1255                    json!({"toolName":"bash","success":true,"timestamp":"2026-09-19T10:00:10Z","metadata":{"executionTime":1000000}}),
1256                )
1257            })
1258            .collect();
1259        fixture.insert("session", "anthropic", json!({}), usage(), json!([]), tools.into());
1260        let mut live = KodeletTranscript::new(&session_path(&fixture.db(), "session"), SpanRetention::Recent);
1261        let mut export = fixture.tracker("session");
1262        live.refresh_all().unwrap();
1263        export.refresh_all().unwrap();
1264        assert_eq!(live.summary().tool_calls, count as u64);
1265        assert_eq!(live.summary().spans.len(), super::super::MAX_SPANS);
1266        assert_eq!(export.summary().spans.len(), count);
1267        assert_eq!(live.summary().spans.iter().next().unwrap().id, "call-0020");
1268    }
1269
1270    #[test]
1271    fn discovery_reads_only_metadata_and_exports_all_valid_ids() {
1272        let fixture = Fixture::new();
1273        for id in ["old", "new", "../escape"] {
1274            fixture.basic(id);
1275        }
1276        fixture.conn.execute("UPDATE conversations SET raw_messages = 'not json', tool_results = 'not json'", []).unwrap();
1277        fixture.conn.execute("UPDATE conversations SET updated_at = '2026-09-18 10:00:00 +0000 UTC' WHERE id = 'old'", []).unwrap();
1278        let sessions = recent_sessions(&fixture.db(), timestamp(CREATED).unwrap());
1279        assert_eq!(sessions.iter().map(|s| s.id.as_str()).collect::<Vec<_>>(), vec!["new"]);
1280        let adapter = fixture.adapter();
1281        assert_eq!(adapter.transcripts().len(), 2);
1282        let attached = HashSet::from([session_path(&fixture.db(), "new")]);
1283        assert_eq!(adapter.unowned(&attached), vec![session_path(&fixture.db(), "old")]);
1284        assert!(!adapter.detect(&fixture.db()));
1285    }
1286
1287    #[test]
1288    fn daemon_receipts_are_exact_and_thin_clients_cannot_steal_them() {
1289        let fixture = Fixture::new();
1290        fixture.runtime();
1291        fixture.basic("parent");
1292        fixture.basic("child");
1293        fixture.basic("idle");
1294        fixture.receipt("parent", "running");
1295        fixture.receipt("child", "accepted");
1296        fixture.receipt("idle", "succeeded");
1297        fixture.json_file(
1298            "server/connection.json",
1299            json!({"schemaVersion":1,"pid":10,"instanceId":"daemon","version":"0.6.17-beta","managed":false}),
1300        );
1301        let (daemon_raw, daemon) = process(10, &["kodelet", "--profile", "work", "serve"]);
1302        let (client_raw, client) = process(20, &["kodelet", "run", "--resume", "parent"]);
1303        let mut adapter = fixture.adapter();
1304        adapter.prepare(&[&client, &daemon], &HashMap::from([(10, &daemon_raw), (20, &client_raw)]));
1305        let attached = HashSet::new();
1306        assert_eq!(adapter.attribute(&client, Some(&client_raw), &context(&attached)).1, Attribution::None);
1307        let (paths, attribution) = adapter.attribute(&daemon, Some(&daemon_raw), &context(&attached));
1308        assert_eq!(attribution, Attribution::HarnessRegistry);
1309        assert_eq!(
1310            paths.into_iter().collect::<HashSet<_>>(),
1311            HashSet::from([session_path(&fixture.db(), "parent"), session_path(&fixture.db(), "child")])
1312        );
1313        assert_eq!(adapter.hints(10).unwrap().version.as_deref(), Some("0.6.17-beta"));
1314        assert_eq!(adapter.hints(10).unwrap().status.as_deref(), Some("busy"));
1315        // File existence and a matching PID are insufficient after PID reuse.
1316        let mut reused = daemon_raw.clone();
1317        reused.start_time = SystemTime::now().duration_since(UNIX_EPOCH).unwrap().as_secs() + 10;
1318        adapter.prepare(&[&daemon], &HashMap::from([(10, &reused)]));
1319        assert!(adapter.attribute(&daemon, Some(&reused), &context(&attached)).0.is_empty());
1320    }
1321
1322    #[test]
1323    fn local_runner_identity_and_freshness_prevent_remote_or_stale_pid_matches() {
1324        let fixture = Fixture::new();
1325        fixture.runtime();
1326        fixture.basic("local");
1327        fixture.basic("remote");
1328        fixture.json_file("runners/host.json", json!({"version":1,"instanceId":"this-host"}));
1329        fixture
1330            .conn
1331            .execute_batch(
1332                "INSERT INTO runner_registrations VALUES ('r-local',10,'this-host','busy',datetime('now'),datetime('now'),'v-local');
1333            INSERT INTO runner_registrations VALUES ('r-remote',20,'other-host','busy',datetime('now'),datetime('now'),'v-remote');
1334            INSERT INTO runner_runs VALUES ('run-local','local','r-local','running');
1335            INSERT INTO runner_runs VALUES ('run-remote','remote','r-remote','running');",
1336            )
1337            .unwrap();
1338        let (local_raw, local) = process(10, &["kodelet", "runner", "start"]);
1339        let (collision_raw, collision) = process(20, &["kodelet", "serve"]);
1340        let mut adapter = fixture.adapter();
1341        adapter.prepare(&[&collision, &local], &HashMap::from([(10, &local_raw), (20, &collision_raw)]));
1342        let attached = HashSet::new();
1343        assert_eq!(
1344            adapter.attribute(&local, Some(&local_raw), &context(&attached)),
1345            (vec![session_path(&fixture.db(), "local")], Attribution::HarnessRegistry)
1346        );
1347        assert!(adapter.attribute(&collision, Some(&collision_raw), &context(&attached)).0.is_empty());
1348        fixture.conn.execute("UPDATE runner_registrations SET last_heartbeat_at = datetime('now','-60 seconds')", []).unwrap();
1349        adapter.prepare(&[&local], &HashMap::from([(10, &local_raw)]));
1350        assert!(adapter.attribute(&local, Some(&local_raw), &context(&attached)).0.is_empty());
1351    }
1352
1353    #[test]
1354    fn pre_registry_cwd_fallback_is_labelled_unambiguous_and_never_for_clients() {
1355        let fixture = Fixture::new();
1356        fixture.basic("session");
1357        let (raw, node) = process(10, &["kodelet", "serve"]);
1358        let (client_raw, client) = process(20, &["kodelet", "run"]);
1359        let mut adapter = fixture.adapter();
1360        let attached = HashSet::new();
1361        assert_eq!(
1362            adapter.attribute(&node, Some(&raw), &context(&attached)),
1363            (vec![session_path(&fixture.db(), "session")], Attribution::CwdHeuristic)
1364        );
1365        assert!(adapter.attribute(&client, Some(&client_raw), &context(&attached)).0.is_empty());
1366        adapter.reserve(&fixture.db(), 30, "session", Attribution::HarnessRegistry);
1367        assert!(adapter.attribute(&node, Some(&raw), &context(&attached)).0.is_empty());
1368        fixture.basic("other");
1369        adapter = fixture.adapter();
1370        assert!(adapter.attribute(&node, Some(&raw), &context(&attached)).0.is_empty(), "ambiguous cwd is not attribution");
1371    }
1372}