use std::collections::HashSet;
use std::path::{Path, PathBuf};
use std::sync::LazyLock;
use anyhow::Result;
use chrono::{DateTime, Utc};
use serde::Deserialize;
use crate::features::usage::{Row, Tool};
use crate::utils::cache::{FileCache, Sig};
use crate::utils::files::{modified_since, num_f, read_lines, walk_ext};
use crate::utils::time::parse_ts;
pub fn sessions_dir() -> Option<PathBuf> {
resolve_sessions_dir(|k| std::env::var_os(k), dirs::home_dir())
}
pub(crate) fn resolve_sessions_dir(
env: impl Fn(&str) -> Option<std::ffi::OsString>,
home: Option<PathBuf>,
) -> Option<PathBuf> {
let non_empty = |k: &str| env(k).filter(|v| !v.is_empty()).map(PathBuf::from);
non_empty("TOKENBURN_PI_SESSIONS")
.or_else(|| non_empty("PI_CODING_AGENT_SESSION_DIR"))
.or_else(|| non_empty("PI_CODING_AGENT_DIR").map(|d| d.join("sessions")))
.or_else(|| home.map(|h| h.join(".pi").join("agent").join("sessions")))
}
pub fn collect_pi(start: DateTime<Utc>) -> Result<Vec<Row>> {
match sessions_dir() {
Some(dir) if dir.is_dir() => collect_pi_from(&dir, start),
_ => Ok(Vec::new()),
}
}
static CACHE: LazyLock<FileCache<Vec<Row>>> = LazyLock::new(FileCache::default);
pub fn collect_pi_from(root: &Path, start: DateTime<Utc>) -> Result<Vec<Row>> {
let files = walk_ext(root, &["jsonl"]);
let live: HashSet<PathBuf> = files.iter().cloned().collect();
let mut rows = Vec::new();
for file in files {
if !modified_since(&file, start) {
continue;
}
let Some(sig) = Sig::of(&file, 0) else {
continue;
};
let parsed = CACHE.get_or_parse(&file, sig, || parse_file(&file));
rows.extend(parsed.iter().filter(|r| r.ts >= start).cloned());
}
CACHE.prune_under(root, &live);
Ok(rows)
}
fn parse_file(file: &Path) -> Vec<Row> {
let Some(lines) = read_lines(file) else {
tracing::warn!("pi: cannot read {}", file.display());
return Vec::new();
};
let project = file
.parent()
.and_then(|p| p.file_name())
.map(|n| n.to_string_lossy().trim_matches('-').to_string())
.unwrap_or_default();
let id = file
.file_name()
.map(|n| n.to_string_lossy().into_owned())
.unwrap_or_default();
lines
.filter_map(|l| parse_line(&l, &project, &id))
.collect()
}
#[derive(Deserialize)]
struct PiLine {
timestamp: Option<String>,
message: Option<PiMessage>,
}
#[derive(Deserialize)]
struct PiMessage {
usage: Option<PiUsage>,
}
#[derive(Deserialize)]
#[serde(rename_all = "camelCase")]
struct PiUsage {
input: Option<f64>,
output: Option<f64>,
cache_read: Option<f64>,
cache_write: Option<f64>,
cost: Option<PiCost>,
}
#[derive(Deserialize)]
struct PiCost {
total: Option<f64>,
}
pub(crate) fn parse_line(line: &str, project: &str, id: &str) -> Option<Row> {
if !line.contains("\"usage\"") {
return None;
}
let rec: PiLine = serde_json::from_str(line).ok()?;
let usage = rec.message?.usage?;
let ts = parse_ts(&rec.timestamp?)?;
Some(Row {
tool: Tool::Pi,
project: project.to_string(),
id: id.to_string(),
ts,
input: num_f(usage.input),
output: num_f(usage.output),
cache_read: num_f(usage.cache_read),
cache_write: num_f(usage.cache_write),
cost: usage.cost.and_then(|c| c.total).unwrap_or(0.0),
})
}
#[cfg(test)]
mod tests {
use super::*;
use std::fs;
static SERIAL: std::sync::Mutex<()> = std::sync::Mutex::new(());
const LINE: &str = r#"{"timestamp":"2026-05-01T10:00:00Z","message":{"usage":{"input":10,"output":5,"cacheRead":3,"cacheWrite":2,"cost":{"total":0.25}}}}"#;
#[test]
fn parses_usage_line() {
let r = parse_line(LINE, "proj", "a.jsonl").unwrap();
assert_eq!(
(r.input, r.output, r.cache_read, r.cache_write),
(10, 5, 3, 2)
);
assert!((r.cost - 0.25).abs() < 1e-9);
assert_eq!(r.tool, Tool::Pi);
}
#[test]
fn skips_lines_without_usage_or_timestamp() {
assert!(parse_line(r#"{"message":{"role":"user"}}"#, "p", "i").is_none());
assert!(parse_line(r#"{"message":{"usage":{"input":1}}}"#, "p", "i").is_none());
assert!(parse_line("not json", "p", "i").is_none());
}
fn env_of<'a>(
pairs: &'a [(&'a str, &'a str)],
) -> impl Fn(&str) -> Option<std::ffi::OsString> + 'a {
move |k| {
pairs
.iter()
.find(|(n, _)| *n == k)
.map(|(_, v)| std::ffi::OsString::from(*v))
}
}
#[test]
fn sessions_dir_defaults_to_home() {
let d = resolve_sessions_dir(env_of(&[]), Some(PathBuf::from("/home/u"))).unwrap();
assert_eq!(d, PathBuf::from("/home/u/.pi/agent/sessions"));
}
#[test]
fn sessions_dir_honours_pi_overrides_in_order() {
let home = Some(PathBuf::from("/home/u"));
let d = resolve_sessions_dir(env_of(&[("PI_CODING_AGENT_DIR", "/agent")]), home.clone());
assert_eq!(d.unwrap(), PathBuf::from("/agent/sessions"));
let d = resolve_sessions_dir(
env_of(&[
("PI_CODING_AGENT_DIR", "/agent"),
("PI_CODING_AGENT_SESSION_DIR", "/sess"),
]),
home.clone(),
);
assert_eq!(d.unwrap(), PathBuf::from("/sess"));
let d = resolve_sessions_dir(
env_of(&[
("PI_CODING_AGENT_SESSION_DIR", "/sess"),
("TOKENBURN_PI_SESSIONS", "/mine"),
]),
home,
);
assert_eq!(d.unwrap(), PathBuf::from("/mine"));
}
#[test]
fn empty_override_is_ignored() {
let d = resolve_sessions_dir(
env_of(&[("TOKENBURN_PI_SESSIONS", "")]),
Some(PathBuf::from("/h")),
);
assert_eq!(d.unwrap(), PathBuf::from("/h/.pi/agent/sessions"));
}
#[test]
fn unchanged_files_are_not_parsed_again_but_growing_ones_are() {
let _serial = SERIAL.lock().unwrap_or_else(|e| e.into_inner());
let dir = tempfile::tempdir().unwrap();
let proj = dir.path().join("--cache-proj--");
fs::create_dir_all(&proj).unwrap();
let f = proj.join("s.jsonl");
fs::write(&f, format!("{LINE}\n")).unwrap();
let start = "1970-01-01T00:00:00Z".parse().unwrap();
let (h0, m0) = CACHE.stats();
assert_eq!(collect_pi_from(dir.path(), start).unwrap().len(), 1);
assert_eq!(collect_pi_from(dir.path(), start).unwrap().len(), 1);
let (h1, m1) = CACHE.stats();
assert_eq!(
(h1 - h0, m1 - m0),
(1, 1),
"second refresh served from the cache"
);
fs::write(&f, format!("{LINE}\n{LINE}\n")).unwrap();
assert_eq!(collect_pi_from(dir.path(), start).unwrap().len(), 2);
let (_, m2) = CACHE.stats();
assert_eq!(m2 - m1, 1);
}
#[test]
fn deleted_files_leave_the_cache() {
let _serial = SERIAL.lock().unwrap_or_else(|e| e.into_inner());
let dir = tempfile::tempdir().unwrap();
let proj = dir.path().join("p");
fs::create_dir_all(&proj).unwrap();
let f = proj.join("s.jsonl");
fs::write(&f, format!("{LINE}\n")).unwrap();
let start = "1970-01-01T00:00:00Z".parse().unwrap();
assert_eq!(collect_pi_from(dir.path(), start).unwrap().len(), 1);
fs::remove_file(&f).unwrap();
assert!(collect_pi_from(dir.path(), start).unwrap().is_empty());
}
#[test]
fn collects_from_directory() {
let _serial = SERIAL.lock().unwrap_or_else(|e| e.into_inner());
let dir = tempfile::tempdir().unwrap();
let proj = dir.path().join("--home-me-proj--");
fs::create_dir_all(&proj).unwrap();
fs::write(proj.join("s.jsonl"), format!("{LINE}\ngarbage\n{LINE}\n")).unwrap();
let start = "1970-01-01T00:00:00Z".parse().unwrap();
let rows = collect_pi_from(dir.path(), start).unwrap();
assert_eq!(rows.len(), 2);
assert_eq!(rows[0].project, "home-me-proj");
}
}