Skip to main content

tokenburn_core/features/opencode/
ops.rs

1use 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
14/// OpenCode data directories: `$OPENCODE_DATA_DIR`, else
15/// `$XDG_DATA_HOME/opencode` and `~/.local/share/opencode` (OpenCode uses the XDG
16/// layout on every OS), plus the platform data dir (`%LOCALAPPDATA%` on Windows).
17pub 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
60/// Read one data directory: SQLite databases first (the current format), then the
61/// legacy `storage/message/**.json` files, de-duplicated by message id.
62pub fn collect_opencode_from(root: &Path, start: DateTime<Utc>) -> Result<Vec<Row>> {
63    Ok(collect_dir(root, start, &mut HashSet::new()))
64}
65
66/// `(message id, row)` pairs per database, valid while the database (and its
67/// `-wal`) are unchanged for the same window start.
68static DB_CACHE: LazyLock<FileCache<Vec<(String, Row)>>> = LazyLock::new(FileCache::default);
69/// The row of one legacy message file, if it is an assistant message with usage.
70static 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        // `seen` also de-duplicates against the legacy files below.
87        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
128/// `opencode.db` and `opencode-*.db`.
129fn 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
182/// One assistant message's JSON (`data` column or a legacy message file).
183fn 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    // OpenCode reports reasoning separately; it is billed as output.
194    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(); // same id as the db row
287        std::fs::write(legacy.join("m2.json"), MSG).unwrap(); // genuinely new
288        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}