diff --git a/src/collector/codex.rs b/src/collector/codex.rs index bac213a..68d9cbe 100644 --- a/src/collector/codex.rs +++ b/src/collector/codex.rs @@ -1,7 +1,7 @@ use super::process::{self, ProcInfo}; use crate::model::{ - AgentSession, ChatMessage, ChatRole, ChildProcess, LaunchSurface, RateLimitInfo, SessionStatus, - ToolCall, MAX_CHAT_MESSAGES, + AgentSession, ChatMessage, ChatRole, ChildProcess, FileAccess, FileOp, LaunchSurface, + RateLimitInfo, SessionStatus, ToolCall, MAX_CHAT_MESSAGES, MAX_FILE_ACCESSES, }; use serde_json::Value; use std::collections::{HashMap, HashSet}; @@ -207,9 +207,7 @@ impl CodexCollector { Self::sort_rollouts_by_mtime_desc(&mut desktop_rollout_paths); for path in desktop_rollout_paths { - let pid = desktop_pid_for_path - .get(&path) - .copied(); + let pid = desktop_pid_for_path.get(&path).copied(); let process_ctx = CodexProcessContext { pid, is_exec: false, @@ -669,7 +667,7 @@ impl CodexCollector { tool_calls: result.tool_calls, pending_since_ms: result.pending_since_ms, thinking_since_ms: result.thinking_since_ms, - file_accesses: vec![], + file_accesses: result.file_accesses, config_root: super::abbrev_path( self.sessions_dir .parent() @@ -906,6 +904,9 @@ struct CodexJSONLResult { waiting_for_user: bool, /// Timestamp when the current model-thinking segment began. thinking_since_ms: u64, + /// Files touched, from the item_completed schema (Codex ≥ ~0.149). + /// The old schema never carried file information. + file_accesses: Vec, } impl CodexJSONLResult { @@ -946,6 +947,163 @@ fn sanitize_tool_arg(arg: &str) -> String { redacted.chars().take(120).collect() } +/// Joined text of an item's `content` array. The schema is not consistent +/// about casing ("Text" on AgentMessage, "text" on UserMessage), so any +/// object with a string `text` field counts. +fn concat_item_text(content: &Value) -> String { + content + .as_array() + .map(|parts| { + parts + .iter() + .filter_map(|p| p["text"].as_str()) + .collect::>() + .join("\n") + }) + .unwrap_or_default() +} + +/// One completed item from the new rollout schema (Codex ≥ ~0.149). +/// +/// UserMessage/AgentMessage carry the chat, CommandExecution and Extension +/// are the tool timeline, and FileChange is the only place file writes appear +/// (the old schema never reported files at all). +fn handle_item_completed(payload: &Value, ts: u64, result: &mut CodexJSONLResult) { + let item = &payload["item"]; + match item["type"].as_str() { + Some("UserMessage") => { + result.task_complete = false; + result.model_generating = true; + result.thinking_since_ms = ts; + let text = concat_item_text(&item["content"]); + if !text.is_empty() { + if result.initial_prompt.is_empty() { + result.initial_prompt = clean_chat_text(&text, 120); + } + push_chat_message( + &mut result.chat_messages, + ChatRole::User, + clean_chat_text(&text, 500), + ); + } + } + Some("AgentMessage") => { + result.turn_count += 1; + // Progress messages do not end a turn. New rollouts can signal + // completion with final_answer even without a task_complete event. + if item["phase"].as_str() == Some("final_answer") { + result.task_complete = true; + result.model_generating = false; + result.thinking_since_ms = 0; + } + let text = concat_item_text(&item["content"]); + push_chat_message( + &mut result.chat_messages, + ChatRole::Assistant, + clean_chat_text(&text, 500), + ); + } + Some("CommandExecution") => { + if result.model_generating { + result.thinking_since_ms = ts; + } + let started = item["started_at_ms"] + .as_u64() + .or_else(|| payload["started_at_ms"].as_u64()) + .unwrap_or(ts); + let completed = item["completed_at_ms"] + .as_u64() + .or_else(|| payload["completed_at_ms"].as_u64()) + .unwrap_or(started); + // Raw argv and parsed commands can contain scripts, file bodies, + // or credentials. Display only a known operation and a path. + let parsed = &item["parsed_cmd"][0]; + let name = match parsed["type"].as_str() { + Some("read") => "read", + Some("write") => "write", + Some("search") => "search", + _ => "exec", + }; + let arg = parsed["path"] + .as_str() + .map(clean_item_path) + .unwrap_or_default(); + if result.tool_calls.len() < 500 { + result.tool_calls.push(ToolCall { + name: name.to_string(), + arg, + duration_ms: completed.saturating_sub(started), + }); + } + if let Some(entries) = item["parsed_cmd"].as_array() { + for pc in entries { + let Some(path) = pc["path"].as_str() else { + continue; + }; + let op = match pc["type"].as_str() { + Some("write") => FileOp::Write, + _ => FileOp::Read, + }; + push_file_access(result, path, op); + } + } + } + Some("FileChange") => { + if result.model_generating { + result.thinking_since_ms = ts; + } + if let Some(changes) = item["changes"].as_object() { + for (path, change) in changes { + let op = match change["type"].as_str() { + Some("add") => FileOp::Write, + _ => FileOp::Edit, + }; + push_file_access(result, path, op); + if result.tool_calls.len() < 500 { + let short = process::last_path_segment(path).unwrap_or(path); + result.tool_calls.push(ToolCall { + name: "edit".to_string(), + arg: clean_item_path(short), + duration_ms: 0, + }); + } + } + } + } + Some("Extension") => { + if result.model_generating { + result.thinking_since_ms = ts; + } + // Extension queries are opaque content, like custom tool inputs. + if result.tool_calls.len() < 500 { + result.tool_calls.push(ToolCall { + name: "extension".to_string(), + arg: String::new(), + duration_ms: 0, + }); + } + } + _ => {} + } +} + +fn clean_item_path(path: &str) -> String { + let safe = super::sanitize_terminal_text(path); + super::redact_secrets(&safe).chars().take(512).collect() +} + +fn push_file_access(result: &mut CodexJSONLResult, path: &str, op: FileOp) { + result.file_accesses.push(FileAccess { + path: clean_item_path(path), + operation: op, + turn_index: result.turn_count, + }); + let len = result.file_accesses.len(); + if len > MAX_FILE_ACCESSES { + result.file_accesses.drain(..len - MAX_FILE_ACCESSES); + } +} + fn push_chat_message(messages: &mut Vec, role: ChatRole, text: String) { if text.is_empty() { return; @@ -1093,6 +1251,7 @@ fn parse_codex_jsonl(path: &Path) -> Option { pending_since_ms: 0, waiting_for_user: false, thinking_since_ms: 0, + file_accesses: Vec::new(), }; let mut call_indices: HashMap = HashMap::new(); let mut call_starts: HashMap = HashMap::new(); @@ -1179,6 +1338,18 @@ fn parse_codex_jsonl(path: &Path) -> Option { result.context_window = cw; } } + // Newer Codex versions replaced the + // user_message/agent_message/function_call vocabulary + // with item_completed wrappers. + Some("item_completed") => { + let ts = event_timestamp_ms(&val).unwrap_or(0); + handle_item_completed(payload, ts, &mut result); + } + Some("turn_aborted") => { + result.task_complete = true; + result.model_generating = false; + result.thinking_since_ms = 0; + } Some("user_message") => { result.task_complete = false; result.model_generating = true; @@ -1790,11 +1961,9 @@ mod tests { #[test] fn desktop_filesystem_only_rollout_is_unknown_without_fd_owner() { let sessions = tempfile::tempdir().unwrap(); - let today = sessions.path().join( - chrono::Local::now() - .format("%Y/%m/%d") - .to_string(), - ); + let today = sessions + .path() + .join(chrono::Local::now().format("%Y/%m/%d").to_string()); fs::create_dir_all(&today).unwrap(); let active = today.join("rollout-active.jsonl"); write_jsonl(&active, &[DESKTOP_SESSION_META]); @@ -2270,7 +2439,9 @@ mod tests { write_lines( &mut file, - &[r#"{"type":"response_item","timestamp":"2026-03-28T15:01:09Z","payload":{"type":"custom_tool_call_output","call_id":"private_call","output":"private-test-value"}}"#], + &[ + r#"{"type":"response_item","timestamp":"2026-03-28T15:01:09Z","payload":{"type":"custom_tool_call_output","call_id":"private_call","output":"private-test-value"}}"#, + ], ); let result = parse_codex_jsonl(file.path()).unwrap(); assert!(result.current_task.is_empty()); @@ -2585,4 +2756,236 @@ mod tests { let file = tempfile::NamedTempFile::new().unwrap(); assert!(parse_codex_jsonl(file.path()).is_none()); } + + #[test] + fn test_item_completed_chat_and_turns() { + // Codex ≥ ~0.149: chat arrives as item_completed wrappers. Note the + // casing drift — "text" on UserMessage content, "Text" on AgentMessage. + let mut file = tempfile::NamedTempFile::new().unwrap(); + write_lines( + &mut file, + &[ + SESSION_META, + r#"{"type":"event_msg","timestamp":"2026-08-21T05:45:27Z","payload":{"type":"item_completed","item":{"type":"UserMessage","id":"u1","content":[{"type":"text","text":"please fix the build"}]}}}"#, + r#"{"type":"event_msg","timestamp":"2026-08-21T05:45:54Z","payload":{"type":"item_completed","item":{"type":"AgentMessage","id":"a1","content":[{"type":"Text","text":"done"}],"phase":"final_answer"}}}"#, + ], + ); + let result = parse_codex_jsonl(file.path()).unwrap(); + assert_eq!(result.chat_messages.len(), 2); + assert_eq!(result.chat_messages[0].text, "please fix the build"); + assert_eq!(result.chat_messages[1].text, "done"); + assert_eq!(result.turn_count, 1); + assert!(!result.model_generating, "AgentMessage ends the turn"); + assert_eq!(result.initial_prompt, "please fix the build"); + } + + #[test] + fn test_item_completed_trailing_user_message_marks_generating() { + let mut file = tempfile::NamedTempFile::new().unwrap(); + write_lines( + &mut file, + &[ + SESSION_META, + r#"{"type":"event_msg","timestamp":"2026-08-21T05:45:27Z","payload":{"type":"item_completed","item":{"type":"UserMessage","id":"u1","content":[{"type":"text","text":"go"}]}}}"#, + ], + ); + let result = parse_codex_jsonl(file.path()).unwrap(); + assert!( + result.model_generating, + "unanswered UserMessage means the model is working" + ); + } + + #[test] + fn test_item_completed_command_execution_tools_and_files() { + let mut file = tempfile::NamedTempFile::new().unwrap(); + write_lines( + &mut file, + &[ + SESSION_META, + r#"{"type":"event_msg","timestamp":"2026-08-21T05:46:11Z","payload":{"type":"item_completed","item":{"type":"CommandExecution","id":"exec-1","command":["/bin/zsh","-lc","sed -n '1,240p' README.md"],"parsed_cmd":[{"type":"read","cmd":"sed -n '1,240p' README.md","name":"README.md","path":"/work/README.md"}],"started_at_ms":1787291167583,"completed_at_ms":1787291168442}}}"#, + ], + ); + let result = parse_codex_jsonl(file.path()).unwrap(); + assert_eq!(result.tool_calls.len(), 1); + assert_eq!(result.tool_calls[0].name, "read"); + assert_eq!(result.tool_calls[0].duration_ms, 859); + assert_eq!(result.file_accesses.len(), 1); + assert_eq!(result.file_accesses[0].path, "/work/README.md"); + assert!(matches!(result.file_accesses[0].operation, FileOp::Read)); + } + + #[test] + fn test_item_completed_file_change_records_accesses() { + let mut file = tempfile::NamedTempFile::new().unwrap(); + write_lines( + &mut file, + &[ + SESSION_META, + r#"{"type":"event_msg","timestamp":"2026-08-21T06:10:59Z","payload":{"type":"item_completed","item":{"type":"FileChange","id":"exec-2","changes":{"/work/new.md":{"type":"add"},"/work/old.rs":{"type":"update"}}}}}"#, + ], + ); + let result = parse_codex_jsonl(file.path()).unwrap(); + assert_eq!(result.file_accesses.len(), 2); + let write = result + .file_accesses + .iter() + .find(|f| f.path == "/work/new.md") + .unwrap(); + let edit = result + .file_accesses + .iter() + .find(|f| f.path == "/work/old.rs") + .unwrap(); + assert!(matches!(write.operation, FileOp::Write)); + assert!(matches!(edit.operation, FileOp::Edit)); + assert_eq!( + result.tool_calls.len(), + 2, + "each change also lands on the tool timeline" + ); + } + + #[test] + fn item_completed_preserves_active_turn_until_final_answer() { + let mut file = tempfile::NamedTempFile::new().unwrap(); + write_lines( + &mut file, + &[ + SESSION_META, + r#"{"type":"event_msg","timestamp":"2026-09-14T00:00:00Z","payload":{"type":"task_started"}}"#, + ], + ); + let collector = CodexCollector::new(); + let processes = HashMap::from([(42, proc_info(42, 1, "codex"))]); + for item in [ + serde_json::json!({"type":"AgentMessage","phase":"commentary","content":[{"type":"Text","text":"Working"}]}), + serde_json::json!({"type":"CommandExecution","command":["echo","done"]}), + serde_json::json!({"type":"FileChange","changes":{"/tmp/report":{"type":"add"}}}), + serde_json::json!({"type":"Extension","kind":"search","query":"private-query"}), + ] { + let event = serde_json::json!({"type":"event_msg","timestamp":"2026-09-14T00:00:01Z","payload":{"type":"item_completed","item":item}}); + write_lines(&mut file, &[&event.to_string()]); + let (session, _) = collector + .load_session_with_rate_limit( + owned_process(42), + file.path(), + &processes, + &HashMap::new(), + &HashMap::new(), + ) + .unwrap(); + assert_eq!(session.status, SessionStatus::Thinking); + } + write_lines( + &mut file, + &[ + r#"{"type":"event_msg","timestamp":"2026-09-14T00:00:02Z","payload":{"type":"item_completed","item":{"type":"AgentMessage","phase":"final_answer","content":[{"type":"Text","text":"Done"}]}}}"#, + ], + ); + let (session, _) = collector + .load_session_with_rate_limit( + owned_process(42), + file.path(), + &processes, + &HashMap::new(), + &HashMap::new(), + ) + .unwrap(); + assert_eq!(session.status, SessionStatus::Waiting); + // A new user item must clear the previous turn's completion flag. + write_lines( + &mut file, + &[ + r#"{"type":"event_msg","timestamp":"2026-09-14T00:00:03Z","payload":{"type":"item_completed","item":{"type":"UserMessage","content":[{"type":"text","text":"Continue"}]}}}"#, + ], + ); + let (session, _) = collector + .load_session_with_rate_limit( + owned_process(42), + file.path(), + &processes, + &HashMap::new(), + &HashMap::new(), + ) + .unwrap(); + assert_eq!(session.status, SessionStatus::Thinking); + write_lines( + &mut file, + &[ + r#"{"type":"event_msg","timestamp":"2026-09-14T00:00:04Z","payload":{"type":"turn_aborted"}}"#, + ], + ); + let (session, _) = collector + .load_session_with_rate_limit( + owned_process(42), + file.path(), + &processes, + &HashMap::new(), + &HashMap::new(), + ) + .unwrap(); + assert_eq!(session.status, SessionStatus::Waiting); + } + + #[test] + fn item_completed_omits_opaque_content_and_sanitizes_metadata() { + let mut file = tempfile::NamedTempFile::new().unwrap(); + write_lines(&mut file, &[SESSION_META]); + for item in [ + serde_json::json!({"type":"CommandExecution","command":["sh","-c","echo PRIVATE_BODY > .env"],"parsed_cmd":[{"type":"read","cmd":"PRIVATE_BODY","path":"/tmp/\u{1b}\u{202e}safe.txt"}]}), + serde_json::json!({"type":"CommandExecution","command":["echo PRIVATE_BODY"],"parsed_cmd":[{"type":"PRIVATE_BODY"}]}), + serde_json::json!({"type":"Extension","kind":"PRIVATE_BODY","query":"PRIVATE_BODY"}), + serde_json::json!({"type":"FileChange","changes":{"/tmp/\u{1b}\u{202e}safe.txt":{"type":"add","diff":"PRIVATE_BODY"}}}), + ] { + let event = serde_json::json!({"type":"event_msg","payload":{"type":"item_completed","item":item,"started_at_ms":100,"completed_at_ms":130}}); + write_lines(&mut file, &[&event.to_string()]); + } + let result = parse_codex_jsonl(file.path()).unwrap(); + assert_eq!(result.tool_calls.len(), 4); + assert_eq!(result.tool_calls[0].arg, "/tmp/safe.txt"); + assert_eq!(result.tool_calls[0].duration_ms, 30); + assert_eq!(result.tool_calls[1].name, "exec"); + assert!(result.tool_calls[1].arg.is_empty()); + assert_eq!(result.tool_calls[2].name, "extension"); + assert!(result.tool_calls[2].arg.is_empty()); + assert_eq!(result.tool_calls[3].arg, "safe.txt"); + for tool in &result.tool_calls { + assert!(!tool.name.contains("PRIVATE_BODY")); + assert!(!tool.arg.contains("PRIVATE_BODY")); + } + assert!(result + .file_accesses + .iter() + .all(|f| f.path == "/tmp/safe.txt")); + assert_eq!(clean_item_path("/tmp/ghp_synthetic"), "/tmp/[REDACTED]"); + assert_eq!(clean_item_path(&"x".repeat(1024)).len(), 512); + } + + #[test] + fn item_completed_tolerates_missing_fields_and_bounds_history() { + let mut file = tempfile::NamedTempFile::new().unwrap(); + write_lines(&mut file, &[SESSION_META]); + for _ in 0..(MAX_FILE_ACCESSES + 10) { + write_lines( + &mut file, + &[ + r#"{"type":"event_msg","payload":{"type":"item_completed","item":{"type":"FileChange","changes":{"/tmp/file":{"type":"add"}}}}}"#, + ], + ); + } + for item in [ + serde_json::Value::Null, + serde_json::json!({"type":"CommandExecution"}), + serde_json::json!({"type":"UnknownFutureItem"}), + ] { + let event = serde_json::json!({"type":"event_msg","payload":{"type":"item_completed","item":item}}); + write_lines(&mut file, &[&event.to_string()]); + } + let result = parse_codex_jsonl(file.path()).unwrap(); + assert_eq!(result.file_accesses.len(), MAX_FILE_ACCESSES); + assert!(result.tool_calls.len() <= 500); + // Once the timeline is full, later items must not grow it. + assert_eq!(result.tool_calls.len(), 500); + } }