use serde_json::Value;
use crate::chat::types::{ConversationItem, FileEdit, ItemDelta, Lifecycle, TurnUsage};
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub(super) enum ItemPhase {
Started,
Completed,
}
pub(super) fn extract_turn_id(params: &Value) -> Option<String> {
params
.get("turnId")
.and_then(Value::as_str)
.or_else(|| {
params
.get("turn")
.and_then(|t| t.get("id"))
.and_then(Value::as_str)
})
.map(ToString::to_string)
}
pub(super) fn extract_thread_id(value: &Value) -> Option<String> {
value
.get("thread")
.and_then(|t| t.get("id"))
.and_then(Value::as_str)
.map(ToString::to_string)
}
pub(super) fn map_turn_status(params: &Value) -> Lifecycle {
match params
.pointer("/turn/status")
.and_then(Value::as_str)
.unwrap_or_default()
{
"interrupted" => Lifecycle::Interrupted,
"failed" => Lifecycle::Failed,
_ => Lifecycle::Completed,
}
}
pub(super) fn map_token_usage(params: &Value) -> TurnUsage {
let total = params.pointer("/tokenUsage/total");
let field = |key: &str| {
total
.and_then(|t| t.get(key))
.and_then(Value::as_u64)
.unwrap_or(0)
};
TurnUsage {
input_tokens: field("inputTokens"),
output_tokens: field("outputTokens"),
reasoning_tokens: Some(field("reasoningOutputTokens")),
cache_read_tokens: Some(field("cachedInputTokens")),
cache_write_tokens: None,
model: None,
cost_usd: None,
}
}
pub(super) fn map_item_id(params: &Value) -> String {
params
.pointer("/item/id")
.and_then(Value::as_str)
.or_else(|| params.get("itemId").and_then(Value::as_str))
.unwrap_or("unknown")
.to_string()
}
pub(super) fn map_item_type(params: &Value) -> &str {
params
.pointer("/item/type")
.and_then(Value::as_str)
.unwrap_or("unknown")
}
pub(super) fn build_item(params: &Value, phase: ItemPhase) -> ConversationItem {
let id = map_item_id(params);
let item_type = map_item_type(params);
let item = params.get("item").unwrap_or(params);
let status = map_item_status(item);
let completed = phase == ItemPhase::Completed;
match item_type {
"commandExecution" => ConversationItem::Command {
id,
command: text_field(item, "command")
.map(|c| vec![c])
.unwrap_or_default(),
cwd: required_text_field(item, "cwd"),
status,
output: if completed {
text_field(item, "aggregatedOutput")
} else {
None
},
exit_code: if completed {
item.get("exitCode")
.and_then(Value::as_i64)
.map(|v| v as i32)
} else {
None
},
duration_ms: if completed {
item.get("durationMs").and_then(Value::as_u64)
} else {
None
},
},
"fileChange" => ConversationItem::File {
id,
changes: parse_file_changes(item),
status,
},
"mcpToolCall" => ConversationItem::Tool {
id,
name: mcp_tool_name(item),
status,
input: item.get("arguments").cloned(),
output: if completed {
mcp_tool_output(item)
} else {
None
},
},
"agentMessage" => ConversationItem::Message {
id,
text: required_text_field(item, "text"),
phase: text_field(item, "phase"),
},
"plan" => ConversationItem::Thought {
id,
text: required_text_field(item, "text"),
},
"reasoning" => ConversationItem::Thought {
id,
text: string_array(item, "summary").join("\n"),
},
_ => ConversationItem::Tool {
id,
name: item_type.to_string(),
status,
input: None,
output: None,
},
}
}
pub(super) fn map_item_delta(method: &str, params: &Value) -> Option<ItemDelta> {
let content = delta_content(params)?;
match method {
"item/commandExecution/outputDelta" | "item/fileChange/outputDelta" => {
Some(ItemDelta::Output { content })
}
"item/plan/delta" => Some(ItemDelta::PlanText { content }),
_ => None,
}
}
pub(super) fn delta_content(params: &Value) -> Option<String> {
params
.get("delta")
.and_then(Value::as_str)
.map(ToString::to_string)
}
fn text_field(value: &Value, key: &str) -> Option<String> {
value
.get(key)
.and_then(Value::as_str)
.map(ToString::to_string)
}
fn required_text_field(value: &Value, key: &str) -> String {
text_field(value, key).unwrap_or_default()
}
fn string_array(value: &Value, key: &str) -> Vec<String> {
value
.get(key)
.and_then(Value::as_array)
.map(|arr| {
arr.iter()
.filter_map(Value::as_str)
.map(ToString::to_string)
.collect()
})
.unwrap_or_default()
}
fn map_item_status(item: &Value) -> Lifecycle {
match item
.get("status")
.and_then(Value::as_str)
.unwrap_or("inProgress")
{
"completed" => Lifecycle::Completed,
"failed" => Lifecycle::Failed,
"declined" => Lifecycle::Failed,
_ => Lifecycle::Running,
}
}
fn mcp_tool_name(item: &Value) -> String {
let tool = required_text_field(item, "tool");
let server = required_text_field(item, "server");
match (server.is_empty(), tool.is_empty()) {
(true, true) => "mcp_tool_call".to_string(),
(true, false) => tool,
(false, true) => server,
(false, false) => format!("{server}/{tool}"),
}
}
fn mcp_tool_output(item: &Value) -> Option<String> {
if let Some(result) = item.get("result").filter(|v| !v.is_null()) {
let texts: Vec<&str> = result
.get("content")
.and_then(Value::as_array)
.map(|blocks| {
blocks
.iter()
.filter_map(|b| b.get("text").and_then(Value::as_str))
.collect()
})
.unwrap_or_default();
if texts.is_empty() {
return Some(result.to_string());
}
return Some(texts.join("\n"));
}
item.pointer("/error/message")
.and_then(Value::as_str)
.map(ToString::to_string)
}
fn parse_file_changes(item: &Value) -> Vec<FileEdit> {
item.get("changes")
.and_then(Value::as_array)
.map(|arr| {
arr.iter()
.map(|c| FileEdit {
path: c
.get("path")
.and_then(Value::as_str)
.unwrap_or_default()
.to_string(),
kind: c
.pointer("/kind/type")
.and_then(Value::as_str)
.map(ToString::to_string),
diff: c
.get("diff")
.and_then(Value::as_str)
.map(ToString::to_string),
})
.collect()
})
.unwrap_or_default()
}
#[cfg(test)]
mod tests {
use serde_json::json;
use super::*;
#[test]
fn build_item_maps_file_change_to_file() {
let params = json!({
"item": {
"id": "item_1",
"type": "fileChange",
"status": "completed",
"changes": [
{"path": "src/main.rs", "kind": {"type": "update", "move_path": null}, "diff": "-a\n+b"}
]
}
});
let item = build_item(¶ms, ItemPhase::Completed);
match item {
ConversationItem::File {
id,
changes,
status,
} => {
assert_eq!(id, "item_1");
assert_eq!(status, Lifecycle::Completed);
assert_eq!(changes.len(), 1);
assert_eq!(changes[0].path, "src/main.rs");
assert_eq!(changes[0].kind.as_deref(), Some("update"));
}
other => panic!("expected file item, got {other:?}"),
}
}
#[test]
fn build_item_maps_agent_message_to_message() {
let params = json!({
"item": {
"id": "msg_1",
"type": "agentMessage",
"text": "Done",
"phase": "final_answer",
"memoryCitation": null
}
});
let item = build_item(¶ms, ItemPhase::Completed);
match item {
ConversationItem::Message { id, text, phase } => {
assert_eq!(id, "msg_1");
assert_eq!(text, "Done");
assert_eq!(phase.as_deref(), Some("final_answer"));
}
other => panic!("expected message item, got {other:?}"),
}
}
#[test]
fn build_item_maps_command_string_to_single_element_argv() {
let params = json!({
"item": {
"id": "cmd_1",
"type": "commandExecution",
"command": "ls -la",
"commandActions": [],
"cwd": "/tmp",
"status": "inProgress"
}
});
let item = build_item(¶ms, ItemPhase::Started);
match item {
ConversationItem::Command {
command,
cwd,
status,
..
} => {
assert_eq!(command, vec!["ls -la".to_string()]);
assert_eq!(cwd, "/tmp");
assert_eq!(status, Lifecycle::Running);
}
other => panic!("expected command item, got {other:?}"),
}
}
#[test]
fn build_item_maps_reasoning_to_thought() {
let params = json!({
"item": {
"id": "rsn_1",
"type": "reasoning",
"summary": ["Check tests", "Then land"]
}
});
let item = build_item(¶ms, ItemPhase::Completed);
match item {
ConversationItem::Thought { id, text } => {
assert_eq!(id, "rsn_1");
assert_eq!(text, "Check tests\nThen land");
}
other => panic!("expected thought item, got {other:?}"),
}
}
#[test]
fn build_item_maps_mcp_tool_call_to_generic_tool() {
let params = json!({
"item": {
"id": "item_4",
"type": "mcpToolCall",
"status": "completed",
"server": "github",
"tool": "search",
"arguments": { "query": "regression" },
"result": { "content": [{"type": "text", "text": "ok"}] }
}
});
let item = build_item(¶ms, ItemPhase::Completed);
match item {
ConversationItem::Tool {
id,
name,
status,
input,
output,
} => {
assert_eq!(id, "item_4");
assert_eq!(name, "github/search");
assert_eq!(status, Lifecycle::Completed);
assert_eq!(input, Some(json!({ "query": "regression" })));
assert_eq!(output.as_deref(), Some("ok"));
}
other => panic!("expected tool item, got {other:?}"),
}
}
#[test]
fn map_token_usage_reads_total_breakdown() {
let usage = map_token_usage(&json!({
"threadId": "t",
"turnId": "u",
"tokenUsage": {
"total": {
"totalTokens": 16070,
"inputTokens": 16065,
"cachedInputTokens": 9600,
"outputTokens": 5,
"reasoningOutputTokens": 0
},
"last": {
"totalTokens": 16070,
"inputTokens": 16065,
"cachedInputTokens": 9600,
"outputTokens": 5,
"reasoningOutputTokens": 0
}
}
}));
assert_eq!(usage.input_tokens, 16065);
assert_eq!(usage.output_tokens, 5);
assert_eq!(usage.reasoning_tokens, Some(0));
assert_eq!(usage.cache_read_tokens, Some(9600));
}
}