1use super::{AttributeContext, ContextLedger, HarnessAdapter, RegistryHints, SessionSummary, SessionTracker, SpanRetention};
72use crate::model::{Activity, Attribution, ContextOrigin, CostBreakdown, Harness, ProcNode, SpanKind, TokenUsage};
73use crate::process::RawProc;
74use rusqlite::{Connection, OpenFlags};
75use std::collections::HashSet;
76use std::path::{Path, PathBuf};
77use std::time::{Duration, SystemTime, UNIX_EPOCH};
78
79pub fn data_dir() -> Option<PathBuf> {
81 if let Some(d) = std::env::var_os("OPENCODE_DATA_DIR") {
82 return Some(PathBuf::from(d));
83 }
84 if let Some(d) = std::env::var_os("XDG_DATA_HOME") {
85 return Some(PathBuf::from(d).join("opencode"));
86 }
87 std::env::var_os("HOME").map(|h| PathBuf::from(h).join(".local/share/opencode"))
88}
89
90pub fn db_path() -> Option<PathBuf> {
92 let p = data_dir()?.join("opencode.db");
93 p.exists().then_some(p)
94}
95
96fn open_ro(db: &Path) -> rusqlite::Result<Connection> {
100 Connection::open_with_flags(db, OpenFlags::SQLITE_OPEN_READ_ONLY)
101}
102
103fn to_ms(t: SystemTime) -> i64 {
104 t.duration_since(UNIX_EPOCH).map(|d| d.as_millis() as i64).unwrap_or(0)
105}
106
107fn from_ms(ms: i64) -> Option<SystemTime> {
108 (ms > 0).then(|| UNIX_EPOCH + Duration::from_millis(ms as u64))
109}
110
111#[derive(Debug, Clone, PartialEq, Eq)]
113pub struct Session {
114 pub id: String,
115 pub directory: PathBuf,
116 pub created: Option<SystemTime>,
117 pub updated: Option<SystemTime>,
118}
119
120pub fn session_path(db: &Path, session_id: &str) -> PathBuf {
123 db.join(session_id)
124}
125
126pub fn session_id_of(path: &Path) -> Option<String> {
128 path.file_name().map(|f| f.to_string_lossy().into_owned())
129}
130
131pub fn recent_sessions(db: &Path, since: SystemTime) -> Vec<Session> {
133 let Ok(conn) = open_ro(db) else { return Vec::new() };
134 let sql = "SELECT id, directory, time_created, time_updated FROM session \
135 WHERE parent_id IS NULL AND time_updated >= ?1 ORDER BY time_updated DESC";
136 let Ok(mut stmt) = conn.prepare(sql) else { return Vec::new() };
137 let rows = stmt.query_map([to_ms(since)], |r| {
138 Ok(Session {
139 id: r.get::<_, String>(0)?,
140 directory: PathBuf::from(r.get::<_, String>(1)?),
141 created: from_ms(r.get::<_, i64>(2)?),
142 updated: from_ms(r.get::<_, i64>(3)?),
143 })
144 });
145 rows.map(|it| it.flatten().collect()).unwrap_or_default()
146}
147
148#[derive(Debug, Clone, Default)]
151pub struct ConfigRoots {
152 pub global: Option<PathBuf>,
154 pub home: Option<PathBuf>,
156}
157
158impl ConfigRoots {
159 pub fn detect() -> ConfigRoots {
160 let home = std::env::var_os("HOME").map(PathBuf::from);
161 let global = std::env::var_os("XDG_CONFIG_HOME")
162 .map(|d| PathBuf::from(d).join("opencode"))
163 .or_else(|| home.as_ref().map(|h| h.join(".config/opencode")));
164 ConfigRoots { global, home }
165 }
166}
167
168fn config_files(directory: &Path, roots: &ConfigRoots) -> Vec<PathBuf> {
174 let mut files = Vec::new();
175 if let Some(g) = &roots.global {
176 for name in ["config.json", "opencode.json", "opencode.jsonc"] {
177 files.push(g.join(name));
178 }
179 }
180 let project = |dir: &Path, files: &mut Vec<PathBuf>| {
181 for name in ["opencode.jsonc", "opencode.json"] {
182 files.push(dir.join(name));
183 files.push(dir.join(".opencode").join(name));
184 }
185 };
186 for dir in directory.ancestors() {
187 project(dir, &mut files);
188 if dir.join(".git").exists() {
189 break;
190 }
191 }
192 if let Some(h) = &roots.home {
193 for name in ["opencode.jsonc", "opencode.json"] {
194 files.push(h.join(".opencode").join(name));
195 }
196 }
197 files
198}
199
200pub fn mcp_server_names(directory: &Path, roots: &ConfigRoots) -> Vec<String> {
204 let mut names: Vec<String> = config_files(directory, roots)
205 .iter()
206 .filter_map(|f| std::fs::read_to_string(f).ok())
207 .filter_map(|text| serde_json::from_str::<serde_json::Value>(&strip_jsonc(&text)).ok())
208 .filter_map(|v| v.get("mcp").and_then(|m| m.as_object()).map(|m| m.keys().cloned().collect::<Vec<_>>()))
209 .flatten()
210 .collect();
211 names.sort();
212 names.dedup();
213 names
214}
215
216fn strip_jsonc(text: &str) -> String {
220 let mut out = String::with_capacity(text.len());
221 let mut chars = text.chars().peekable();
222 let mut in_string = false;
223 while let Some(c) = chars.next() {
224 if in_string {
225 out.push(c);
226 match c {
227 '\\' => {
228 if let Some(n) = chars.next() {
229 out.push(n);
230 }
231 }
232 '"' => in_string = false,
233 _ => {}
234 }
235 continue;
236 }
237 match (c, chars.peek()) {
238 ('"', _) => {
239 in_string = true;
240 out.push(c);
241 }
242 ('/', Some('/')) => {
243 for n in chars.by_ref() {
244 if n == '\n' {
245 out.push('\n');
246 break;
247 }
248 }
249 }
250 ('/', Some('*')) => {
251 chars.next();
252 let mut prev = ' ';
253 for n in chars.by_ref() {
254 if prev == '*' && n == '/' {
255 break;
256 }
257 prev = n;
258 }
259 }
260 _ => out.push(c),
261 }
262 }
263 let bytes: Vec<char> = out.chars().collect();
265 let mut cleaned = String::with_capacity(out.len());
266 let mut in_string = false;
267 let mut i = 0;
268 while i < bytes.len() {
269 let c = bytes[i];
270 if in_string {
271 cleaned.push(c);
272 if c == '\\' && i + 1 < bytes.len() {
273 cleaned.push(bytes[i + 1]);
274 i += 1;
275 } else if c == '"' {
276 in_string = false;
277 }
278 } else if c == '"' {
279 in_string = true;
280 cleaned.push(c);
281 } else if c == ',' {
282 let next = bytes[i + 1..].iter().find(|n| !n.is_whitespace());
283 if !matches!(next, Some('}') | Some(']')) {
284 cleaned.push(c);
285 }
286 } else {
287 cleaned.push(c);
288 }
289 i += 1;
290 }
291 cleaned
292}
293
294fn sanitize(name: &str) -> String {
297 name.chars().map(|c| if c.is_ascii_alphanumeric() || c == '_' || c == '-' { c } else { '_' }).collect()
298}
299
300pub fn mcp_server_of<'a>(tool_name: &str, servers: &'a [String]) -> Option<&'a str> {
306 servers
307 .iter()
308 .map(|s| (s, sanitize(s)))
309 .filter(|(_, key)| {
310 tool_name.len() > key.len() + 1 && tool_name.starts_with(key.as_str()) && tool_name.as_bytes()[key.len()] == b'_'
311 })
312 .max_by_key(|(_, key)| key.len())
313 .map(|(s, _)| s.as_str())
314}
315
316fn model_id(raw: &str) -> Option<String> {
319 let raw = raw.trim();
320 if raw.is_empty() {
321 return None;
322 }
323 match serde_json::from_str::<serde_json::Value>(raw) {
324 Ok(v) => v.get("id").and_then(|x| x.as_str()).map(str::to_string).or_else(|| Some(raw.to_string())),
325 Err(_) => Some(raw.to_string()),
326 }
327}
328
329#[derive(Default)]
331pub struct OpenCodeAdapter {
332 db: Option<PathBuf>,
333 recent: Vec<Session>,
334}
335
336impl OpenCodeAdapter {
337 fn db(&self) -> Option<PathBuf> {
338 self.db.clone().or_else(db_path)
339 }
340}
341
342impl HarnessAdapter for OpenCodeAdapter {
343 fn harness(&self) -> Harness {
344 Harness::OpenCode
345 }
346
347 fn rescan(&mut self, since: SystemTime) {
348 self.db = db_path();
349 self.recent = match &self.db {
350 Some(db) => recent_sessions(db, since),
351 None => Vec::new(),
352 };
353 }
354
355 fn hints(&self, _pid: u32) -> Option<RegistryHints> {
356 None
357 }
358
359 fn attribute(&self, _root: &ProcNode, _raw: Option<&RawProc>, ctx: &AttributeContext) -> (Vec<PathBuf>, Attribution) {
363 let (Some(cwd), Some(db)) = (ctx.cwd, self.db()) else { return (Vec::new(), Attribution::None) };
364 let slack = Duration::from_secs(60);
365 let mut mine: Vec<&Session> = self
366 .recent
367 .iter()
368 .filter(|s| s.directory == cwd)
369 .filter(|s| s.created.is_none_or(|c| c + slack >= ctx.proc_start))
370 .filter(|s| !ctx.attached.contains(&session_path(&db, &s.id)))
371 .collect();
372 mine.sort_by_key(|s| std::cmp::Reverse(s.updated));
373 match mine.first() {
374 Some(s) => (vec![session_path(&db, &s.id)], Attribution::CwdHeuristic),
375 None => (Vec::new(), Attribution::None),
376 }
377 }
378
379 fn unowned(&self, attached: &HashSet<PathBuf>) -> Vec<PathBuf> {
380 let Some(db) = self.db() else { return Vec::new() };
381 self.recent.iter().map(|s| session_path(&db, &s.id)).filter(|p| !attached.contains(p)).collect()
382 }
383
384 fn open(&self, path: &Path, spans: SpanRetention) -> Box<dyn SessionTracker> {
385 Box::new(OpenCodeTranscript::new(path, spans))
386 }
387
388 fn detect(&self, _path: &Path) -> bool {
391 false
392 }
393
394 fn transcripts(&self) -> Vec<(String, PathBuf)> {
395 let Some(db) = self.db() else { return Vec::new() };
396 recent_sessions(&db, UNIX_EPOCH).into_iter().map(|s| (s.id.clone(), session_path(&db, &s.id))).collect()
397 }
398}
399
400pub struct OpenCodeTranscript {
407 db: PathBuf,
408 session_id: String,
409 virtual_path: PathBuf,
410 conn: Option<Connection>,
411 retention: SpanRetention,
412 summary: SessionSummary,
413 last_updated: Option<i64>,
415 config: ConfigRoots,
417}
418
419impl OpenCodeTranscript {
420 pub fn new(virtual_path: &Path, retention: SpanRetention) -> Self {
421 let session_id = session_id_of(virtual_path).unwrap_or_default();
423 let db = virtual_path.parent().map(Path::to_path_buf).unwrap_or_default();
424 OpenCodeTranscript {
425 db,
426 session_id,
427 virtual_path: virtual_path.to_path_buf(),
428 conn: None,
429 retention,
430 summary: SessionSummary {
431 harness: Some(Harness::OpenCode),
432 folds_child_usage: true,
433 spans: retention.log(),
434 ..Default::default()
435 },
436 last_updated: None,
437 config: ConfigRoots::detect(),
438 }
439 }
440
441 pub fn with_config(mut self, config: ConfigRoots) -> Self {
443 self.config = config;
444 self
445 }
446
447 fn conn(&mut self) -> Option<&Connection> {
448 if self.conn.is_none() {
449 self.conn = open_ro(&self.db).ok();
450 }
451 self.conn.as_ref()
452 }
453
454 fn reload(&mut self) -> rusqlite::Result<()> {
455 let retention = self.retention;
456 let id = self.session_id.clone();
457 let config = self.config.clone();
458 let Some(conn) = self.conn() else { return Ok(()) };
459
460 let mut summary = SessionSummary {
466 harness: Some(Harness::OpenCode),
467 folds_child_usage: true,
470 price_source: Some(crate::model::PriceSource::Harness),
471 spans: retention.log(),
472 ..Default::default()
473 };
474 let mut ids: Vec<(String, bool)> = Vec::new(); {
476 let sql = "SELECT id, directory, agent, model, cost, tokens_input, tokens_output, tokens_reasoning, \
477 tokens_cache_read, tokens_cache_write, time_created, time_updated, version, parent_id \
478 FROM session WHERE id = ?1 OR parent_id = ?1 ORDER BY (parent_id IS NOT NULL), time_created";
479 let mut stmt = conn.prepare(sql)?;
480 let mut rows = stmt.query([&id])?;
481 while let Some(r) = rows.next()? {
482 let row_id: String = r.get(0)?;
483 let parent: Option<String> = r.get(13)?;
484 let is_sub = parent.is_some();
485 ids.push((row_id.clone(), is_sub));
486
487 let usage = TokenUsage {
488 cache_write_unsplit: 0,
489 input: r.get::<_, i64>(5)? as u64,
490 output: (r.get::<_, i64>(6)? + r.get::<_, i64>(7)?) as u64, cache_read: r.get::<_, i64>(8)? as u64,
492 cache_write_5m: r.get::<_, i64>(9)? as u64,
493 cache_write_1h: 0,
494 };
495 summary.usage.add(&usage);
496 summary.cost_usd += r.get::<_, f64>(4)?;
497
498 if !is_sub {
499 summary.session_id = Some(row_id.clone());
500 summary.cwd = Some(PathBuf::from(r.get::<_, String>(1)?));
501 summary.model = r.get::<_, Option<String>>(3)?.as_deref().and_then(model_id);
502 summary.harness_version = r.get::<_, Option<String>>(12)?;
503 summary.started_at = from_ms(r.get::<_, i64>(10)?);
504 summary.last_activity = from_ms(r.get::<_, i64>(11)?);
505 }
506 }
507 }
508 if ids.is_empty() {
509 self.summary = summary;
511 return Ok(());
512 }
513
514 for (sid, is_sub) in &ids {
516 let n: i64 = conn.query_row(
517 "SELECT count(*) FROM message WHERE session_id = ?1 AND json_extract(data,'$.role') = 'assistant'",
518 [sid],
519 |r| r.get(0),
520 )?;
521 summary.turns += n as u64;
522 summary.health.billable_messages += n as u64;
523 if *is_sub {
524 summary.subagent_turns += n as u64;
525 }
526 }
527 summary.health.usage_records = summary.health.billable_messages;
528 if summary.usage.total() == 0 {
529 summary.health.empty_usage_records = summary.health.usage_records;
530 }
531
532 for (sid, _is_sub) in &ids {
536 summary.tool_calls +=
537 conn.query_row("SELECT count(*) FROM part WHERE session_id = ?1 AND json_extract(data,'$.type') = 'tool'", [sid], |r| {
538 r.get::<_, i64>(0)
539 })? as u64;
540 }
541 let limit = match retention {
542 SpanRetention::All => -1,
543 SpanRetention::Recent => super::MAX_SPANS as i64,
544 };
545 let sub_ids: HashSet<&str> = ids.iter().filter(|(_, s)| *s).map(|(i, _)| i.as_str()).collect();
546 let placeholders = ids.iter().map(|_| "?").collect::<Vec<_>>().join(",");
547
548 let servers = summary.cwd.as_deref().map(|d| mcp_server_names(d, &config)).unwrap_or_default();
551 for (sid, is_sub) in &ids {
552 let ledger = context_ledger(conn, sid, &servers)?;
553 if *is_sub {
554 summary.context.merge(&ledger);
555 } else {
556 let subs = std::mem::take(&mut summary.context);
557 summary.context = ledger;
558 summary.context.merge(&subs);
559 }
560 }
561
562 if !servers.is_empty() {
567 let sql = format!(
568 "SELECT json_extract(data,'$.tool'), json_extract(data,'$.state.status'), \
569 json_extract(data,'$.state.time.end'), json_extract(data,'$.state.time.start') \
570 FROM part WHERE json_extract(data,'$.type') = 'tool' AND session_id IN ({placeholders})"
571 );
572 let mut stmt = conn.prepare(&sql)?;
573 let params: Vec<&dyn rusqlite::ToSql> = ids.iter().map(|(i, _)| i as &dyn rusqlite::ToSql).collect();
574 let mut rows = stmt.query(params.as_slice())?;
575 while let Some(r) = rows.next()? {
576 let Some(name) = r.get::<_, Option<String>>(0)? else { continue };
577 let Some(server) = mcp_server_of(&name, &servers) else { continue };
578 let error = r.get::<_, Option<String>>(1)?.as_deref() == Some("error");
579 let at = r.get::<_, Option<i64>>(2)?.or(r.get::<_, Option<i64>>(3)?).and_then(from_ms);
580 let u = summary.mcp.entry(server.to_string()).or_default();
581 u.calls += 1;
582 u.errors += u64::from(error);
583 u.last_call = u.last_call.max(at);
584 }
585 }
586
587 struct Pending {
590 id: String,
591 name: String,
592 kind: SpanKind,
593 start: SystemTime,
594 end: Option<SystemTime>,
597 sidechain: bool,
598 error: bool,
599 }
600 let mut pending: Vec<Pending> = Vec::new();
601
602 {
604 let sql = format!(
605 "SELECT session_id, json_extract(data,'$.tool'), json_extract(data,'$.callID'), \
606 json_extract(data,'$.state.status'), json_extract(data,'$.state.time.start'), \
607 json_extract(data,'$.state.time.end') \
608 FROM part WHERE json_extract(data,'$.type') = 'tool' AND session_id IN ({placeholders}) \
609 ORDER BY time_created DESC LIMIT ?{}",
610 ids.len() + 1
611 );
612 let mut stmt = conn.prepare(&sql)?;
613 let params: Vec<&dyn rusqlite::ToSql> =
614 ids.iter().map(|(i, _)| i as &dyn rusqlite::ToSql).chain(std::iter::once(&limit as &dyn rusqlite::ToSql)).collect();
615 let mut rows = stmt.query(params.as_slice())?;
616 let mut i = 0;
617 while let Some(r) = rows.next()? {
618 let sid: String = r.get(0)?;
619 let name: String = r.get::<_, Option<String>>(1)?.unwrap_or_else(|| "tool".into());
620 let call_id: String = r.get::<_, Option<String>>(2)?.unwrap_or_default();
621 let status: Option<String> = r.get(3)?;
622 let Some(start) = r.get::<_, Option<i64>>(4)?.and_then(from_ms) else { continue };
623 let end = r.get::<_, Option<i64>>(5)?.and_then(from_ms);
624 let id = if call_id.is_empty() { format!("oc-tool-{i}") } else { call_id };
625 i += 1;
626 pending.push(Pending {
627 id,
628 name,
629 kind: SpanKind::Tool,
630 start,
631 end: Some(end.unwrap_or(start).max(start)),
632 sidechain: sub_ids.contains(sid.as_str()),
633 error: status.as_deref() == Some("error"),
634 });
635 }
636 }
637
638 {
644 let sql = format!(
645 "SELECT session_id, json_extract(data,'$.role'), json_extract(data,'$.time.created'), \
646 json_extract(data,'$.time.completed') FROM message WHERE session_id IN ({placeholders}) \
647 ORDER BY json_extract(data,'$.time.created') DESC LIMIT ?{}",
648 ids.len() + 1
649 );
650 let mut stmt = conn.prepare(&sql)?;
651 let params: Vec<&dyn rusqlite::ToSql> =
652 ids.iter().map(|(i, _)| i as &dyn rusqlite::ToSql).chain(std::iter::once(&limit as &dyn rusqlite::ToSql)).collect();
653 let mut rows = stmt.query(params.as_slice())?;
654 struct Msg {
655 sid: String,
656 role: String,
657 created: SystemTime,
658 completed: Option<SystemTime>,
659 }
660 let mut msgs: Vec<Msg> = Vec::new();
661 while let Some(r) = rows.next()? {
662 let sid: String = r.get(0)?;
663 let role: String = r.get::<_, Option<String>>(1)?.unwrap_or_default();
664 let Some(created) = r.get::<_, Option<i64>>(2)?.and_then(from_ms) else { continue };
665 let completed = r.get::<_, Option<i64>>(3)?.and_then(from_ms);
666 msgs.push(Msg { sid, role, created, completed });
667 }
668 msgs.sort_by_key(|m| m.created);
669 match msgs.last() {
672 Some(m) if m.role == "assistant" && m.completed.is_some() => summary.activity = Activity::Waiting,
673 Some(_) => summary.activity = Activity::Working,
674 None => {}
675 }
676 let sessions: Vec<String> = ids.iter().map(|(i, _)| i.clone()).collect();
678 for sess in &sessions {
679 let sidechain = sub_ids.contains(sess.as_str());
680 let mut turn_idx = 0u64;
681 let mut turn: Option<(String, SystemTime, Option<SystemTime>)> = None; let mut inf_idx = 0u64;
683 let flush = |turn: &mut Option<(String, SystemTime, Option<SystemTime>)>, pending: &mut Vec<Pending>| {
684 if let Some((id, start, end)) = turn.take() {
685 pending.push(Pending { id, name: "turn".into(), kind: SpanKind::Turn, start, end, sidechain, error: false });
686 }
687 };
688 for m in msgs.iter().filter(|m| &m.sid == sess) {
689 if m.role == "user" {
690 flush(&mut turn, &mut pending);
691 turn_idx += 1;
692 turn = Some((format!("turn:{sess}:{turn_idx}"), m.created, None));
693 } else if m.role == "assistant" {
694 inf_idx += 1;
695 pending.push(Pending {
696 id: format!("inf:{sess}:{inf_idx}"),
697 name: "inference".into(),
698 kind: SpanKind::Inference,
699 start: m.created,
700 end: m.completed,
701 sidechain,
702 error: false,
703 });
704 let end = m.completed.unwrap_or(m.created);
707 match &mut turn {
708 Some((_, _, e)) => *e = Some(end),
709 None => {
710 turn_idx += 1;
711 turn = Some((format!("turn:{sess}:{turn_idx}"), m.created, Some(end)));
712 }
713 }
714 }
715 }
716 flush(&mut turn, &mut pending);
717 }
718 }
719
720 pending.sort_by_key(|p| p.start);
723 for p in pending {
724 summary.spans.open_kind(p.id.clone(), p.name, p.start, p.sidechain, p.kind);
725 if let Some(end) = p.end {
726 if p.error {
727 summary.spans.close(&p.id, end.max(p.start), true);
728 } else {
729 summary.spans.end_at(&p.id, end.max(p.start));
730 }
731 }
732 }
733
734 self.summary = summary;
736 Ok(())
737 }
738}
739
740fn context_ledger(conn: &Connection, session_id: &str, servers: &[String]) -> rusqlite::Result<ContextLedger> {
744 let mut calls: std::collections::HashMap<String, Vec<(String, String)>> = std::collections::HashMap::new();
745 {
746 let mut stmt = conn.prepare(
747 "SELECT message_id, json_extract(data,'$.callID'), json_extract(data,'$.tool') FROM part \
748 WHERE session_id = ?1 AND json_extract(data,'$.type') = 'tool' ORDER BY id",
749 )?;
750 let mut rows = stmt.query([session_id])?;
751 while let Some(r) = rows.next()? {
752 let (Some(msg), Some(call), Some(tool)) =
753 (r.get::<_, Option<String>>(0)?, r.get::<_, Option<String>>(1)?, r.get::<_, Option<String>>(2)?)
754 else {
755 continue;
756 };
757 calls.entry(msg).or_default().push((call, tool));
758 }
759 }
760 let mut ledger = ContextLedger::default();
761 let mut stmt = conn.prepare(
762 "SELECT id, json_extract(data,'$.tokens.input'), json_extract(data,'$.tokens.output'), \
763 json_extract(data,'$.tokens.reasoning'), json_extract(data,'$.tokens.cache.read'), \
764 json_extract(data,'$.tokens.cache.write'), json_extract(data,'$.cost'), json_extract(data,'$.modelID'), \
765 json_extract(data,'$.summary'), json_extract(data,'$.mode') FROM message \
766 WHERE session_id = ?1 AND json_extract(data,'$.role') = 'assistant' ORDER BY time_created, id",
767 )?;
768 let mut rows = stmt.query([session_id])?;
769 let n = |v: Option<i64>| v.unwrap_or(0).max(0) as u64;
770 while let Some(r) = rows.next()? {
771 let id: String = r.get(0)?;
772 let usage = TokenUsage {
773 cache_write_unsplit: 0,
774 input: n(r.get(1)?),
775 output: n(r.get(2)?) + n(r.get(3)?),
776 cache_read: n(r.get(4)?),
777 cache_write_5m: n(r.get(5)?),
778 cache_write_1h: 0,
779 };
780 let total: f64 = r.get::<_, Option<f64>>(6)?.unwrap_or(0.0);
781 let cost = r
782 .get::<_, Option<String>>(7)?
783 .and_then(|m| crate::pricing::table().lookup(&m))
784 .map(|p| scaled(p.breakdown(&usage), total))
785 .unwrap_or_default();
786 ledger.response(&usage, &cost);
787 let summary = r.get::<_, Option<i64>>(8)?.unwrap_or(0) != 0 || r.get::<_, Option<String>>(9)?.as_deref() == Some("compaction");
788 if summary {
789 ledger.compacted();
790 }
791 for (call, tool) in calls.remove(&id).unwrap_or_default() {
792 match mcp_server_of(&tool, servers) {
793 Some(server) => ledger.result(&call, ContextOrigin::Mcp, server),
794 None => ledger.result(&call, ContextOrigin::Tool, &tool),
795 }
796 }
797 }
798 Ok(ledger)
799}
800
801fn scaled(b: CostBreakdown, total: f64) -> CostBreakdown {
806 let list = b.total();
807 if total <= 0.0 || list <= 0.0 {
808 return b;
809 }
810 let k = total / list;
811 CostBreakdown {
812 input: b.input * k,
813 cache_write_5m: b.cache_write_5m * k,
814 cache_write_1h: b.cache_write_1h * k,
815 cache_write_unsplit: b.cache_write_unsplit * k,
816 cache_read: b.cache_read * k,
817 output: b.output * k,
818 web_search: b.web_search * k,
819 }
820}
821
822impl SessionTracker for OpenCodeTranscript {
823 fn refresh(&mut self) -> anyhow::Result<bool> {
824 let id = self.session_id.clone();
825 let updated: Option<i64> =
826 self.conn().and_then(|c| c.query_row("SELECT time_updated FROM session WHERE id = ?1", [&id], |r| r.get(0)).ok());
827 if updated.is_some() && updated == self.last_updated {
829 return Ok(false);
830 }
831 self.last_updated = updated;
832 let _ = self.reload();
835 Ok(false)
836 }
837
838 fn summary(&self) -> &SessionSummary {
839 &self.summary
840 }
841
842 fn path(&self) -> &Path {
843 &self.virtual_path
844 }
845}
846
847#[cfg(test)]
848mod tests {
849 use super::*;
850 use rusqlite::Connection;
851
852 fn make_db(dir: &Path) -> PathBuf {
854 let db = dir.join("opencode.db");
855 let conn = Connection::open(&db).unwrap();
856 conn.execute_batch(
857 "CREATE TABLE session (id TEXT PRIMARY KEY, project_id TEXT, parent_id TEXT, directory TEXT, agent TEXT, \
858 model TEXT, cost REAL DEFAULT 0, tokens_input INTEGER DEFAULT 0, tokens_output INTEGER DEFAULT 0, \
859 tokens_reasoning INTEGER DEFAULT 0, tokens_cache_read INTEGER DEFAULT 0, tokens_cache_write INTEGER DEFAULT 0, \
860 time_created INTEGER, time_updated INTEGER, version TEXT);
861 CREATE TABLE message (id TEXT PRIMARY KEY, session_id TEXT, time_created INTEGER, data TEXT);
862 CREATE TABLE part (id TEXT PRIMARY KEY, message_id TEXT, session_id TEXT, time_created INTEGER, data TEXT);",
863 )
864 .unwrap();
865 conn.execute(
867 "INSERT INTO session VALUES ('ses_parent', 'p', NULL, '/tmp/proj', 'build', \
868 '{\"id\":\"deepseek-v4-pro\",\"providerID\":\"deepseek\"}', 0.25, 1000, 200, 50, 900000, 0, 1000, 5000, '1.18.15')",
869 [],
870 )
871 .unwrap();
872 conn.execute(
873 "INSERT INTO session VALUES ('ses_child', 'p', 'ses_parent', '/tmp/proj', 'explore', \
874 '{\"id\":\"deepseek-v4-pro\"}', 0.05, 300, 40, 10, 1000, 0, 2000, 3000, '1.18.15')",
875 [],
876 )
877 .unwrap();
878 let msg = |role: &str, created: i64, completed: Option<i64>| match completed {
882 Some(c) => format!("{{\"role\":\"{role}\",\"time\":{{\"created\":{created},\"completed\":{c}}}}}"),
883 None => format!("{{\"role\":\"{role}\",\"time\":{{\"created\":{created}}}}}"),
884 };
885 for (i, sid, data) in [
886 (1, "ses_parent", msg("user", 100, None)),
887 (2, "ses_parent", msg("assistant", 110, Some(200))),
888 (3, "ses_parent", msg("assistant", 210, Some(300))),
889 (4, "ses_child", msg("assistant", 150, Some(180))),
890 ] {
891 conn.execute("INSERT INTO message VALUES (?1, ?2, ?3, ?4)", rusqlite::params![format!("m{i}"), sid, 1000 + i as i64, data])
892 .unwrap();
893 }
894 let tool = |tool: &str, call: &str, status: &str, start: i64, end: i64| {
896 format!(
897 "{{\"type\":\"tool\",\"tool\":\"{tool}\",\"callID\":\"{call}\",\"state\":{{\"status\":\"{status}\",\"time\":{{\"start\":{start},\"end\":{end}}}}}}}"
898 )
899 };
900 for (i, sid, data) in [
901 (1, "ses_parent", tool("read", "c1", "completed", 1000, 1300)),
902 (2, "ses_parent", tool("bash", "c2", "error", 1400, 2400)),
903 (3, "ses_child", tool("grep", "c3", "completed", 1500, 1600)),
904 ] {
905 conn.execute("INSERT INTO part VALUES (?1, 'm', ?2, ?3, ?4)", rusqlite::params![format!("p{i}"), sid, 1000 + i as i64, data])
906 .unwrap();
907 }
908 db
909 }
910
911 #[test]
912 fn reads_a_session_folds_its_subagent_and_builds_tool_spans() {
913 let dir = std::env::temp_dir().join(format!("agent-top-oc-{}", std::process::id()));
914 let _ = std::fs::remove_dir_all(&dir);
915 std::fs::create_dir_all(&dir).unwrap();
916 let db = make_db(&dir);
917 let path = session_path(&db, "ses_parent");
918 let mut t = OpenCodeTranscript::new(&path, SpanRetention::All);
919 t.refresh().unwrap();
920 let s = t.summary();
921 assert_eq!(s.session_id.as_deref(), Some("ses_parent"));
922 assert_eq!(s.model.as_deref(), Some("deepseek-v4-pro"));
923 assert_eq!(s.cwd.as_deref(), Some(Path::new("/tmp/proj")));
924 assert_eq!(s.harness_version.as_deref(), Some("1.18.15"));
925 assert_eq!(s.usage.input, 1300);
927 assert_eq!(s.usage.output, 300);
928 assert_eq!(s.usage.cache_read, 901000);
929 assert!((s.cost_usd - 0.30).abs() < 1e-9, "{}", s.cost_usd);
931 assert_eq!(s.unpriced_tokens, 0, "OpenCode prices its own session");
932 assert_eq!(
933 s.price_source,
934 Some(crate::model::PriceSource::Harness),
935 "the cost came from OpenCode, so no price table may be named as its source"
936 );
937 assert_eq!(s.turns, 3, "two assistant turns in the parent, one in the subagent");
938 assert_eq!(s.subagent_turns, 1);
939 assert!(s.folds_child_usage, "a child's usage is inside this summary, so the share is worth breaking out");
940 assert_eq!(s.tool_calls, 3);
941 let tools: Vec<_> = s.spans.iter().filter(|sp| sp.kind == SpanKind::Tool).collect();
942 assert_eq!(tools.len(), 3);
943 assert_eq!(tools[0].name, "read");
944 assert_eq!(tools[0].duration_ms, Some(300));
945 let bash = tools.iter().find(|sp| sp.name == "bash").unwrap();
946 assert!(bash.error);
947 let grep = tools.iter().find(|sp| sp.name == "grep").unwrap();
948 assert!(grep.sidechain, "the subagent's tool call is a sidechain");
949
950 let inf: Vec<_> = s.spans.iter().filter(|sp| sp.kind == SpanKind::Inference).collect();
952 assert_eq!(inf.len(), 3, "two in the parent, one in the subagent");
953 let parent_inf: Vec<_> = inf.iter().filter(|sp| !sp.sidechain).collect();
954 assert_eq!(parent_inf[0].duration_ms, Some(90), "110 to 200");
955 assert_eq!(parent_inf[1].duration_ms, Some(90), "210 to 300");
956 assert_eq!(inf.iter().find(|sp| sp.sidechain).unwrap().duration_ms, Some(30), "subagent inference 150 to 180");
957 let turns: Vec<_> = s.spans.iter().filter(|sp| sp.kind == SpanKind::Turn).collect();
959 assert_eq!(turns.len(), 2, "one parent turn, one subagent turn");
960 let parent_turn = turns.iter().find(|sp| !sp.sidechain).unwrap();
961 assert_eq!(parent_turn.duration_ms, Some(200), "user at 100 to the last reply completing at 300");
962 assert!(turns.iter().any(|sp| sp.sidechain), "the subagent has its own turn");
963 assert_eq!(s.activity, Activity::Waiting, "the newest message is a completed reply");
964 assert!(!s.health.fields_unrecognised());
965 let _ = std::fs::remove_dir_all(&dir);
966 }
967
968 #[test]
969 fn recent_sessions_lists_only_top_level_and_attributes_by_directory() {
970 let dir = std::env::temp_dir().join(format!("agent-top-oc-attr-{}", std::process::id()));
971 let _ = std::fs::remove_dir_all(&dir);
972 std::fs::create_dir_all(&dir).unwrap();
973 let db = make_db(&dir);
974 let found = recent_sessions(&db, UNIX_EPOCH);
975 assert_eq!(found.len(), 1, "only the parent, not the subagent");
976 assert_eq!(found[0].id, "ses_parent");
977 assert_eq!(found[0].directory, PathBuf::from("/tmp/proj"));
978
979 let adapter = OpenCodeAdapter { db: Some(db.clone()), recent: found };
980 let ctx = AttributeContext {
981 cwd: Some(Path::new("/tmp/proj")),
982 proc_start: UNIX_EPOCH + Duration::from_secs(3),
983 now: SystemTime::now(),
984 attached: &HashSet::new(),
985 activity_timeout: Duration::from_secs(900),
986 };
987 let root = ProcNode {
988 pid: 1,
989 ppid: None,
990 name: "opencode".into(),
991 cmdline: "opencode".into(),
992 kind: crate::model::ProcKind::Agent,
993 harness: Some(Harness::OpenCode),
994 cpu_percent: 0.0,
995 rss_bytes: 0,
996 age_secs: 0,
997 cwd: None,
998 children: Vec::new(),
999 };
1000 let (paths, attribution) = adapter.attribute(&root, None, &ctx);
1001 assert_eq!(paths, vec![session_path(&db, "ses_parent")]);
1002 assert_eq!(attribution, Attribution::CwdHeuristic);
1003 let ctx2 = AttributeContext { cwd: Some(Path::new("/tmp/other")), ..ctx };
1005 assert!(adapter.attribute(&root, None, &ctx2).0.is_empty());
1006 let _ = std::fs::remove_dir_all(&dir);
1007 }
1008
1009 #[test]
1015 fn counts_mcp_calls_per_configured_server() {
1016 let dir = std::env::temp_dir().join(format!("agent-top-oc-mcp-{}", std::process::id()));
1017 let _ = std::fs::remove_dir_all(&dir);
1018 let project = dir.join("repo/sub");
1019 std::fs::create_dir_all(&project).unwrap();
1020 std::fs::create_dir_all(dir.join("repo/.git")).unwrap();
1021 std::fs::write(
1022 dir.join("repo/opencode.jsonc"),
1023 r#"{
1024 // servers for this repo
1025 "mcp": {
1026 "scratch_fs": { "type": "local", "command": ["npx", "-y", "@modelcontextprotocol/server-filesystem", "/tmp/x"] },
1027 "scratch": { "type": "remote", "url": "https://example.invalid/mcp" },
1028 "my.api": { "type": "remote", "url": "https://example.invalid/api", },
1029 },
1030}"#,
1031 )
1032 .unwrap();
1033 std::fs::write(dir.join("opencode.json"), r#"{"mcp":{"outside":{}}}"#).unwrap();
1035 let db = make_db(&dir);
1036 let conn = Connection::open(&db).unwrap();
1037 conn.execute("UPDATE session SET directory = ?1", [project.to_string_lossy()]).unwrap();
1038 let tool = |tool: &str, call: &str, status: &str, start: i64, end: i64| {
1039 format!(
1040 "{{\"type\":\"tool\",\"tool\":\"{tool}\",\"callID\":\"{call}\",\"state\":{{\"status\":\"{status}\",\"time\":{{\"start\":{start},\"end\":{end}}}}}}}"
1041 )
1042 };
1043 for (i, sid, data) in [
1044 (
1045 10,
1046 "ses_parent",
1047 tool("scratch_fs_list_directory", "call_00_7aG02jznvnkRRCWvUZzh2824", "completed", 1789297396598, 1789297396604),
1048 ),
1049 (
1050 11,
1051 "ses_parent",
1052 tool("scratch_fs_read_text_file", "call_00_1pkIe0AHF5h5X8VTC0EH3246", "error", 1789297397807, 1789297397810),
1053 ),
1054 (12, "ses_child", tool("scratch_search", "call_00_child", "completed", 1789297398000, 1789297398100)),
1055 (13, "ses_parent", tool("my_api_get", "call_00_api", "completed", 1789297399000, 1789297399050)),
1056 (14, "ses_parent", tool("outside_thing", "call_00_out", "completed", 1789297399100, 1789297399150)),
1057 ] {
1058 conn.execute("INSERT INTO part VALUES (?1, 'm', ?2, ?3, ?4)", rusqlite::params![format!("p{i}"), sid, 1000 + i as i64, data])
1059 .unwrap();
1060 }
1061 drop(conn);
1062
1063 let roots = ConfigRoots { global: Some(dir.join("no-global")), home: Some(dir.join("no-home")) };
1064 assert_eq!(mcp_server_names(&project, &roots), vec!["my.api".to_string(), "scratch".into(), "scratch_fs".into()]);
1065
1066 let mut t = OpenCodeTranscript::new(&session_path(&db, "ses_parent"), SpanRetention::Recent).with_config(roots);
1067 t.refresh().unwrap();
1068 let s = t.summary();
1069 assert_eq!(s.tool_calls, 8, "every tool part is still a tool call");
1070 assert_eq!(s.mcp.len(), 3, "{:?}", s.mcp);
1071 let fs = &s.mcp["scratch_fs"];
1072 assert_eq!((fs.calls, fs.errors), (2, 1), "the longer server name wins over its prefix");
1073 assert_eq!(fs.last_call, from_ms(1789297397810));
1074 assert_eq!(s.mcp["scratch"].calls, 1, "a subagent's call counts for the session");
1075 assert_eq!(s.mcp["my.api"].calls, 1, "matched through the sanitised name, filed under the configured one");
1076 assert!(!s.mcp.contains_key("outside"));
1077 let _ = std::fs::remove_dir_all(&dir);
1078 }
1079
1080 #[test]
1087 fn context_by_source_from_messages_in_order() {
1088 let dir = std::env::temp_dir().join(format!("agent-top-oc-ctx-{}", std::process::id()));
1089 let _ = std::fs::remove_dir_all(&dir);
1090 std::fs::create_dir_all(dir.join(".git")).unwrap();
1091 std::fs::write(dir.join("opencode.json"), r#"{"mcp":{"scratch_fs":{}}}"#).unwrap();
1092 let db = dir.join("opencode.db");
1093 let conn = Connection::open(&db).unwrap();
1094 conn.execute_batch(
1095 "CREATE TABLE session (id TEXT PRIMARY KEY, project_id TEXT, parent_id TEXT, directory TEXT, agent TEXT, \
1096 model TEXT, cost REAL DEFAULT 0, tokens_input INTEGER DEFAULT 0, tokens_output INTEGER DEFAULT 0, \
1097 tokens_reasoning INTEGER DEFAULT 0, tokens_cache_read INTEGER DEFAULT 0, tokens_cache_write INTEGER DEFAULT 0, \
1098 time_created INTEGER, time_updated INTEGER, version TEXT);
1099 CREATE TABLE message (id TEXT PRIMARY KEY, session_id TEXT, time_created INTEGER, data TEXT);
1100 CREATE TABLE part (id TEXT PRIMARY KEY, message_id TEXT, session_id TEXT, time_created INTEGER, data TEXT);",
1101 )
1102 .unwrap();
1103 let model = "claude-opus-5";
1104 let price = crate::pricing::table().lookup(model).expect("the test model is priced");
1105 conn.execute(
1106 "INSERT INTO session VALUES ('s', 'p', NULL, ?1, 'build', ?2, 0, 0, 0, 0, 0, 0, 1, 9, '1.18.15')",
1107 rusqlite::params![dir.to_string_lossy(), format!("{{\"id\":\"{model}\"}}")],
1108 )
1109 .unwrap();
1110 conn.execute(
1111 "INSERT INTO session VALUES ('sub', 'p', 's', ?1, 'explore', NULL, 0, 0, 0, 0, 0, 0, 2, 8, '1.18.15')",
1112 [dir.to_string_lossy()],
1113 )
1114 .unwrap();
1115 struct Reply {
1116 id: &'static str,
1117 sid: &'static str,
1118 created: i64,
1119 input: u64,
1120 output: u64,
1121 reasoning: u64,
1122 read: u64,
1123 cost: f64,
1124 summary: bool,
1125 }
1126 let reply = |id, sid, created, input, output, reasoning, read, cost, summary| Reply {
1127 id,
1128 sid,
1129 created,
1130 input,
1131 output,
1132 reasoning,
1133 read,
1134 cost,
1135 summary,
1136 };
1137 let replies = [
1138 reply("m1", "s", 10, 1000, 40, 10, 0, 0.02, false),
1139 reply("m2", "s", 20, 400, 30, 0, 1050, 0.01, false),
1140 reply("m3", "s", 30, 300, 20, 0, 1480, 0.01, false),
1141 reply("m4", "s", 40, 0, 0, 0, 0, 0.0, true),
1142 reply("m5", "s", 50, 200, 10, 0, 0, 0.005, false),
1143 reply("k1", "sub", 25, 500, 5, 0, 0, 0.004, false),
1144 ];
1145 for r in &replies {
1146 let mode = if r.summary { "compaction" } else { "build" };
1147 let data = format!(
1148 "{{\"role\":\"assistant\",\"mode\":\"{mode}\",\"summary\":{},\"modelID\":\"{model}\",\"cost\":{},\
1149 \"tokens\":{{\"input\":{},\"output\":{},\"reasoning\":{},\"cache\":{{\"read\":{},\"write\":0}}}},\
1150 \"time\":{{\"created\":{},\"completed\":{}}}}}",
1151 r.summary,
1152 r.cost,
1153 r.input,
1154 r.output,
1155 r.reasoning,
1156 r.read,
1157 r.created,
1158 r.created + 5
1159 );
1160 conn.execute("INSERT INTO message VALUES (?1, ?2, ?3, ?4)", rusqlite::params![r.id, r.sid, r.created, data]).unwrap();
1161 }
1162 let tool = |id: &str, msg: &str, sid: &str, name: &str, call: &str| {
1163 let data = format!(
1164 "{{\"type\":\"tool\",\"tool\":\"{name}\",\"callID\":\"{call}\",\"state\":{{\"status\":\"completed\",\"time\":{{\"start\":1,\"end\":2}}}}}}"
1165 );
1166 conn.execute("INSERT INTO part VALUES (?1, ?2, ?3, 1, ?4)", rusqlite::params![id, msg, sid, data]).unwrap();
1167 };
1168 tool("p1", "m1", "s", "read", "c1");
1169 tool("p2", "m2", "s", "scratch_fs_list_directory", "c2");
1170 tool("p3", "m3", "s", "bash", "c3");
1171 drop(conn);
1172
1173 let roots = ConfigRoots { global: Some(dir.join("no-global")), home: Some(dir.join("no-home")) };
1174 let mut t = OpenCodeTranscript::new(&session_path(&db, "s"), SpanRetention::Recent).with_config(roots);
1175 t.refresh().unwrap();
1176 let rows: std::collections::HashMap<String, crate::model::ContextSource> =
1177 t.summary().context.sources().into_iter().map(|c| (c.name.clone(), c)).collect();
1178
1179 assert_eq!((rows["read"].calls, rows["read"].tokens), (1, 400), "{rows:?}");
1181 assert_eq!((rows["scratch_fs"].calls, rows["scratch_fs"].tokens, rows["scratch_fs"].origin), (1, 300, ContextOrigin::Mcp));
1183 assert_eq!((rows["bash"].calls, rows["bash"].tokens), (1, 200), "{rows:?}");
1187 assert_eq!(rows[ContextLedger::OTHER].tokens, 1580, "{rows:?}");
1189
1190 let prompt_side = |input: u64, output: u64, read: u64, total: f64| {
1192 let usage = TokenUsage { input, output, cache_read: read, ..Default::default() };
1193 let b = scaled(price.breakdown(&usage), total);
1194 b.input + b.cache_read + b.cache_write_5m + b.cache_write_1h
1195 };
1196 let expected: f64 =
1197 replies.iter().filter(|r| r.input + r.read > 0).map(|r| prompt_side(r.input, r.output + r.reasoning, r.read, r.cost)).sum();
1198 let got: f64 = rows.values().map(|c| c.cost_usd).sum();
1199 assert!((expected - got).abs() < 1e-9, "rows {got} vs prompt-side {expected}");
1200 let _ = std::fs::remove_dir_all(&dir);
1201 }
1202
1203 #[test]
1204 fn no_configured_server_means_no_mcp_guess() {
1205 assert_eq!(mcp_server_of("scratch_fs_list_directory", &[]), None);
1206 let servers = vec!["github".to_string(), "github_enterprise".to_string()];
1207 assert_eq!(mcp_server_of("github_enterprise_search", &servers), Some("github_enterprise"));
1208 assert_eq!(mcp_server_of("github_search", &servers), Some("github"));
1209 assert_eq!(mcp_server_of("github_", &servers), None, "a server name with no tool after it");
1210 assert_eq!(mcp_server_of("githubsearch", &servers), None);
1211 assert_eq!(mcp_server_of("todowrite", &servers), None);
1212 }
1213
1214 #[test]
1215 fn jsonc_comments_and_trailing_commas_are_stripped_but_strings_are_kept() {
1216 let text = r#"{ "a": "http://x//y", /* block, */ "b": [1, 2,], // line
1217 "c": "a \"quoted, }\" value", }"#;
1218 let v: serde_json::Value = serde_json::from_str(&strip_jsonc(text)).unwrap();
1219 assert_eq!(v["a"], "http://x//y");
1220 assert_eq!(v["b"], serde_json::json!([1, 2]));
1221 assert_eq!(v["c"], "a \"quoted, }\" value");
1222 }
1223
1224 #[test]
1225 fn model_id_is_pulled_from_the_json_blob() {
1226 assert_eq!(model_id(r#"{"id":"deepseek-v4-pro","providerID":"deepseek"}"#), Some("deepseek-v4-pro".into()));
1227 assert_eq!(model_id("claude-fable-5-1"), Some("claude-fable-5-1".into()));
1228 assert_eq!(model_id(""), None);
1229 }
1230}