tokenburn_core/features/opencode/
ops.rs1use std::collections::HashSet;
2use std::path::{Path, PathBuf};
3use std::sync::LazyLock;
4
5use anyhow::Result;
6use chrono::{DateTime, TimeZone, Utc};
7use rusqlite::{Connection, OpenFlags};
8use serde_json::Value;
9
10use crate::features::usage::{Row, Tool};
11use crate::utils::cache::{FileCache, Sig};
12use crate::utils::files::{dir_basename, env_path, modified_since, num, walk_ext};
13
14pub fn roots() -> Vec<PathBuf> {
18 resolve_roots(
19 env_path("OPENCODE_DATA_DIR"),
20 env_path("XDG_DATA_HOME"),
21 dirs::home_dir(),
22 dirs::data_local_dir(),
23 )
24}
25
26pub(crate) fn resolve_roots(
27 over: Option<PathBuf>,
28 xdg_data: Option<PathBuf>,
29 home: Option<PathBuf>,
30 data_local: Option<PathBuf>,
31) -> Vec<PathBuf> {
32 if let Some(o) = over {
33 return vec![o];
34 }
35 let mut out = Vec::new();
36 for p in [
37 xdg_data.map(|d| d.join("opencode")),
38 home.map(|h| h.join(".local/share/opencode")),
39 data_local.map(|d| d.join("opencode")),
40 ]
41 .into_iter()
42 .flatten()
43 {
44 if !out.contains(&p) {
45 out.push(p);
46 }
47 }
48 out
49}
50
51pub fn collect_opencode(start: DateTime<Utc>) -> Result<Vec<Row>> {
52 let mut rows = Vec::new();
53 let mut seen = HashSet::new();
54 for root in roots().into_iter().filter(|r| r.is_dir()) {
55 rows.extend(collect_dir(&root, start, &mut seen));
56 }
57 Ok(rows)
58}
59
60pub fn collect_opencode_from(root: &Path, start: DateTime<Utc>) -> Result<Vec<Row>> {
63 Ok(collect_dir(root, start, &mut HashSet::new()))
64}
65
66static DB_CACHE: LazyLock<FileCache<Vec<(String, Row)>>> = LazyLock::new(FileCache::default);
69static FILE_CACHE: LazyLock<FileCache<Option<Row>>> = LazyLock::new(FileCache::default);
71
72fn collect_dir(root: &Path, start: DateTime<Utc>, seen: &mut HashSet<String>) -> Vec<Row> {
73 let mut rows = Vec::new();
74 for db in databases(root) {
75 let wal = PathBuf::from(format!("{}-wal", db.display()));
76 let Some(sig) = Sig::of_all(&[db.clone(), wal], start.timestamp_millis()) else {
77 continue;
78 };
79 let parsed = DB_CACHE.get_or_parse(&db, sig, || match read_db(&db, start) {
80 Ok(r) => r,
81 Err(e) => {
82 tracing::warn!("opencode: {}: {e}", db.display());
83 Vec::new()
84 }
85 });
86 rows.extend(
88 parsed
89 .iter()
90 .filter(|(id, _)| seen.insert(id.clone()))
91 .map(|(_, r)| r.clone()),
92 );
93 }
94
95 let messages = root.join("storage/message");
96 let files = walk_ext(&messages, &["json"]);
97 let live: HashSet<PathBuf> = files.iter().cloned().collect();
98 for file in files {
99 if !modified_since(&file, start) {
100 continue;
101 }
102 let id = file
103 .file_stem()
104 .map(|s| s.to_string_lossy().into_owned())
105 .unwrap_or_default();
106 let Some(sig) = Sig::of(&file, 0) else {
107 continue;
108 };
109 let parsed = FILE_CACHE.get_or_parse(&file, sig, || {
110 let session = file
111 .parent()
112 .and_then(|p| p.file_name())
113 .map(|n| n.to_string_lossy().into_owned())
114 .unwrap_or_default();
115 std::fs::read_to_string(&file)
116 .ok()
117 .and_then(|text| parse_message(&text, &session, None))
118 });
119 if let Some(row) = parsed.as_ref().clone().filter(|_| seen.insert(id)) {
120 rows.push(row);
121 }
122 }
123 FILE_CACHE.prune_under(&messages, &live);
124 rows.retain(|r| r.ts >= start);
125 rows
126}
127
128fn databases(root: &Path) -> Vec<PathBuf> {
130 let Ok(rd) = std::fs::read_dir(root) else {
131 return Vec::new();
132 };
133 let mut dbs: Vec<PathBuf> = rd
134 .flatten()
135 .map(|e| e.path())
136 .filter(|p| {
137 p.file_name().and_then(|n| n.to_str()).is_some_and(|n| {
138 n == "opencode.db" || (n.starts_with("opencode-") && n.ends_with(".db"))
139 })
140 })
141 .collect();
142 dbs.sort();
143 dbs
144}
145
146fn read_db(db: &Path, start: DateTime<Utc>) -> Result<Vec<(String, Row)>> {
147 let conn = Connection::open_with_flags(
148 db,
149 OpenFlags::SQLITE_OPEN_READ_ONLY | OpenFlags::SQLITE_OPEN_NO_MUTEX,
150 )?;
151 let has_message: bool = conn
152 .query_row(
153 "SELECT 1 FROM sqlite_master WHERE type='table' AND name='message'",
154 [],
155 |_| Ok(true),
156 )
157 .unwrap_or(false);
158 if !has_message {
159 return Ok(Vec::new());
160 }
161 let mut stmt = conn.prepare(
162 "SELECT id, session_id, time_created, data FROM message WHERE time_created >= ?1",
163 )?;
164 let mapped = stmt.query_map([start.timestamp_millis()], |r| {
165 Ok((
166 r.get::<_, String>(0)?,
167 r.get::<_, String>(1)?,
168 r.get::<_, i64>(2)?,
169 r.get::<_, String>(3)?,
170 ))
171 })?;
172 let mut rows = Vec::new();
173 for item in mapped.flatten() {
174 let (id, session, created, data) = item;
175 if let Some(row) = parse_message(&data, &session, Some(created)) {
176 rows.push((id, row));
177 }
178 }
179 Ok(rows)
180}
181
182fn parse_message(json: &str, session: &str, created_ms: Option<i64>) -> Option<Row> {
184 let v: Value = serde_json::from_str(json).ok()?;
185 if v.get("role")
186 .and_then(Value::as_str)
187 .is_some_and(|r| r != "assistant")
188 {
189 return None;
190 }
191 let t = v.get("tokens").filter(|t| t.is_object())?;
192 let cache = t.get("cache");
193 let output = num(t.get("output")) + num(t.get("reasoning"));
195 let input = num(t.get("input"));
196 let cache_read = num(cache.and_then(|c| c.get("read")));
197 let cache_write = num(cache.and_then(|c| c.get("write")));
198 if input + output + cache_read + cache_write == 0 {
199 return None;
200 }
201 let ms = created_ms.or_else(|| v.pointer("/time/created").and_then(Value::as_i64))?;
202 let ts = Utc.timestamp_millis_opt(ms).single()?;
203 let project = v
204 .pointer("/path/cwd")
205 .and_then(Value::as_str)
206 .and_then(dir_basename)
207 .or_else(|| v.get("modelID").and_then(Value::as_str).map(str::to_string))
208 .unwrap_or_else(|| "unknown".into());
209 Some(Row {
210 tool: Tool::OpenCode,
211 project,
212 id: session.to_string(),
213 ts,
214 input,
215 output,
216 cache_read,
217 cache_write,
218 cost: v.get("cost").and_then(Value::as_f64).unwrap_or(0.0),
219 })
220}
221
222#[cfg(test)]
223mod tests {
224 use super::*;
225
226 const MSG: &str = r#"{"role":"assistant","modelID":"claude-x","providerID":"anthropic",
227 "time":{"created":1777629600000},"path":{"cwd":"/home/me/app"},"cost":0.25,
228 "tokens":{"input":100,"output":20,"reasoning":5,"cache":{"read":40,"write":8}}}"#;
229
230 #[test]
231 fn parses_tokens_cost_project_and_time() {
232 let r = parse_message(MSG, "ses_1", None).unwrap();
233 assert_eq!(
234 (r.input, r.output, r.cache_read, r.cache_write),
235 (100, 25, 40, 8)
236 );
237 assert!((r.cost - 0.25).abs() < 1e-9);
238 assert_eq!(r.project, "app");
239 assert_eq!(r.ts.timestamp_millis(), 1777629600000);
240 assert_eq!(r.tool, Tool::OpenCode);
241 }
242
243 #[test]
244 fn user_messages_and_empty_usage_are_skipped() {
245 assert!(parse_message(r#"{"role":"user","tokens":{"input":5}}"#, "s", Some(1)).is_none());
246 assert!(parse_message(
247 r#"{"role":"assistant","tokens":{"input":0,"output":0}}"#,
248 "s",
249 Some(1)
250 )
251 .is_none());
252 assert!(parse_message("garbage", "s", Some(1)).is_none());
253 }
254
255 fn make_db(dir: &Path) {
256 let c = Connection::open(dir.join("opencode.db")).unwrap();
257 c.execute(
258 "CREATE TABLE message (id TEXT, session_id TEXT, time_created INTEGER, data TEXT)",
259 [],
260 )
261 .unwrap();
262 c.execute(
263 "INSERT INTO message VALUES ('m1','s1',1777629600000,?1)",
264 [MSG],
265 )
266 .unwrap();
267 c.execute("INSERT INTO message VALUES ('m0','s1',1000,?1)", [MSG])
268 .unwrap();
269 }
270
271 #[test]
272 fn reads_the_sqlite_database_and_honours_the_start() {
273 let d = tempfile::tempdir().unwrap();
274 make_db(d.path());
275 let start = Utc.timestamp_millis_opt(1_700_000_000_000).unwrap();
276 let rows = collect_opencode_from(d.path(), start).unwrap();
277 assert_eq!(rows.len(), 1, "the 1970 message is before `start`");
278 }
279
280 #[test]
281 fn legacy_json_is_read_and_deduplicated_behind_the_database() {
282 let d = tempfile::tempdir().unwrap();
283 make_db(d.path());
284 let legacy = d.path().join("storage/message/s1");
285 std::fs::create_dir_all(&legacy).unwrap();
286 std::fs::write(legacy.join("m1.json"), MSG).unwrap(); std::fs::write(legacy.join("m2.json"), MSG).unwrap(); let rows = collect_opencode_from(d.path(), DateTime::<Utc>::UNIX_EPOCH).unwrap();
289 assert_eq!(
290 rows.len(),
291 3,
292 "m1 (db) + m0 (db) + m2 (legacy); the m1 legacy copy is dropped"
293 );
294 }
295
296 #[test]
297 fn a_database_without_a_message_table_is_not_an_error() {
298 let d = tempfile::tempdir().unwrap();
299 Connection::open(d.path().join("opencode.db"))
300 .unwrap()
301 .execute("CREATE TABLE other (x)", [])
302 .unwrap();
303 assert!(collect_opencode_from(d.path(), DateTime::<Utc>::UNIX_EPOCH)
304 .unwrap()
305 .is_empty());
306 }
307
308 #[test]
309 fn roots_prefer_the_override_then_xdg_and_home() {
310 assert_eq!(
311 resolve_roots(Some("/o".into()), None, None, None),
312 [PathBuf::from("/o")]
313 );
314 let r = resolve_roots(
315 None,
316 Some("/x".into()),
317 Some("/h".into()),
318 Some("/h/.local/share".into()),
319 );
320 assert_eq!(
321 r,
322 [
323 PathBuf::from("/x/opencode"),
324 PathBuf::from("/h/.local/share/opencode")
325 ]
326 );
327 }
328}