403 lines
14 KiB
Rust
403 lines
14 KiB
Rust
//! Codex app-server notifications into the common event model.
|
|
//!
|
|
//! The CLI promises JSONL but deliberately leaves room for new item kinds. We
|
|
//! therefore match only the records that have a useful common equivalent and
|
|
//! ignore the rest; an added Codex item must not make a live session go deaf.
|
|
|
|
use std::collections::HashSet;
|
|
|
|
use serde_json::{Value, json};
|
|
|
|
use super::super::driver::{Event, SessionStatus};
|
|
|
|
#[derive(Default)]
|
|
pub(super) struct Translator {
|
|
pub(super) thread_id: Option<String>,
|
|
completed: bool,
|
|
limited: bool,
|
|
streamed_messages: HashSet<String>,
|
|
pending_usage: Option<u64>,
|
|
}
|
|
|
|
impl Translator {
|
|
pub(super) fn translate(&mut self, line: &Value) -> Vec<Event> {
|
|
let method = line.get("method").and_then(Value::as_str);
|
|
let body = method.and_then(|_| line.get("params")).unwrap_or(line);
|
|
let kind = method.or_else(|| line.get("type").and_then(Value::as_str));
|
|
match kind {
|
|
Some("thread.started") => {
|
|
self.thread_id = line
|
|
.get("thread_id")
|
|
.and_then(Value::as_str)
|
|
.map(str::to_string);
|
|
Vec::new()
|
|
}
|
|
Some("thread/started") => {
|
|
self.thread_id = body
|
|
.pointer("/thread/id")
|
|
.and_then(Value::as_str)
|
|
.map(str::to_string);
|
|
Vec::new()
|
|
}
|
|
Some("turn.started") | Some("turn/started") => {
|
|
self.completed = false;
|
|
self.limited = false;
|
|
self.pending_usage = None;
|
|
self.streamed_messages.clear();
|
|
vec![Event::Status {
|
|
state: SessionStatus::Running,
|
|
}]
|
|
}
|
|
Some("item.started") | Some("item/started") => start_item(&body["item"]),
|
|
Some("item.updated") => update_item(&line["item"]),
|
|
Some("item.completed") | Some("item/completed") => {
|
|
complete_item(&body["item"], &self.streamed_messages)
|
|
}
|
|
Some("item/agentMessage/delta") => {
|
|
let Some(delta) = body.get("delta").and_then(Value::as_str) else {
|
|
return Vec::new();
|
|
};
|
|
if let Some(id) = body.get("itemId").and_then(Value::as_str) {
|
|
self.streamed_messages.insert(id.to_string());
|
|
}
|
|
vec![Event::AssistantText {
|
|
delta: delta.to_string(),
|
|
}]
|
|
}
|
|
Some("item/commandExecution/outputDelta") => {
|
|
let (Some(id), Some(output)) = (
|
|
body.get("itemId").and_then(Value::as_str),
|
|
body.get("delta").and_then(Value::as_str),
|
|
) else {
|
|
return Vec::new();
|
|
};
|
|
vec![Event::ToolUpdate {
|
|
id: id.to_string(),
|
|
output: output.to_string(),
|
|
}]
|
|
}
|
|
Some("thread/tokenUsage/updated") => {
|
|
self.pending_usage = body
|
|
.pointer("/tokenUsage/last/totalTokens")
|
|
.and_then(Value::as_u64);
|
|
Vec::new()
|
|
}
|
|
Some("turn.completed") | Some("turn/completed") => {
|
|
self.completed = true;
|
|
let mut events = Vec::new();
|
|
if let Some(tokens) = self.pending_usage.take() {
|
|
events.push(Event::UsageDelta {
|
|
tokens,
|
|
context: None,
|
|
});
|
|
} else if let Some(usage) = line.get("usage") {
|
|
let input = number(usage, "input_tokens");
|
|
let output = number(usage, "output_tokens");
|
|
if input.is_some() || output.is_some() {
|
|
events.push(Event::UsageDelta {
|
|
tokens: input.unwrap_or(0) + output.unwrap_or(0),
|
|
// `exec` reports the sum across every model call in
|
|
// a turn, not the final call's context.
|
|
context: None,
|
|
});
|
|
}
|
|
}
|
|
if let Some(error) = body.pointer("/turn/error")
|
|
&& !error.is_null()
|
|
{
|
|
events.extend(self.failure(error));
|
|
}
|
|
events.push(Event::Status {
|
|
state: SessionStatus::Idle,
|
|
});
|
|
events
|
|
}
|
|
Some("turn.failed") | Some("error") => self.failure(body),
|
|
_ => Vec::new(),
|
|
}
|
|
}
|
|
|
|
#[cfg(test)]
|
|
pub(super) fn completed(&self) -> bool {
|
|
self.completed
|
|
}
|
|
|
|
#[cfg(test)]
|
|
pub(super) fn limited(&self) -> bool {
|
|
self.limited
|
|
}
|
|
|
|
fn failure(&mut self, line: &Value) -> Vec<Event> {
|
|
let detail = error_text(line);
|
|
if is_limit(&detail) {
|
|
if self.limited {
|
|
return Vec::new();
|
|
}
|
|
self.limited = true;
|
|
vec![Event::LimitReached {
|
|
resets_at: find_reset(line),
|
|
}]
|
|
} else {
|
|
vec![Event::Error { message: detail }]
|
|
}
|
|
}
|
|
}
|
|
|
|
fn start_item(item: &Value) -> Vec<Event> {
|
|
let Some((id, tool, input)) = tool(item) else {
|
|
return Vec::new();
|
|
};
|
|
vec![Event::ToolStart { id, tool, input }]
|
|
}
|
|
|
|
fn update_item(item: &Value) -> Vec<Event> {
|
|
let Some(id) = item.get("id").and_then(Value::as_str) else {
|
|
return Vec::new();
|
|
};
|
|
let output = item
|
|
.get("aggregated_output")
|
|
.or_else(|| item.get("output"))
|
|
.and_then(value_text);
|
|
output
|
|
.filter(|text| !text.is_empty())
|
|
.map(|output| {
|
|
vec![Event::ToolUpdate {
|
|
id: id.to_string(),
|
|
output,
|
|
}]
|
|
})
|
|
.unwrap_or_default()
|
|
}
|
|
|
|
fn complete_item(item: &Value, streamed_messages: &HashSet<String>) -> Vec<Event> {
|
|
match item.get("type").and_then(Value::as_str) {
|
|
Some("agent_message" | "agentMessage") => item
|
|
.get("text")
|
|
.and_then(Value::as_str)
|
|
.filter(|text| !text.is_empty())
|
|
.filter(|_| {
|
|
item.get("id")
|
|
.and_then(Value::as_str)
|
|
.is_none_or(|id| !streamed_messages.contains(id))
|
|
})
|
|
.map(|delta| {
|
|
vec![Event::AssistantText {
|
|
delta: delta.to_string(),
|
|
}]
|
|
})
|
|
.unwrap_or_default(),
|
|
Some("reasoning" | "userMessage") => Vec::new(),
|
|
_ => {
|
|
let Some((id, _, _)) = tool(item) else {
|
|
return Vec::new();
|
|
};
|
|
let output = tool_output(item);
|
|
vec![Event::ToolEnd { id, output }]
|
|
}
|
|
}
|
|
}
|
|
|
|
fn tool(item: &Value) -> Option<(String, String, Value)> {
|
|
let id = item.get("id")?.as_str()?.to_string();
|
|
let kind = item.get("type")?.as_str()?;
|
|
let (name, input) = match kind {
|
|
"command_execution" | "commandExecution" => (
|
|
"exec_command".to_string(),
|
|
json!({"command": item.get("command").cloned().unwrap_or(Value::Null)}),
|
|
),
|
|
"file_change" | "fileChange" => (
|
|
"apply_patch".to_string(),
|
|
item.get("changes").cloned().unwrap_or(Value::Null),
|
|
),
|
|
"mcp_tool_call" | "mcpToolCall" => (
|
|
item.get("tool")
|
|
.or_else(|| item.get("name"))
|
|
.and_then(Value::as_str)
|
|
.map(|tool| format!("mcp:{tool}"))
|
|
.unwrap_or_else(|| "mcp".to_string()),
|
|
item.get("arguments").cloned().unwrap_or(Value::Null),
|
|
),
|
|
"web_search" | "webSearch" => (
|
|
"web_search".to_string(),
|
|
json!({"query": item.get("query").cloned().unwrap_or(Value::Null)}),
|
|
),
|
|
"todo_list" | "todoList" | "plan" => ("update_plan".to_string(), item.clone()),
|
|
_ => return None,
|
|
};
|
|
Some((id, name, input))
|
|
}
|
|
|
|
fn tool_output(item: &Value) -> String {
|
|
for key in [
|
|
"aggregated_output",
|
|
"aggregatedOutput",
|
|
"output",
|
|
"result",
|
|
"error",
|
|
] {
|
|
if let Some(text) = item.get(key).and_then(value_text)
|
|
&& !text.is_empty()
|
|
{
|
|
return text;
|
|
}
|
|
}
|
|
match item.get("status").and_then(Value::as_str) {
|
|
Some(status) => status.to_string(),
|
|
None => String::new(),
|
|
}
|
|
}
|
|
|
|
fn value_text(value: &Value) -> Option<String> {
|
|
value
|
|
.as_str()
|
|
.map(str::to_string)
|
|
.or_else(|| (!value.is_null()).then(|| value.to_string()))
|
|
}
|
|
|
|
fn number(value: &Value, key: &str) -> Option<u64> {
|
|
value.get(key).and_then(Value::as_u64)
|
|
}
|
|
|
|
fn error_text(line: &Value) -> String {
|
|
line.get("message")
|
|
.and_then(Value::as_str)
|
|
.or_else(|| line.pointer("/error/message").and_then(Value::as_str))
|
|
.or_else(|| line.get("error").and_then(Value::as_str))
|
|
.unwrap_or("Codex ended the turn with an unknown error")
|
|
.to_string()
|
|
}
|
|
|
|
fn is_limit(detail: &str) -> bool {
|
|
let lower = detail.to_ascii_lowercase();
|
|
lower.contains("usage limit")
|
|
|| lower.contains("rate limit")
|
|
|| lower.contains("quota exceeded")
|
|
|| lower.contains("credits depleted")
|
|
}
|
|
|
|
fn find_reset(value: &Value) -> Option<f64> {
|
|
for key in ["resets_at", "resetsAt", "reset_at", "resetAt"] {
|
|
if let Some(at) = value.get(key).and_then(Value::as_f64) {
|
|
return Some(at);
|
|
}
|
|
}
|
|
value.as_object()?.values().find_map(find_reset)
|
|
}
|
|
|
|
#[cfg(test)]
|
|
mod tests {
|
|
use super::*;
|
|
|
|
fn line(text: &str) -> Value {
|
|
serde_json::from_str(text).expect("fixture")
|
|
}
|
|
|
|
#[test]
|
|
fn translates_the_observed_minimal_stream() {
|
|
let mut translator = Translator::default();
|
|
assert!(
|
|
translator
|
|
.translate(&line(r#"{"type":"thread.started","thread_id":"thread-1"}"#))
|
|
.is_empty()
|
|
);
|
|
assert_eq!(translator.thread_id.as_deref(), Some("thread-1"));
|
|
assert_eq!(
|
|
translator.translate(&line(
|
|
r#"{"type":"item.completed","item":{"id":"item_0","type":"agent_message","text":"hello"}}"#
|
|
)),
|
|
vec![Event::AssistantText {
|
|
delta: "hello".to_string()
|
|
}]
|
|
);
|
|
let events = translator.translate(&line(
|
|
r#"{"type":"turn.completed","usage":{"input_tokens":13,"cached_input_tokens":8,"output_tokens":5}}"#,
|
|
));
|
|
assert_eq!(
|
|
events[0],
|
|
Event::UsageDelta {
|
|
tokens: 18,
|
|
context: None
|
|
}
|
|
);
|
|
assert!(translator.completed());
|
|
}
|
|
|
|
#[test]
|
|
fn translates_tools_and_limits_without_matching_whole_records() {
|
|
let mut translator = Translator::default();
|
|
let started = translator.translate(&line(
|
|
r#"{"type":"item.started","item":{"id":"item_1","type":"command_execution","command":"pwd","status":"in_progress"}}"#,
|
|
));
|
|
assert!(matches!(&started[0], Event::ToolStart { tool, .. } if tool == "exec_command"));
|
|
let ended = translator.translate(&line(
|
|
r#"{"type":"item.completed","item":{"id":"item_1","type":"command_execution","command":"pwd","aggregated_output":"/tmp\n","exit_code":0,"status":"completed"}}"#,
|
|
));
|
|
assert_eq!(
|
|
ended,
|
|
vec![Event::ToolEnd {
|
|
id: "item_1".to_string(),
|
|
output: "/tmp\n".to_string()
|
|
}]
|
|
);
|
|
let limit = translator.translate(&line(
|
|
r#"{"type":"turn.failed","error":{"message":"usage limit reached","resetsAt":1234}}"#,
|
|
));
|
|
assert_eq!(
|
|
limit,
|
|
vec![Event::LimitReached {
|
|
resets_at: Some(1234.0)
|
|
}]
|
|
);
|
|
assert!(translator.limited());
|
|
}
|
|
|
|
#[test]
|
|
fn translates_native_app_server_streaming_without_repeating_the_final_item() {
|
|
let mut translator = Translator::default();
|
|
assert_eq!(
|
|
translator.translate(&line(
|
|
r#"{"method":"turn/started","params":{"threadId":"thread-1","turn":{"id":"turn-1"}}}"#
|
|
)),
|
|
vec![Event::Status {
|
|
state: SessionStatus::Running
|
|
}]
|
|
);
|
|
assert_eq!(
|
|
translator.translate(&line(
|
|
r#"{"method":"item/agentMessage/delta","params":{"itemId":"message-1","delta":"hello"}}"#
|
|
)),
|
|
vec![Event::AssistantText {
|
|
delta: "hello".to_string()
|
|
}]
|
|
);
|
|
assert!(
|
|
translator
|
|
.translate(&line(
|
|
r#"{"method":"item/completed","params":{"item":{"id":"message-1","type":"agentMessage","text":"hello"}}}"#
|
|
))
|
|
.is_empty()
|
|
);
|
|
assert!(
|
|
translator
|
|
.translate(&line(
|
|
r#"{"method":"thread/tokenUsage/updated","params":{"tokenUsage":{"last":{"totalTokens":42}}}}"#
|
|
))
|
|
.is_empty()
|
|
);
|
|
assert_eq!(
|
|
translator.translate(&line(
|
|
r#"{"method":"turn/completed","params":{"turn":{"id":"turn-1","status":"completed","error":null}}}"#
|
|
)),
|
|
vec![
|
|
Event::UsageDelta {
|
|
tokens: 42,
|
|
context: None
|
|
},
|
|
Event::Status {
|
|
state: SessionStatus::Idle
|
|
}
|
|
]
|
|
);
|
|
}
|
|
}
|