Initial commit

This commit is contained in:
xggz
2026-03-06 22:56:13 +08:00
commit 54d1097b41
273 changed files with 92457 additions and 0 deletions
+983
View File
@@ -0,0 +1,983 @@
use std::fs;
use std::io::{BufRead, BufReader};
use std::path::PathBuf;
use std::sync::OnceLock;
use chrono::{DateTime, Utc};
use regex::Regex;
use crate::models::*;
use crate::parsers::{folder_name_from_path, truncate_str, AgentParser, ParseError};
/// Regex that matches Claude Code system-injected XML tags and their content.
/// These tags are internal metadata and should not be displayed to users.
/// Note: Rust regex doesn't support backreferences, so each tag is listed explicitly.
fn system_tag_regex() -> &'static Regex {
static RE: OnceLock<Regex> = OnceLock::new();
RE.get_or_init(|| {
Regex::new(concat!(
r"(?s)",
r"<system-reminder>.*?</system-reminder>",
r"|<local-command-caveat>.*?</local-command-caveat>",
r"|<command-name>.*?</command-name>",
r"|<command-message>.*?</command-message>",
r"|<command-args>.*?</command-args>",
r"|<local-command-stdout>.*?</local-command-stdout>",
r"|<user-prompt-submit-hook>.*?</user-prompt-submit-hook>",
))
.unwrap()
})
}
/// Regex that matches an optional model capacity suffix like `[1M]` / `[500k]`.
fn model_capacity_suffix_regex() -> &'static Regex {
static RE: OnceLock<Regex> = OnceLock::new();
RE.get_or_init(|| {
Regex::new(r"(?i)\[\s*([0-9]+(?:\.[0-9]+)?)\s*([km])\s*\]\s*$")
.expect("valid model capacity regex")
})
}
/// Strip system-injected XML tags from text content.
/// Returns None if the text becomes empty after stripping.
fn strip_system_tags(text: &str) -> Option<String> {
let cleaned = system_tag_regex().replace_all(text, "");
let trimmed = cleaned.trim();
if trimmed.is_empty() {
None
} else {
Some(trimmed.to_string())
}
}
/// Check if a JSONL entry is a system meta message (isMeta: true).
fn is_meta_message(value: &serde_json::Value) -> bool {
value
.get("isMeta")
.and_then(|v| v.as_bool())
.unwrap_or(false)
}
fn parse_model_capacity_suffix(model: &str) -> Option<u64> {
let captures = model_capacity_suffix_regex().captures(model.trim())?;
let value = captures.get(1)?.as_str().parse::<f64>().ok()?;
if !value.is_finite() || value <= 0.0 {
return None;
}
let unit = captures
.get(2)
.map(|m| m.as_str().to_ascii_lowercase())
.unwrap_or_default();
let multiplier = match unit.as_str() {
"m" => 1_000_000.0,
"k" => 1_000.0,
_ => return None,
};
Some((value * multiplier) as u64)
}
fn claude_context_window_max_tokens_for_model(model: Option<&str>) -> Option<u64> {
let model = model?.trim();
if model.is_empty() {
return None;
}
// If user/model config contains an explicit capacity suffix, prefer it.
if let Some(suffixed_limit) = parse_model_capacity_suffix(model) {
return Some(suffixed_limit);
}
// Claude models default to 200k when no explicit capacity is provided.
if model.to_ascii_lowercase().starts_with("claude") {
return Some(200_000);
}
None
}
fn claude_context_window_used_tokens_from_usage(usage: &TurnUsage) -> Option<u64> {
let used_tokens = usage
.input_tokens
.saturating_add(usage.cache_creation_input_tokens)
.saturating_add(usage.cache_read_input_tokens);
if used_tokens > 0 {
Some(used_tokens)
} else {
None
}
}
fn latest_claude_context_window_used_tokens(turns: &[MessageTurn]) -> Option<u64> {
turns.iter().rev().find_map(|turn| {
turn.usage
.as_ref()
.and_then(claude_context_window_used_tokens_from_usage)
})
}
fn merge_claude_context_window_stats(
stats: Option<SessionStats>,
used_tokens: Option<u64>,
max_tokens: Option<u64>,
) -> Option<SessionStats> {
if used_tokens.is_none() && max_tokens.is_none() {
return stats;
}
let usage_percent = match (used_tokens, max_tokens) {
(Some(used), Some(max)) if max > 0 => Some((used as f64 / max as f64) * 100.0),
_ => None,
};
match stats {
Some(mut s) => {
s.context_window_used_tokens = used_tokens;
s.context_window_max_tokens = max_tokens;
s.context_window_usage_percent = usage_percent;
Some(s)
}
None => Some(SessionStats {
total_usage: None,
total_tokens: None,
total_duration_ms: 0,
context_window_used_tokens: used_tokens,
context_window_max_tokens: max_tokens,
context_window_usage_percent: usage_percent,
}),
}
}
pub struct ClaudeParser {
base_dir: PathBuf,
}
impl ClaudeParser {
pub fn new() -> Self {
let base_dir = resolve_claude_config_dir().join("projects");
Self { base_dir }
}
fn decode_folder_path(encoded: &str) -> String {
encoded.replace('-', "/")
}
fn parse_jsonl_summary(
&self,
path: &PathBuf,
) -> Result<Option<ConversationSummary>, ParseError> {
let file = fs::File::open(path)?;
let reader = BufReader::new(file);
let mut conversation_id: Option<String> = None;
let mut cwd: Option<String> = None;
let mut git_branch: Option<String> = None;
let mut model: Option<String> = None;
let mut title: Option<String> = None;
let mut first_timestamp: Option<DateTime<Utc>> = None;
let mut last_timestamp: Option<DateTime<Utc>> = None;
let mut message_count: u32 = 0;
for line in reader.lines() {
let line = match line {
Ok(l) => l,
Err(_) => continue,
};
if line.trim().is_empty() {
continue;
}
let value: serde_json::Value = match serde_json::from_str(&line) {
Ok(v) => v,
Err(_) => continue,
};
let msg_type = value.get("type").and_then(|t| t.as_str()).unwrap_or("");
// Skip non-conversation entries
if msg_type == "file-history-snapshot" || msg_type == "progress" {
continue;
}
// Skip system meta messages (e.g. local-command-caveat injections)
if is_meta_message(&value) {
continue;
}
if conversation_id.is_none() {
conversation_id = value
.get("sessionId")
.and_then(|s| s.as_str())
.map(|s| s.to_string());
}
if cwd.is_none() {
cwd = value
.get("cwd")
.and_then(|s| s.as_str())
.map(|s| s.to_string());
}
if git_branch.is_none() {
git_branch = value
.get("gitBranch")
.and_then(|s| s.as_str())
.map(|s| s.to_string());
}
if let Some(ts_str) = value.get("timestamp").and_then(|t| t.as_str()) {
if let Ok(ts) = ts_str.parse::<DateTime<Utc>>() {
if first_timestamp.is_none() {
first_timestamp = Some(ts);
}
last_timestamp = Some(ts);
}
}
if msg_type == "user" || msg_type == "assistant" {
message_count += 1;
// Extract model from assistant messages
if msg_type == "assistant" && model.is_none() {
model = value
.get("message")
.and_then(|m| m.get("model"))
.and_then(|m| m.as_str())
.map(|s| s.to_string());
}
// Extract title from first user message
if msg_type == "user" && title.is_none() {
title = extract_user_text(&value).map(|t| truncate_str(&t, 100));
}
}
}
let started_at = match first_timestamp {
Some(ts) => ts,
None => return Ok(None),
};
// Use filename (without .jsonl) as ID fallback
let id = conversation_id.unwrap_or_else(|| {
path.file_stem()
.unwrap_or_default()
.to_string_lossy()
.to_string()
});
let folder_path = cwd.clone();
let folder_name = folder_path.as_ref().map(|p| folder_name_from_path(p));
Ok(Some(ConversationSummary {
id,
agent_type: AgentType::ClaudeCode,
folder_path,
folder_name,
title,
started_at,
ended_at: last_timestamp,
message_count,
model,
git_branch,
}))
}
}
fn resolve_claude_config_dir() -> PathBuf {
resolve_claude_config_dir_from(
std::env::var_os("CLAUDE_CONFIG_DIR"),
dirs::home_dir(),
)
}
fn resolve_claude_config_dir_from(
claude_config_dir_env: Option<std::ffi::OsString>,
home_dir: Option<PathBuf>,
) -> PathBuf {
claude_config_dir_env
.filter(|value| !value.is_empty())
.map(PathBuf::from)
.unwrap_or_else(|| home_dir.unwrap_or_default().join(".claude"))
}
impl AgentParser for ClaudeParser {
fn list_conversations(&self) -> Result<Vec<ConversationSummary>, ParseError> {
let mut conversations = Vec::new();
if !self.base_dir.exists() {
return Ok(conversations);
}
let entries = fs::read_dir(&self.base_dir)?;
for entry in entries {
let entry = match entry {
Ok(e) => e,
Err(_) => continue,
};
let project_dir = entry.path();
if !project_dir.is_dir() {
continue;
}
let jsonl_files = fs::read_dir(&project_dir)?;
for file_entry in jsonl_files {
let file_entry = match file_entry {
Ok(e) => e,
Err(_) => continue,
};
let file_path = file_entry.path();
if file_path.extension().and_then(|e| e.to_str()) != Some("jsonl") {
continue;
}
match self.parse_jsonl_summary(&file_path) {
Ok(Some(mut summary)) => {
// If folder_path is still None, derive from directory name
if summary.folder_path.is_none() {
let dir_name = project_dir
.file_name()
.unwrap_or_default()
.to_string_lossy()
.to_string();
let decoded = Self::decode_folder_path(&dir_name);
summary.folder_path = Some(decoded.clone());
summary.folder_name = Some(folder_name_from_path(&decoded));
}
conversations.push(summary);
}
Ok(None) => continue,
Err(_) => continue,
}
}
}
conversations.sort_by(|a, b| b.started_at.cmp(&a.started_at));
Ok(conversations)
}
fn get_conversation(&self, conversation_id: &str) -> Result<ConversationDetail, ParseError> {
// Find the conversation file by searching all directories
if !self.base_dir.exists() {
return Err(ParseError::ConversationNotFound(
conversation_id.to_string(),
));
}
for entry in fs::read_dir(&self.base_dir)? {
let entry = match entry {
Ok(e) => e,
Err(_) => continue,
};
let project_dir = entry.path();
if !project_dir.is_dir() {
continue;
}
let file_path = project_dir.join(format!("{}.jsonl", conversation_id));
if file_path.exists() {
return self.parse_conversation_detail(&file_path, conversation_id);
}
}
Err(ParseError::ConversationNotFound(
conversation_id.to_string(),
))
}
}
impl ClaudeParser {
fn parse_conversation_detail(
&self,
path: &PathBuf,
conversation_id: &str,
) -> Result<ConversationDetail, ParseError> {
let file = fs::File::open(path)?;
let reader = BufReader::new(file);
let mut messages = Vec::new();
let mut cwd: Option<String> = None;
let mut git_branch: Option<String> = None;
let mut model: Option<String> = None;
let mut title: Option<String> = None;
let mut first_timestamp: Option<DateTime<Utc>> = None;
let mut last_timestamp: Option<DateTime<Utc>> = None;
for line in reader.lines() {
let line = match line {
Ok(l) => l,
Err(_) => continue,
};
if line.trim().is_empty() {
continue;
}
let value: serde_json::Value = match serde_json::from_str(&line) {
Ok(v) => v,
Err(_) => continue,
};
let msg_type = value.get("type").and_then(|t| t.as_str()).unwrap_or("");
if msg_type == "file-history-snapshot" || msg_type == "progress" {
continue;
}
// Skip system meta messages
if is_meta_message(&value) {
continue;
}
if cwd.is_none() {
cwd = value
.get("cwd")
.and_then(|s| s.as_str())
.map(|s| s.to_string());
}
if git_branch.is_none() {
git_branch = value
.get("gitBranch")
.and_then(|s| s.as_str())
.map(|s| s.to_string());
}
if let Some(ts_str) = value.get("timestamp").and_then(|t| t.as_str()) {
if let Ok(ts) = ts_str.parse::<DateTime<Utc>>() {
if first_timestamp.is_none() {
first_timestamp = Some(ts);
}
last_timestamp = Some(ts);
}
}
match msg_type {
"user" => {
let content = extract_user_content(&value);
// Skip user messages that are empty after system tag stripping
if content.is_empty() {
continue;
}
let timestamp = parse_timestamp(&value).unwrap_or_else(Utc::now);
let uuid = value
.get("uuid")
.and_then(|u| u.as_str())
.unwrap_or("")
.to_string();
if title.is_none() {
if let Some(first_text) = content.iter().find_map(|c| match c {
ContentBlock::Text { text } => Some(text.clone()),
_ => None,
}) {
title = Some(truncate_str(&first_text, 100));
}
}
messages.push(UnifiedMessage {
id: uuid,
role: MessageRole::User,
content,
timestamp,
usage: None,
duration_ms: None,
model: None,
});
}
"assistant" => {
let timestamp = parse_timestamp(&value).unwrap_or_else(Utc::now);
let uuid = value
.get("uuid")
.and_then(|u| u.as_str())
.unwrap_or("")
.to_string();
let msg_model = value
.get("message")
.and_then(|m| m.get("model"))
.and_then(|m| m.as_str())
.map(|s| s.to_string());
if model.is_none() {
model = msg_model.clone();
}
let content = extract_assistant_content(&value);
let usage = extract_usage(&value);
messages.push(UnifiedMessage {
id: uuid,
role: MessageRole::Assistant,
content,
timestamp,
usage,
duration_ms: None,
model: msg_model,
});
}
"system" => {
let subtype = value.get("subtype").and_then(|s| s.as_str()).unwrap_or("");
if subtype == "turn_duration" {
if let Some(duration) = value.get("durationMs").and_then(|d| d.as_u64()) {
// Attach to the last assistant message
if let Some(last) = messages
.iter_mut()
.rev()
.find(|m| matches!(m.role, MessageRole::Assistant))
{
last.duration_ms = Some(duration);
}
}
}
}
_ => {}
}
}
let folder_path = cwd.clone();
let folder_name = folder_path.as_ref().map(|p| folder_name_from_path(p));
let turns = group_into_turns(messages);
let context_window_used_tokens = latest_claude_context_window_used_tokens(&turns);
let context_window_max_tokens =
claude_context_window_max_tokens_for_model(model.as_deref());
let session_stats = merge_claude_context_window_stats(
super::compute_session_stats(&turns),
context_window_used_tokens,
context_window_max_tokens,
);
let summary = ConversationSummary {
id: conversation_id.to_string(),
agent_type: AgentType::ClaudeCode,
folder_path,
folder_name,
title,
started_at: first_timestamp.unwrap_or_else(Utc::now),
ended_at: last_timestamp,
message_count: turns.len() as u32,
model,
git_branch,
};
Ok(ConversationDetail {
summary,
turns,
session_stats,
})
}
}
fn parse_timestamp(value: &serde_json::Value) -> Option<DateTime<Utc>> {
value
.get("timestamp")
.and_then(|t| t.as_str())
.and_then(|s| s.parse::<DateTime<Utc>>().ok())
}
fn extract_user_text(value: &serde_json::Value) -> Option<String> {
let message = value.get("message")?;
let content = message.get("content")?;
if let Some(text) = content.as_str() {
return strip_system_tags(text);
}
if let Some(arr) = content.as_array() {
for item in arr {
if item.get("type").and_then(|t| t.as_str()) == Some("text") {
if let Some(text) = item.get("text").and_then(|t| t.as_str()) {
if let Some(cleaned) = strip_system_tags(text) {
return Some(cleaned);
}
}
}
}
}
None
}
fn extract_user_content(value: &serde_json::Value) -> Vec<ContentBlock> {
let mut blocks = Vec::new();
let message = match value.get("message") {
Some(m) => m,
None => return blocks,
};
let content = match message.get("content") {
Some(c) => c,
None => return blocks,
};
if let Some(text) = content.as_str() {
if let Some(cleaned) = strip_system_tags(text) {
blocks.push(ContentBlock::Text { text: cleaned });
}
return blocks;
}
if let Some(arr) = content.as_array() {
for item in arr {
let block_type = item.get("type").and_then(|t| t.as_str()).unwrap_or("");
match block_type {
"text" => {
if let Some(text) = item.get("text").and_then(|t| t.as_str()) {
if let Some(cleaned) = strip_system_tags(text) {
blocks.push(ContentBlock::Text { text: cleaned });
}
}
}
"tool_result" | "server_tool_result" => {
let tool_use_id = item
.get("tool_use_id")
.and_then(|n| n.as_str())
.map(|s| s.to_string());
let output = extract_tool_result_text(item);
let is_error = item
.get("is_error")
.and_then(|e| e.as_bool())
.unwrap_or(false);
blocks.push(ContentBlock::ToolResult {
tool_use_id,
output_preview: output,
is_error,
});
}
_ => {}
}
}
}
blocks
}
fn extract_assistant_content(value: &serde_json::Value) -> Vec<ContentBlock> {
let mut blocks = Vec::new();
let message = match value.get("message") {
Some(m) => m,
None => return blocks,
};
let content = match message.get("content") {
Some(c) => c,
None => return blocks,
};
if let Some(arr) = content.as_array() {
for item in arr {
let block_type = item.get("type").and_then(|t| t.as_str()).unwrap_or("");
match block_type {
"text" => {
if let Some(text) = item.get("text").and_then(|t| t.as_str()) {
blocks.push(ContentBlock::Text {
text: text.to_string(),
});
}
}
"thinking" => {
if let Some(text) = item.get("thinking").and_then(|t| t.as_str()) {
blocks.push(ContentBlock::Thinking {
text: text.to_string(),
});
}
}
"tool_use" | "server_tool_use" => {
let tool_use_id = item
.get("id")
.and_then(|n| n.as_str())
.map(|s| s.to_string());
let tool_name = item
.get("name")
.and_then(|n| n.as_str())
.unwrap_or("unknown")
.to_string();
let input_preview = item.get("input").map(|i| i.to_string());
blocks.push(ContentBlock::ToolUse {
tool_use_id,
tool_name,
input_preview,
});
}
_ => {}
}
}
}
blocks
}
fn extract_usage(value: &serde_json::Value) -> Option<TurnUsage> {
let usage = value.get("message")?.get("usage")?;
Some(TurnUsage {
input_tokens: usage
.get("input_tokens")
.and_then(|v| v.as_u64())
.unwrap_or(0),
output_tokens: usage
.get("output_tokens")
.and_then(|v| v.as_u64())
.unwrap_or(0),
cache_creation_input_tokens: usage
.get("cache_creation_input_tokens")
.and_then(|v| v.as_u64())
.unwrap_or(0),
cache_read_input_tokens: usage
.get("cache_read_input_tokens")
.and_then(|v| v.as_u64())
.unwrap_or(0),
})
}
fn extract_tool_result_text(item: &serde_json::Value) -> Option<String> {
let content = item.get("content")?;
if let Some(text) = content.as_str() {
return Some(text.to_string());
}
if let Some(arr) = content.as_array() {
let texts: Vec<String> = arr
.iter()
.filter_map(|c| {
if c.get("type").and_then(|t| t.as_str()) == Some("text") {
c.get("text")
.and_then(|t| t.as_str())
.map(|s| s.to_string())
} else {
None
}
})
.collect();
if !texts.is_empty() {
return Some(texts.join("\n"));
}
}
None
}
/// Check if a user message contains ONLY tool_result blocks (no text).
/// In Claude Code, tool results come back as "user" messages.
fn is_tool_result_only(msg: &UnifiedMessage) -> bool {
matches!(msg.role, MessageRole::User)
&& !msg.content.is_empty()
&& msg
.content
.iter()
.all(|b| matches!(b, ContentBlock::ToolResult { .. }))
}
/// Group flat messages into conversation turns.
/// Claude Code rule: assistant msg + following tool-result-only user msgs
/// merge into one Assistant turn.
fn group_into_turns(messages: Vec<UnifiedMessage>) -> Vec<MessageTurn> {
let mut turns = Vec::new();
let mut i = 0;
while i < messages.len() {
let msg = &messages[i];
if matches!(msg.role, MessageRole::Assistant) {
let mut blocks: Vec<ContentBlock> = msg.content.clone();
let timestamp = msg.timestamp;
let id = format!("turn-{}", turns.len());
let usage = msg.usage.clone();
let duration_ms = msg.duration_ms;
let turn_model = msg.model.clone();
i += 1;
// Absorb consecutive assistant msgs AND tool-result-only user msgs
while i < messages.len()
&& (matches!(messages[i].role, MessageRole::Assistant)
|| is_tool_result_only(&messages[i]))
{
blocks.extend(messages[i].content.clone());
i += 1;
}
turns.push(MessageTurn {
id,
role: TurnRole::Assistant,
blocks,
timestamp,
usage,
duration_ms,
model: turn_model,
});
} else if matches!(msg.role, MessageRole::System) {
turns.push(MessageTurn {
id: format!("turn-{}", turns.len()),
role: TurnRole::System,
blocks: msg.content.clone(),
timestamp: msg.timestamp,
usage: None,
duration_ms: None,
model: None,
});
i += 1;
} else {
turns.push(MessageTurn {
id: format!("turn-{}", turns.len()),
role: TurnRole::User,
blocks: msg.content.clone(),
timestamp: msg.timestamp,
usage: None,
duration_ms: None,
model: None,
});
i += 1;
}
}
turns
}
#[cfg(test)]
mod tests {
use std::io::Write;
use super::*;
#[test]
fn parses_model_capacity_suffix() {
assert_eq!(
parse_model_capacity_suffix("claude-sonnet-4-6[1.5M]"),
Some(1_500_000)
);
assert_eq!(
parse_model_capacity_suffix("claude-opus-4-6 [500k]"),
Some(500_000)
);
assert_eq!(parse_model_capacity_suffix("claude-sonnet-4-6"), None);
}
#[test]
fn defaults_context_limit_for_claude_models() {
assert_eq!(
claude_context_window_max_tokens_for_model(Some("claude-sonnet-4-6")),
Some(200_000)
);
assert_eq!(
claude_context_window_max_tokens_for_model(Some("custom-model-x")),
None
);
}
#[test]
fn uses_latest_assistant_usage_for_context_tokens() {
let timestamp = Utc::now();
let turns = vec![
MessageTurn {
id: "turn-0".to_string(),
role: TurnRole::Assistant,
blocks: vec![],
timestamp,
usage: Some(TurnUsage {
input_tokens: 100,
output_tokens: 20,
cache_creation_input_tokens: 30,
cache_read_input_tokens: 40,
}),
duration_ms: None,
model: None,
},
MessageTurn {
id: "turn-1".to_string(),
role: TurnRole::Assistant,
blocks: vec![],
timestamp,
usage: Some(TurnUsage {
input_tokens: 250,
output_tokens: 60,
cache_creation_input_tokens: 70,
cache_read_input_tokens: 80,
}),
duration_ms: None,
model: None,
},
];
assert_eq!(
latest_claude_context_window_used_tokens(&turns),
Some(250 + 70 + 80)
);
}
#[test]
fn parse_detail_sets_claude_context_window_stats() {
let path = std::env::temp_dir().join(format!(
"codeg-claude-parser-{}.jsonl",
uuid::Uuid::new_v4()
));
let mut file = fs::File::create(&path).expect("create temp jsonl");
writeln!(
file,
"{}",
serde_json::json!({
"type": "user",
"sessionId": "session-test",
"timestamp": "2026-03-01T10:00:00Z",
"uuid": "u1",
"cwd": "/tmp/demo",
"gitBranch": "main",
"message": {
"content": [{"type": "text", "text": "hello"}]
}
})
)
.expect("write user line");
writeln!(
file,
"{}",
serde_json::json!({
"type": "assistant",
"sessionId": "session-test",
"timestamp": "2026-03-01T10:00:02Z",
"uuid": "a1",
"message": {
"model": "claude-sonnet-4-6",
"content": [{"type": "text", "text": "world"}],
"usage": {
"input_tokens": 1000,
"output_tokens": 200,
"cache_creation_input_tokens": 300,
"cache_read_input_tokens": 400
}
}
})
)
.expect("write assistant line");
let parser = ClaudeParser {
base_dir: PathBuf::new(),
};
let detail = parser
.parse_conversation_detail(&path, "session-test")
.expect("parse conversation detail");
fs::remove_file(&path).expect("cleanup temp jsonl");
let stats = detail.session_stats.expect("session stats");
assert_eq!(stats.context_window_used_tokens, Some(1700));
assert_eq!(stats.context_window_max_tokens, Some(200_000));
let percent = stats
.context_window_usage_percent
.expect("context window usage percent");
assert!((percent - 0.85).abs() < f64::EPSILON);
}
#[test]
fn claude_config_dir_env_overrides_home() {
let resolved = resolve_claude_config_dir_from(
Some(std::ffi::OsString::from("/tmp/claude-config")),
Some(PathBuf::from("/Users/default")),
);
assert_eq!(resolved, PathBuf::from("/tmp/claude-config"));
}
#[test]
fn claude_config_dir_defaults_to_home_dot_claude() {
let resolved = resolve_claude_config_dir_from(
None,
Some(PathBuf::from("/Users/default")),
);
assert_eq!(resolved, PathBuf::from("/Users/default/.claude"));
}
}
File diff suppressed because it is too large Load Diff
+715
View File
@@ -0,0 +1,715 @@
use std::fs;
use std::path::{Path, PathBuf};
use chrono::{DateTime, Utc};
use walkdir::WalkDir;
use crate::models::*;
use crate::parsers::{folder_name_from_path, truncate_str, AgentParser, ParseError};
pub struct GeminiParser {
base_dir: PathBuf,
}
impl GeminiParser {
pub fn new() -> Self {
let base_dir = resolve_gemini_base_dir();
Self { base_dir }
}
#[cfg(test)]
fn with_base_dir(base_dir: PathBuf) -> Self {
Self { base_dir }
}
fn tmp_dir(&self) -> PathBuf {
self.base_dir.join("tmp")
}
fn history_dir(&self) -> PathBuf {
self.base_dir.join("history")
}
fn projects_json_path(&self) -> PathBuf {
self.base_dir.join("projects.json")
}
fn is_chat_file(path: &Path) -> bool {
if path.extension().and_then(|e| e.to_str()) != Some("json") {
return false;
}
let file_name = path.file_name().and_then(|n| n.to_str()).unwrap_or("");
if !file_name.starts_with("session-") {
return false;
}
path.parent()
.and_then(|p| p.file_name())
.and_then(|n| n.to_str())
== Some("chats")
}
fn list_chat_files(&self) -> Vec<PathBuf> {
let tmp_dir = self.tmp_dir();
if !tmp_dir.exists() {
return Vec::new();
}
let mut files: Vec<PathBuf> = WalkDir::new(&tmp_dir)
.into_iter()
.filter_map(|e| e.ok())
.map(|e| e.path().to_path_buf())
.filter(|p| p.is_file() && Self::is_chat_file(p))
.collect();
files.sort();
files
}
fn project_alias_from_chat_path(path: &Path) -> Option<String> {
path.parent()?
.parent()?
.file_name()
.map(|n| n.to_string_lossy().to_string())
}
fn read_project_root_file(path: PathBuf) -> Option<String> {
let raw = fs::read_to_string(path).ok()?;
let trimmed = raw.trim();
if trimmed.is_empty() {
None
} else {
Some(trimmed.to_string())
}
}
fn resolve_project_root(&self, alias: &str) -> Option<String> {
let tmp_root = self.tmp_dir().join(alias).join(".project_root");
if let Some(path) = Self::read_project_root_file(tmp_root) {
return Some(path);
}
let history_root = self.history_dir().join(alias).join(".project_root");
if let Some(path) = Self::read_project_root_file(history_root) {
return Some(path);
}
self.resolve_project_root_from_projects_json(alias)
}
fn resolve_project_root_from_projects_json(&self, alias: &str) -> Option<String> {
let raw = fs::read_to_string(self.projects_json_path()).ok()?;
let value: serde_json::Value = serde_json::from_str(&raw).ok()?;
let projects = value.get("projects")?.as_object()?;
projects
.iter()
.find_map(|(path, mapped_alias)| (mapped_alias.as_str() == Some(alias)).then(|| path))
.map(|s| s.to_string())
}
fn parse_timestamp(value: Option<&serde_json::Value>) -> Option<DateTime<Utc>> {
value.and_then(|v| v.as_str()?.parse::<DateTime<Utc>>().ok())
}
fn extract_text(value: &serde_json::Value) -> Option<String> {
match value {
serde_json::Value::String(text) => {
let trimmed = text.trim();
if trimmed.is_empty() {
None
} else {
Some(trimmed.to_string())
}
}
serde_json::Value::Array(items) => {
let mut parts = Vec::new();
for item in items {
if let Some(text) = item.get("text").and_then(Self::extract_text) {
parts.push(text);
} else if let Some(text) = Self::extract_text(item) {
parts.push(text);
}
}
if parts.is_empty() {
None
} else {
Some(parts.join("\n"))
}
}
serde_json::Value::Object(map) => {
if let Some(text) = map.get("text").and_then(Self::extract_text) {
return Some(text);
}
if let Some(text) = map.get("message").and_then(Self::extract_text) {
return Some(text);
}
None
}
_ => None,
}
}
fn extract_message_text(message: &serde_json::Value) -> Option<String> {
message
.get("content")
.and_then(Self::extract_text)
.or_else(|| message.get("message").and_then(Self::extract_text))
}
fn parse_summary_from_value(
&self,
path: &Path,
value: &serde_json::Value,
) -> Option<ConversationSummary> {
let id = value.get("sessionId").and_then(|v| v.as_str())?.to_string();
let messages = value
.get("messages")
.and_then(|m| m.as_array())
.cloned()
.unwrap_or_default();
let first_message_ts = messages
.first()
.and_then(|m| Self::parse_timestamp(m.get("timestamp")));
let last_message_ts = messages
.iter()
.rev()
.find_map(|m| Self::parse_timestamp(m.get("timestamp")));
let started_at = Self::parse_timestamp(value.get("startTime"))
.or(first_message_ts)
.unwrap_or_else(Utc::now);
let ended_at = Self::parse_timestamp(value.get("lastUpdated")).or(last_message_ts);
let title = messages
.iter()
.filter(|m| m.get("type").and_then(|t| t.as_str()) == Some("user"))
.find_map(Self::extract_message_text)
.map(|t| truncate_str(&t, 100));
let model = messages.iter().rev().find_map(|m| {
m.get("model")
.and_then(|v| v.as_str())
.map(|s| s.to_string())
});
let folder_alias = Self::project_alias_from_chat_path(path);
let folder_path = folder_alias
.as_deref()
.and_then(|alias| self.resolve_project_root(alias));
let folder_name = folder_path
.as_ref()
.map(|p| folder_name_from_path(p))
.or(folder_alias);
Some(ConversationSummary {
id,
agent_type: AgentType::Gemini,
folder_path,
folder_name,
title,
started_at,
ended_at,
message_count: messages.len() as u32,
model,
git_branch: None,
})
}
fn result_preview(result: Option<&serde_json::Value>) -> Option<String> {
let v = result?;
if let Some(s) = v.as_str() {
let trimmed = s.trim();
if trimmed.is_empty() {
return None;
}
return Some(trimmed.to_string());
}
serde_json::to_string(v).ok()
}
fn tool_call_is_error(call: &serde_json::Value, output_preview: Option<&str>) -> bool {
if call
.get("status")
.and_then(|v| v.as_str())
.map(|s| {
matches!(
s.to_ascii_lowercase().as_str(),
"error" | "failed" | "failure" | "cancelled" | "canceled"
)
})
.unwrap_or(false)
{
return true;
}
if call
.get("result")
.and_then(|r| r.as_array())
.map(|items| {
items.iter().any(|item| {
item.get("functionResponse")
.and_then(|fr| fr.get("response"))
.and_then(|resp| resp.get("error"))
.is_some()
})
})
.unwrap_or(false)
{
return true;
}
output_preview
.map(|s| s.trim_start().to_ascii_lowercase().starts_with("error"))
.unwrap_or(false)
}
fn parse_assistant_blocks(message: &serde_json::Value) -> Vec<ContentBlock> {
let mut blocks: Vec<ContentBlock> = Vec::new();
if let Some(thoughts) = message.get("thoughts").and_then(|v| v.as_array()) {
for thought in thoughts {
let subject = thought
.get("subject")
.and_then(|v| v.as_str())
.map(str::trim)
.filter(|s| !s.is_empty());
let description = thought
.get("description")
.and_then(|v| v.as_str())
.map(str::trim)
.filter(|s| !s.is_empty());
let text = match (subject, description) {
(Some(sub), Some(desc)) => format!("{sub}: {desc}"),
(Some(sub), None) => sub.to_string(),
(None, Some(desc)) => desc.to_string(),
(None, None) => continue,
};
blocks.push(ContentBlock::Thinking { text });
}
}
if let Some(tool_calls) = message.get("toolCalls").and_then(|v| v.as_array()) {
for call in tool_calls {
let tool_use_id = call
.get("id")
.and_then(|v| v.as_str())
.map(|s| s.to_string());
let tool_name = call
.get("displayName")
.or_else(|| call.get("name"))
.and_then(|v| v.as_str())
.unwrap_or("unknown")
.to_string();
let input_preview = call
.get("args")
.and_then(|v| serde_json::to_string(v).ok())
.or_else(|| {
call.get("input")
.and_then(|v| Self::result_preview(Some(v)))
});
blocks.push(ContentBlock::ToolUse {
tool_use_id: tool_use_id.clone(),
tool_name,
input_preview,
});
let output_preview = call
.get("resultDisplay")
.and_then(|v| v.as_str())
.map(str::trim)
.filter(|s| !s.is_empty())
.map(|s| s.to_string())
.or_else(|| Self::result_preview(call.get("result")));
let is_error = Self::tool_call_is_error(call, output_preview.as_deref());
blocks.push(ContentBlock::ToolResult {
tool_use_id,
output_preview,
is_error,
});
}
}
if let Some(text) = Self::extract_message_text(message) {
blocks.push(ContentBlock::Text { text });
}
blocks
}
fn parse_usage(message: &serde_json::Value) -> Option<TurnUsage> {
let tokens = message.get("tokens")?;
let input_tokens = tokens.get("input").and_then(|v| v.as_u64()).unwrap_or(0);
let output_tokens = tokens.get("output").and_then(|v| v.as_u64()).unwrap_or(0);
let cached_tokens = tokens.get("cached").and_then(|v| v.as_u64()).unwrap_or(0);
Some(TurnUsage {
input_tokens,
output_tokens,
cache_creation_input_tokens: 0,
cache_read_input_tokens: cached_tokens,
})
}
fn parse_conversation_detail(
&self,
path: &Path,
value: &serde_json::Value,
conversation_id: &str,
) -> Result<ConversationDetail, ParseError> {
let mut summary = self
.parse_summary_from_value(path, value)
.ok_or_else(|| ParseError::ConversationNotFound(conversation_id.to_string()))?;
let messages_raw = value
.get("messages")
.and_then(|m| m.as_array())
.cloned()
.unwrap_or_default();
let mut messages: Vec<UnifiedMessage> = Vec::new();
for raw in messages_raw {
let msg_id = raw
.get("id")
.and_then(|v| v.as_str())
.map(|s| s.to_string())
.unwrap_or_else(|| format!("msg-{}", messages.len()));
let timestamp =
Self::parse_timestamp(raw.get("timestamp")).unwrap_or(summary.started_at);
let msg_type = raw
.get("type")
.and_then(|v| v.as_str())
.unwrap_or("")
.to_ascii_lowercase();
match msg_type.as_str() {
"user" => {
let Some(text) = Self::extract_message_text(&raw) else {
continue;
};
messages.push(UnifiedMessage {
id: msg_id,
role: MessageRole::User,
content: vec![ContentBlock::Text { text }],
timestamp,
usage: None,
duration_ms: None,
model: None,
});
}
"gemini" | "assistant" | "model" => {
let blocks = Self::parse_assistant_blocks(&raw);
if blocks.is_empty() {
continue;
}
messages.push(UnifiedMessage {
id: msg_id,
role: MessageRole::Assistant,
content: blocks,
timestamp,
usage: Self::parse_usage(&raw),
duration_ms: None,
model: raw
.get("model")
.and_then(|v| v.as_str())
.map(|s| s.to_string()),
});
}
"system" => {
let Some(text) = Self::extract_message_text(&raw) else {
continue;
};
messages.push(UnifiedMessage {
id: msg_id,
role: MessageRole::System,
content: vec![ContentBlock::Text { text }],
timestamp,
usage: None,
duration_ms: None,
model: None,
});
}
_ => {}
}
}
// Approximate duration for assistant messages from adjacent timestamps
for i in 0..messages.len() {
if matches!(messages[i].role, MessageRole::Assistant)
&& messages[i].duration_ms.is_none()
{
if let Some(next) = messages.get(i + 1) {
let dur = (next.timestamp - messages[i].timestamp).num_milliseconds();
if dur > 0 && dur < 300_000 {
messages[i].duration_ms = Some(dur as u64);
}
}
}
}
let turns = group_into_turns(messages);
summary.message_count = turns.len() as u32;
summary.id = conversation_id.to_string();
let context_window_used_tokens = super::latest_turn_total_usage_tokens(&turns);
let context_window_max_tokens =
super::infer_context_window_max_tokens(summary.model.as_deref());
let session_stats = super::merge_context_window_stats(
super::compute_session_stats(&turns),
context_window_used_tokens,
context_window_max_tokens,
);
Ok(ConversationDetail {
summary,
turns,
session_stats,
})
}
}
fn resolve_gemini_base_dir() -> PathBuf {
resolve_gemini_base_dir_from(
std::env::var_os("GEMINI_CLI_HOME"),
dirs::home_dir(),
)
}
fn resolve_gemini_base_dir_from(
gemini_cli_home_env: Option<std::ffi::OsString>,
home_dir: Option<PathBuf>,
) -> PathBuf {
gemini_cli_home_env
.filter(|value| !value.is_empty())
.map(PathBuf::from)
.unwrap_or_else(|| home_dir.unwrap_or_default())
.join(".gemini")
}
impl AgentParser for GeminiParser {
fn list_conversations(&self) -> Result<Vec<ConversationSummary>, ParseError> {
let mut conversations = Vec::new();
for chat_file in self.list_chat_files() {
let raw = match fs::read_to_string(&chat_file) {
Ok(raw) => raw,
Err(_) => continue,
};
let value: serde_json::Value = match serde_json::from_str(&raw) {
Ok(value) => value,
Err(_) => continue,
};
if let Some(summary) = self.parse_summary_from_value(&chat_file, &value) {
conversations.push(summary);
}
}
conversations.sort_by(|a, b| b.started_at.cmp(&a.started_at));
Ok(conversations)
}
fn get_conversation(&self, conversation_id: &str) -> Result<ConversationDetail, ParseError> {
for chat_file in self.list_chat_files() {
let raw = match fs::read_to_string(&chat_file) {
Ok(raw) => raw,
Err(_) => continue,
};
if !raw.contains(conversation_id) {
continue;
}
let value: serde_json::Value = match serde_json::from_str(&raw) {
Ok(value) => value,
Err(_) => continue,
};
let session_id = value.get("sessionId").and_then(|v| v.as_str());
if session_id != Some(conversation_id) {
continue;
}
return self.parse_conversation_detail(&chat_file, &value, conversation_id);
}
Err(ParseError::ConversationNotFound(
conversation_id.to_string(),
))
}
}
fn group_into_turns(messages: Vec<UnifiedMessage>) -> Vec<MessageTurn> {
let mut turns = Vec::new();
let mut i = 0;
while i < messages.len() {
let msg = &messages[i];
if matches!(msg.role, MessageRole::User) {
turns.push(MessageTurn {
id: format!("turn-{}", turns.len()),
role: TurnRole::User,
blocks: msg.content.clone(),
timestamp: msg.timestamp,
usage: None,
duration_ms: None,
model: None,
});
i += 1;
continue;
}
if matches!(msg.role, MessageRole::System) {
turns.push(MessageTurn {
id: format!("turn-{}", turns.len()),
role: TurnRole::System,
blocks: msg.content.clone(),
timestamp: msg.timestamp,
usage: None,
duration_ms: None,
model: None,
});
i += 1;
continue;
}
let mut blocks = msg.content.clone();
let mut usage = msg.usage.clone();
let mut duration_ms = msg.duration_ms;
let mut models: Vec<String> = msg.model.iter().cloned().collect();
let timestamp = msg.timestamp;
i += 1;
while i < messages.len()
&& (matches!(messages[i].role, MessageRole::Assistant)
|| matches!(messages[i].role, MessageRole::Tool))
{
blocks.extend(messages[i].content.clone());
if usage.is_none() {
usage = messages[i].usage.clone();
}
if duration_ms.is_none() {
duration_ms = messages[i].duration_ms;
}
if let Some(model) = &messages[i].model {
models.push(model.clone());
}
i += 1;
}
let model = models.pop();
turns.push(MessageTurn {
id: format!("turn-{}", turns.len()),
role: TurnRole::Assistant,
blocks,
timestamp,
usage,
duration_ms,
model,
});
}
turns
}
#[cfg(test)]
mod tests {
use super::GeminiParser;
use super::resolve_gemini_base_dir_from;
use crate::parsers::AgentParser;
use std::env;
use std::fs;
use std::path::PathBuf;
use std::time::{SystemTime, UNIX_EPOCH};
#[test]
fn parses_gemini_session_detail_from_chat_json() {
let nanos = SystemTime::now()
.duration_since(UNIX_EPOCH)
.expect("system time ok")
.as_nanos();
let base: PathBuf = env::temp_dir().join(format!("codeg-gemini-test-{nanos}"));
let chats_dir = base.join("tmp").join("codeg").join("chats");
fs::create_dir_all(&chats_dir).expect("create chat dir");
fs::write(
base.join("tmp").join("codeg").join(".project_root"),
"/Users/test/workspace/demo",
)
.expect("write project root");
let file_path = chats_dir.join("session-2026-03-02T04-30-32c7d221.json");
let content = r#"{
"sessionId": "32c7d221-0553-46c8-ba50-e664719cae7f",
"projectHash": "abc",
"startTime": "2026-03-02T04:30:20.796Z",
"lastUpdated": "2026-03-02T04:33:13.631Z",
"messages": [
{
"id": "u1",
"timestamp": "2026-03-02T04:30:20.796Z",
"type": "user",
"content": [{"text": "你会做什么"}]
},
{
"id": "a1",
"timestamp": "2026-03-02T04:33:13.631Z",
"type": "gemini",
"content": "我是一个助手",
"toolCalls": [
{
"id": "cli_help-1",
"name": "cli_help",
"args": {"question": "你会做什么"},
"resultDisplay": "ok",
"status": "success"
}
],
"tokens": {"input": 12, "output": 34, "cached": 5},
"model": "gemini-3.1-pro-preview"
}
]
}"#;
fs::write(&file_path, content).expect("write chat file");
let parser = GeminiParser::with_base_dir(base.clone());
let summaries = parser.list_conversations().expect("list conversations");
assert_eq!(summaries.len(), 1);
assert_eq!(
summaries[0].id,
"32c7d221-0553-46c8-ba50-e664719cae7f".to_string()
);
let detail = parser
.get_conversation("32c7d221-0553-46c8-ba50-e664719cae7f")
.expect("get conversation");
assert_eq!(detail.turns.len(), 2);
assert_eq!(
detail.summary.folder_path.as_deref(),
Some("/Users/test/workspace/demo")
);
assert!(detail.session_stats.is_some());
let stats = detail.session_stats.expect("session stats");
assert_eq!(stats.context_window_used_tokens, Some(51));
assert_eq!(stats.context_window_max_tokens, Some(1_000_000));
let percent = stats
.context_window_usage_percent
.expect("context window percent");
assert!((percent - 0.0051).abs() < 1e-9);
let _ = fs::remove_dir_all(base);
}
#[test]
fn gemini_cli_home_env_overrides_user_home() {
let resolved = resolve_gemini_base_dir_from(
Some(std::ffi::OsString::from("/tmp/gemini-home")),
Some(PathBuf::from("/Users/default")),
);
assert_eq!(resolved, PathBuf::from("/tmp/gemini-home/.gemini"));
}
#[test]
fn gemini_defaults_to_home_dot_gemini() {
let resolved = resolve_gemini_base_dir_from(
None,
Some(PathBuf::from("/Users/default")),
);
assert_eq!(resolved, PathBuf::from("/Users/default/.gemini"));
}
}
+356
View File
@@ -0,0 +1,356 @@
pub mod claude;
pub mod codex;
pub mod gemini;
pub mod opencode;
use std::sync::OnceLock;
use regex::Regex;
use crate::models::{
ConversationDetail, ConversationSummary, MessageTurn, SessionStats, TurnUsage,
};
#[derive(Debug, thiserror::Error)]
pub enum ParseError {
#[error("IO error: {0}")]
Io(#[from] std::io::Error),
#[error("JSON parse error: {0}")]
Json(#[from] serde_json::Error),
#[error("Database error: {0}")]
Db(#[from] sea_orm::DbErr),
#[error("Conversation not found: {0}")]
ConversationNotFound(String),
#[allow(dead_code)]
#[error("Invalid data: {0}")]
InvalidData(String),
}
pub trait AgentParser {
fn list_conversations(&self) -> Result<Vec<ConversationSummary>, ParseError>;
fn get_conversation(&self, conversation_id: &str) -> Result<ConversationDetail, ParseError>;
}
/// Truncate a string to `max_len` characters, appending "..." if truncated.
pub fn truncate_str(s: &str, max_len: usize) -> String {
if s.chars().count() <= max_len {
s.to_string()
} else {
let truncated: String = s.chars().take(max_len).collect();
format!("{}...", truncated)
}
}
/// Aggregate turn-level usage and duration into a single `SessionStats`.
pub fn compute_session_stats(turns: &[MessageTurn]) -> Option<SessionStats> {
let mut total_in = 0u64;
let mut total_out = 0u64;
let mut total_cache_create = 0u64;
let mut total_cache_read = 0u64;
let mut total_duration = 0u64;
let mut has_data = false;
for turn in turns {
if let Some(ref u) = turn.usage {
total_in += u.input_tokens;
total_out += u.output_tokens;
total_cache_create += u.cache_creation_input_tokens;
total_cache_read += u.cache_read_input_tokens;
has_data = true;
}
if let Some(d) = turn.duration_ms {
total_duration += d;
}
}
if !has_data {
return None;
}
Some(SessionStats {
total_usage: Some(TurnUsage {
input_tokens: total_in,
output_tokens: total_out,
cache_creation_input_tokens: total_cache_create,
cache_read_input_tokens: total_cache_read,
}),
total_tokens: Some(total_in + total_out + total_cache_create + total_cache_read),
total_duration_ms: total_duration,
context_window_used_tokens: None,
context_window_max_tokens: None,
context_window_usage_percent: None,
})
}
fn model_capacity_suffix_regex() -> &'static Regex {
static RE: OnceLock<Regex> = OnceLock::new();
RE.get_or_init(|| {
Regex::new(r"(?i)\[\s*([0-9]+(?:\.[0-9]+)?)\s*([km])\s*\]\s*$")
.expect("valid model capacity regex")
})
}
fn parse_model_capacity_suffix(model: &str) -> Option<u64> {
let captures = model_capacity_suffix_regex().captures(model.trim())?;
let value = captures.get(1)?.as_str().parse::<f64>().ok()?;
if !value.is_finite() || value <= 0.0 {
return None;
}
let unit = captures
.get(2)
.map(|m| m.as_str().to_ascii_lowercase())
.unwrap_or_default();
let multiplier = match unit.as_str() {
"m" => 1_000_000.0,
"k" => 1_000.0,
_ => return None,
};
Some((value * multiplier) as u64)
}
pub fn infer_context_window_max_tokens(model: Option<&str>) -> Option<u64> {
let raw = model?.trim();
if raw.is_empty() {
return None;
}
if let Some(suffixed_limit) = parse_model_capacity_suffix(raw) {
return Some(suffixed_limit);
}
let normalized = raw
.rsplit('/')
.next()
.unwrap_or(raw)
.split(':')
.next()
.unwrap_or(raw)
.trim()
.to_ascii_lowercase();
if normalized.starts_with("claude") {
return Some(200_000);
}
if normalized.starts_with("gemini") {
return Some(1_000_000);
}
match normalized.as_str() {
"gpt-5.2-codex" | "gpt-5.1-codex-max" | "gpt-5.1-codex-mini" | "gpt-5.2" => Some(258_000),
"gpt-5.1" | "gpt-5.1-codex" | "gpt-4o" | "gpt-4o-mini" | "gpt-4-turbo" | "o1-mini"
| "o1-preview" => Some(128_000),
"gpt-4" => Some(8_192),
"o3" | "o3-mini" | "o1" => Some(200_000),
_ => {
if normalized.starts_with("gpt-5") {
Some(258_000)
} else if normalized.starts_with("gpt-4o")
|| normalized.starts_with("gpt-4.1")
|| normalized.starts_with("gpt-4-turbo")
{
Some(128_000)
} else if normalized.starts_with("o3") || normalized == "o1" {
Some(200_000)
} else if normalized.starts_with("o1-mini") || normalized.starts_with("o1-preview") {
Some(128_000)
} else {
None
}
}
}
}
pub fn latest_turn_total_usage_tokens(turns: &[MessageTurn]) -> Option<u64> {
turns.iter().rev().find_map(|turn| {
turn.usage.as_ref().map(|usage| {
usage
.input_tokens
.saturating_add(usage.output_tokens)
.saturating_add(usage.cache_creation_input_tokens)
.saturating_add(usage.cache_read_input_tokens)
})
})
}
pub fn merge_context_window_stats(
stats: Option<SessionStats>,
used_tokens: Option<u64>,
max_tokens: Option<u64>,
) -> Option<SessionStats> {
if used_tokens.is_none() && max_tokens.is_none() {
return stats;
}
let usage_percent = match (used_tokens, max_tokens) {
(Some(used), Some(max)) if max > 0 => Some((used as f64 / max as f64) * 100.0),
_ => None,
};
match stats {
Some(mut s) => {
s.context_window_used_tokens = used_tokens;
s.context_window_max_tokens = max_tokens;
s.context_window_usage_percent = usage_percent;
Some(s)
}
None => Some(SessionStats {
total_usage: None,
total_tokens: None,
total_duration_ms: 0,
context_window_used_tokens: used_tokens,
context_window_max_tokens: max_tokens,
context_window_usage_percent: usage_percent,
}),
}
}
/// Extract the last path component as the folder name.
pub fn folder_name_from_path(path: &str) -> String {
path.rsplit(['/', '\\']).next().unwrap_or(path).to_string()
}
/// Normalize a filesystem path string for tolerant cross-platform comparison.
/// This intentionally does not hit the filesystem (no canonicalize), and only
/// normalizes separators/casing differences that commonly break exact matching.
pub fn normalize_path_for_matching(path: &str) -> String {
let mut normalized = path.trim().replace('\\', "/");
#[cfg(target_os = "windows")]
{
if let Some(stripped) = normalized.strip_prefix("//?/") {
normalized = stripped.to_string();
}
normalized = normalized.to_ascii_lowercase();
}
while normalized.ends_with('/') {
if normalized == "/" {
break;
}
// Keep Windows drive root such as "c:/" intact.
if normalized.len() == 3
&& normalized.as_bytes().get(1) == Some(&b':')
&& normalized.as_bytes().get(2) == Some(&b'/')
{
break;
}
normalized.pop();
}
normalized
}
pub fn path_eq_for_matching(left: &str, right: &str) -> bool {
normalize_path_for_matching(left) == normalize_path_for_matching(right)
}
#[cfg(test)]
mod tests {
use chrono::Utc;
use super::{
infer_context_window_max_tokens, latest_turn_total_usage_tokens, merge_context_window_stats,
path_eq_for_matching,
};
use crate::models::{MessageTurn, SessionStats, TurnRole, TurnUsage};
#[test]
fn infers_model_context_limits() {
assert_eq!(
infer_context_window_max_tokens(Some("claude-sonnet-4-6")),
Some(200_000)
);
assert_eq!(
infer_context_window_max_tokens(Some("gemini-2.5-pro")),
Some(1_000_000)
);
assert_eq!(
infer_context_window_max_tokens(Some("claude-sonnet-4-6 [1.5M]")),
Some(1_500_000)
);
assert_eq!(infer_context_window_max_tokens(Some("unknown-model")), None);
}
#[test]
fn picks_latest_turn_usage_total_tokens() {
let timestamp = Utc::now();
let turns = vec![
MessageTurn {
id: "turn-0".to_string(),
role: TurnRole::Assistant,
blocks: vec![],
timestamp,
usage: Some(TurnUsage {
input_tokens: 10,
output_tokens: 20,
cache_creation_input_tokens: 30,
cache_read_input_tokens: 40,
}),
duration_ms: None,
model: None,
},
MessageTurn {
id: "turn-1".to_string(),
role: TurnRole::Assistant,
blocks: vec![],
timestamp,
usage: Some(TurnUsage {
input_tokens: 11,
output_tokens: 21,
cache_creation_input_tokens: 31,
cache_read_input_tokens: 41,
}),
duration_ms: None,
model: None,
},
];
assert_eq!(latest_turn_total_usage_tokens(&turns), Some(104));
}
#[test]
fn merges_context_window_stats() {
let merged = merge_context_window_stats(None, Some(1500), Some(3000))
.expect("context stats should exist");
assert_eq!(merged.context_window_used_tokens, Some(1500));
assert_eq!(merged.context_window_max_tokens, Some(3000));
assert!(merged.total_usage.is_none());
let percent = merged
.context_window_usage_percent
.expect("usage percent should exist");
assert!((percent - 50.0).abs() < f64::EPSILON);
let existing = Some(SessionStats {
total_usage: Some(TurnUsage {
input_tokens: 1,
output_tokens: 2,
cache_creation_input_tokens: 3,
cache_read_input_tokens: 4,
}),
total_tokens: Some(10),
total_duration_ms: 100,
context_window_used_tokens: None,
context_window_max_tokens: None,
context_window_usage_percent: None,
});
let merged_existing =
merge_context_window_stats(existing, Some(200), Some(1000)).expect("merged");
assert_eq!(merged_existing.total_tokens, Some(10));
assert_eq!(merged_existing.context_window_used_tokens, Some(200));
assert_eq!(merged_existing.context_window_max_tokens, Some(1000));
}
#[test]
fn path_matching_handles_separator_differences() {
assert!(path_eq_for_matching(
"/Users/demo/workspace/codeg",
"/Users/demo/workspace/codeg/"
));
assert!(path_eq_for_matching(
"C:\\Users\\demo\\workspace\\codeg",
"C:/Users/demo/workspace/codeg"
));
}
}
+646
View File
@@ -0,0 +1,646 @@
use std::future::Future;
use std::path::PathBuf;
use std::time::Duration;
use chrono::{DateTime, TimeZone, Utc};
use sea_orm::{
ConnectOptions, ConnectionTrait, Database, DatabaseConnection, DbBackend, QueryResult,
Statement,
};
use crate::models::*;
use crate::parsers::{folder_name_from_path, AgentParser, ParseError};
pub struct OpenCodeParser {
base_dir: PathBuf,
}
impl OpenCodeParser {
pub fn new() -> Self {
let base_dir = resolve_opencode_base_dir();
Self { base_dir }
}
fn sqlite_db_path(&self) -> PathBuf {
self.base_dir.join("opencode.db")
}
fn block_on<F, T>(&self, fut: F) -> Result<T, ParseError>
where
F: Future<Output = Result<T, ParseError>>,
{
let runtime = tokio::runtime::Builder::new_current_thread()
.enable_all()
.build()
.map_err(|e| ParseError::InvalidData(format!("failed to build runtime: {e}")))?;
runtime.block_on(fut)
}
async fn open_sqlite_connection(&self) -> Result<DatabaseConnection, ParseError> {
let db_path = self.sqlite_db_path();
let db_url = format!(
"sqlite:{}?mode=ro",
urlencoding::encode(&db_path.to_string_lossy())
);
let mut opts = ConnectOptions::new(db_url);
opts.max_connections(1)
.min_connections(1)
.connect_timeout(Duration::from_secs(5))
.idle_timeout(Duration::from_secs(30))
.sqlx_logging(false);
let conn = Database::connect(opts).await?;
conn.execute(Statement::from_string(
DbBackend::Sqlite,
"PRAGMA busy_timeout=3000;".to_owned(),
))
.await?;
Ok(conn)
}
fn parse_sqlite_summary_row(row: &QueryResult) -> Result<ConversationSummary, ParseError> {
let id: String = row.try_get("", "id")?;
let directory: Option<String> = row.try_get("", "directory")?;
let title: Option<String> = row.try_get("", "title")?;
let created_ms: i64 = row.try_get("", "created_ms")?;
let updated_ms: i64 = row.try_get("", "updated_ms")?;
let message_count_i64: i64 = row.try_get("", "message_count")?;
let model: Option<String> = row.try_get("", "model")?;
let folder_path = normalize_optional_string(directory);
let folder_name = folder_path.as_ref().map(|p| folder_name_from_path(p));
let message_count = if message_count_i64 <= 0 {
0
} else {
u32::try_from(message_count_i64).unwrap_or(u32::MAX)
};
Ok(ConversationSummary {
id,
agent_type: AgentType::OpenCode,
folder_path,
folder_name,
title: normalize_optional_string(title),
started_at: millis_to_datetime(created_ms),
ended_at: (updated_ms > 0).then(|| millis_to_datetime(updated_ms)),
message_count,
model: normalize_optional_string(model),
git_branch: None,
})
}
async fn list_conversations_from_sqlite(&self) -> Result<Vec<ConversationSummary>, ParseError> {
let conn = self.open_sqlite_connection().await?;
let rows = conn
.query_all(Statement::from_string(
DbBackend::Sqlite,
r#"
SELECT
s.id AS id,
s.directory AS directory,
s.title AS title,
s.time_created AS created_ms,
s.time_updated AS updated_ms,
COALESCE((
SELECT COUNT(*)
FROM message m
WHERE m.session_id = s.id
), 0) AS message_count,
(
SELECT json_extract(m2.data, '$.modelID')
FROM message m2
WHERE m2.session_id = s.id
AND json_extract(m2.data, '$.role') = 'assistant'
ORDER BY m2.time_created DESC
LIMIT 1
) AS model
FROM session s
ORDER BY s.time_created DESC
"#
.to_string(),
))
.await?;
let mut conversations = Vec::with_capacity(rows.len());
for row in rows {
conversations.push(Self::parse_sqlite_summary_row(&row)?);
}
Ok(conversations)
}
async fn sqlite_summary_by_id(
&self,
conn: &DatabaseConnection,
conversation_id: &str,
) -> Result<Option<ConversationSummary>, ParseError> {
let row = conn
.query_one(Statement::from_sql_and_values(
DbBackend::Sqlite,
r#"
SELECT
s.id AS id,
s.directory AS directory,
s.title AS title,
s.time_created AS created_ms,
s.time_updated AS updated_ms,
COALESCE((
SELECT COUNT(*)
FROM message m
WHERE m.session_id = s.id
), 0) AS message_count,
(
SELECT json_extract(m2.data, '$.modelID')
FROM message m2
WHERE m2.session_id = s.id
AND json_extract(m2.data, '$.role') = 'assistant'
ORDER BY m2.time_created DESC
LIMIT 1
) AS model
FROM session s
WHERE s.id = ?
LIMIT 1
"#,
[conversation_id.into()],
))
.await?;
row.map(|r| Self::parse_sqlite_summary_row(&r)).transpose()
}
async fn get_conversation_from_sqlite(
&self,
conversation_id: &str,
) -> Result<ConversationDetail, ParseError> {
let conn = self.open_sqlite_connection().await?;
let summary = self
.sqlite_summary_by_id(&conn, conversation_id)
.await?
.ok_or_else(|| ParseError::ConversationNotFound(conversation_id.to_string()))?;
let messages = self.load_sqlite_messages(&conn, conversation_id).await?;
let turns = group_into_turns(messages);
let context_window_used_tokens = super::latest_turn_total_usage_tokens(&turns);
let context_window_max_tokens =
super::infer_context_window_max_tokens(summary.model.as_deref());
let session_stats = super::merge_context_window_stats(
super::compute_session_stats(&turns),
context_window_used_tokens,
context_window_max_tokens,
);
Ok(ConversationDetail {
summary,
turns,
session_stats,
})
}
async fn load_sqlite_messages(
&self,
conn: &DatabaseConnection,
conversation_id: &str,
) -> Result<Vec<UnifiedMessage>, ParseError> {
let rows = conn
.query_all(Statement::from_sql_and_values(
DbBackend::Sqlite,
r#"
SELECT id, time_created, data
FROM message
WHERE session_id = ?
ORDER BY time_created ASC, id ASC
"#,
[conversation_id.into()],
))
.await?;
let mut messages = Vec::with_capacity(rows.len());
for row in rows {
let msg_id: String = row.try_get("", "id")?;
let row_time_created: i64 = row.try_get("", "time_created")?;
let data_raw: String = row.try_get("", "data")?;
let value: serde_json::Value = match serde_json::from_str(&data_raw) {
Ok(v) => v,
Err(_) => continue,
};
let role = match value.get("role").and_then(|r| r.as_str()) {
Some("user") => MessageRole::User,
Some("assistant") => MessageRole::Assistant,
Some("system") => MessageRole::System,
Some("tool") => MessageRole::Tool,
_ => continue,
};
let created_ms = value
.get("time")
.and_then(|t| t.get("created"))
.and_then(|c| c.as_i64())
.unwrap_or(row_time_created);
let timestamp = millis_to_datetime(created_ms);
let is_assistant = matches!(role, MessageRole::Assistant);
let msg_model = if is_assistant {
value
.get("modelID")
.and_then(|m| m.as_str())
.map(|s| s.to_string())
} else {
None
};
let (content_blocks, usage_from_step_finish) =
self.load_sqlite_parts(conn, &msg_id).await?;
let usage = if is_assistant {
extract_opencode_usage(&value).or(usage_from_step_finish)
} else {
None
};
let duration_ms = if is_assistant {
let completed_ms = value
.get("time")
.and_then(|t| t.get("completed"))
.and_then(|c| c.as_i64());
match completed_ms {
Some(done) if done > created_ms => Some((done - created_ms) as u64),
_ => None,
}
} else {
None
};
messages.push(UnifiedMessage {
id: msg_id,
role,
content: content_blocks,
timestamp,
usage,
duration_ms,
model: msg_model,
});
}
Ok(messages)
}
async fn load_sqlite_parts(
&self,
conn: &DatabaseConnection,
message_id: &str,
) -> Result<(Vec<ContentBlock>, Option<TurnUsage>), ParseError> {
let rows = conn
.query_all(Statement::from_sql_and_values(
DbBackend::Sqlite,
r#"
SELECT data
FROM part
WHERE message_id = ?
ORDER BY time_created ASC, id ASC
"#,
[message_id.into()],
))
.await?;
let mut blocks = Vec::new();
let mut usage_from_step_finish: Option<TurnUsage> = None;
for row in rows {
let data_raw: String = row.try_get("", "data")?;
let value: serde_json::Value = match serde_json::from_str(&data_raw) {
Ok(v) => v,
Err(_) => continue,
};
let part_type = value.get("type").and_then(|t| t.as_str()).unwrap_or("");
match part_type {
"text" => {
if let Some(text) = value
.get("text")
.and_then(|t| t.as_str())
.map(str::trim)
.filter(|s| !s.is_empty())
{
blocks.push(ContentBlock::Text {
text: text.to_string(),
});
}
}
"reasoning" => {
if let Some(text) = value
.get("text")
.and_then(|t| t.as_str())
.map(str::trim)
.filter(|s| !s.is_empty())
{
blocks.push(ContentBlock::Thinking {
text: text.to_string(),
});
}
}
"tool" => {
let tool_name = value
.get("tool")
.and_then(|t| t.as_str())
.unwrap_or("unknown")
.to_string();
let call_id = value
.get("callID")
.and_then(|c| c.as_str())
.map(|s| s.to_string());
let status = value
.get("state")
.and_then(|s| s.get("status"))
.and_then(|s| s.as_str())
.unwrap_or("");
let input_preview = value
.get("state")
.and_then(|s| s.get("input"))
.and_then(|v| value_to_preview(Some(v)));
blocks.push(ContentBlock::ToolUse {
tool_use_id: call_id.clone(),
tool_name,
input_preview,
});
let output_preview = value
.get("state")
.and_then(|s| s.get("output"))
.and_then(|v| value_to_preview(Some(v)));
let has_error_field = value.get("state").and_then(|s| s.get("error")).is_some();
blocks.push(ContentBlock::ToolResult {
tool_use_id: call_id,
output_preview,
is_error: is_error_status(status) || has_error_field,
});
}
"file" => {
if let Some(file_ref) = extract_file_reference(&value) {
blocks.push(ContentBlock::Text {
text: format!("@{}", file_ref),
});
}
}
"patch" => {
let files = value
.get("files")
.and_then(|v| v.as_array())
.map(|arr| {
arr.iter()
.filter_map(|item| item.as_str())
.collect::<Vec<_>>()
})
.unwrap_or_default();
if !files.is_empty() {
blocks.push(ContentBlock::Text {
text: format!("Applied patch: {}", files.join(", ")),
});
}
}
"step-finish" => {
if usage_from_step_finish.is_none() {
usage_from_step_finish = value
.get("tokens")
.and_then(extract_opencode_usage_from_tokens);
}
}
_ => {}
}
}
Ok((blocks, usage_from_step_finish))
}
}
impl AgentParser for OpenCodeParser {
fn list_conversations(&self) -> Result<Vec<ConversationSummary>, ParseError> {
if !self.sqlite_db_path().exists() {
return Ok(Vec::new());
}
self.block_on(self.list_conversations_from_sqlite())
}
fn get_conversation(&self, conversation_id: &str) -> Result<ConversationDetail, ParseError> {
if !self.sqlite_db_path().exists() {
return Err(ParseError::ConversationNotFound(
conversation_id.to_string(),
));
}
self.block_on(self.get_conversation_from_sqlite(conversation_id))
}
}
fn resolve_opencode_base_dir() -> PathBuf {
resolve_xdg_data_home(std::env::var_os("XDG_DATA_HOME"), dirs::home_dir())
.map(|xdg_data_home| xdg_data_home.join("opencode"))
.unwrap_or_else(|| PathBuf::from("opencode"))
}
fn resolve_xdg_data_home(
xdg_data_home_env: Option<std::ffi::OsString>,
home_dir: Option<PathBuf>,
) -> Option<PathBuf> {
xdg_data_home_env
.filter(|value| !value.is_empty())
.map(PathBuf::from)
.or_else(|| home_dir.map(|home| home.join(".local").join("share")))
}
fn normalize_optional_string(value: Option<String>) -> Option<String> {
value.and_then(|s| {
let trimmed = s.trim();
if trimmed.is_empty() {
None
} else {
Some(trimmed.to_string())
}
})
}
fn value_to_preview(value: Option<&serde_json::Value>) -> Option<String> {
let v = value?;
if v.is_null() {
return None;
}
if let Some(s) = v.as_str() {
let trimmed = s.trim();
if trimmed.is_empty() {
None
} else {
Some(trimmed.to_string())
}
} else {
serde_json::to_string(v).ok()
}
}
fn extract_file_reference(value: &serde_json::Value) -> Option<String> {
value
.get("source")
.and_then(|s| s.get("path"))
.and_then(|v| v.as_str())
.or_else(|| value.get("filename").and_then(|v| v.as_str()))
.map(|s| s.trim())
.filter(|s| !s.is_empty())
.map(|s| s.to_string())
}
fn is_error_status(status: &str) -> bool {
matches!(
status.to_ascii_lowercase().as_str(),
"error" | "failed" | "failure" | "cancelled" | "canceled"
)
}
fn extract_opencode_usage(value: &serde_json::Value) -> Option<TurnUsage> {
value
.get("tokens")
.and_then(extract_opencode_usage_from_tokens)
}
fn extract_opencode_usage_from_tokens(tokens: &serde_json::Value) -> Option<TurnUsage> {
let input = tokens.get("input").and_then(|v| v.as_u64()).unwrap_or(0);
let output = tokens.get("output").and_then(|v| v.as_u64()).unwrap_or(0);
let cache = tokens.get("cache");
let cache_write = cache
.and_then(|c| c.get("write"))
.and_then(|v| v.as_u64())
.unwrap_or(0);
let cache_read = cache
.and_then(|c| c.get("read"))
.and_then(|v| v.as_u64())
.unwrap_or(0);
if input == 0 && output == 0 && cache_write == 0 && cache_read == 0 {
return None;
}
Some(TurnUsage {
input_tokens: input,
output_tokens: output,
cache_creation_input_tokens: cache_write,
cache_read_input_tokens: cache_read,
})
}
fn millis_to_datetime(ms: i64) -> DateTime<Utc> {
let secs = ms / 1000;
let nsecs = ((ms.rem_euclid(1000)) * 1_000_000) as u32;
Utc.timestamp_opt(secs, nsecs)
.single()
.unwrap_or_else(Utc::now)
}
/// Group flat messages into conversation turns (same strategy as Codex).
fn group_into_turns(messages: Vec<UnifiedMessage>) -> Vec<MessageTurn> {
let mut turns = Vec::new();
let mut i = 0;
while i < messages.len() {
let msg = &messages[i];
if matches!(msg.role, MessageRole::User) {
turns.push(MessageTurn {
id: format!("turn-{}", turns.len()),
role: TurnRole::User,
blocks: msg.content.clone(),
timestamp: msg.timestamp,
usage: None,
duration_ms: None,
model: None,
});
i += 1;
} else if matches!(msg.role, MessageRole::System) {
turns.push(MessageTurn {
id: format!("turn-{}", turns.len()),
role: TurnRole::System,
blocks: msg.content.clone(),
timestamp: msg.timestamp,
usage: None,
duration_ms: None,
model: None,
});
i += 1;
} else {
let mut blocks: Vec<ContentBlock> = msg.content.clone();
let mut usage = msg.usage.clone();
let mut duration_ms = msg.duration_ms;
let mut turn_model = msg.model.clone();
let timestamp = msg.timestamp;
i += 1;
while i < messages.len()
&& (matches!(messages[i].role, MessageRole::Assistant)
|| matches!(messages[i].role, MessageRole::Tool))
{
blocks.extend(messages[i].content.clone());
if usage.is_none() {
usage = messages[i].usage.clone();
}
if duration_ms.is_none() {
duration_ms = messages[i].duration_ms;
}
if turn_model.is_none() {
turn_model = messages[i].model.clone();
}
i += 1;
}
turns.push(MessageTurn {
id: format!("turn-{}", turns.len()),
role: TurnRole::Assistant,
blocks,
timestamp,
usage,
duration_ms,
model: turn_model,
});
}
}
turns
}
#[cfg(test)]
mod tests {
use super::resolve_xdg_data_home;
use std::path::PathBuf;
#[test]
fn xdg_data_home_env_overrides_home_fallback() {
let resolved = resolve_xdg_data_home(
Some(std::ffi::OsString::from("/tmp/xdg-data")),
Some(PathBuf::from("/Users/default")),
);
assert_eq!(resolved, Some(PathBuf::from("/tmp/xdg-data")));
}
#[test]
fn xdg_data_home_falls_back_to_home_local_share() {
let resolved = resolve_xdg_data_home(
None,
Some(PathBuf::from("/Users/default")),
);
assert_eq!(
resolved,
Some(PathBuf::from("/Users/default/.local/share"))
);
}
}