diff --git a/bt-daemon/src/translate/codex.rs b/bt-daemon/src/translate/codex.rs index 7f06b2c..f554931 100644 --- a/bt-daemon/src/translate/codex.rs +++ b/bt-daemon/src/translate/codex.rs @@ -26,6 +26,8 @@ use super::{ use crate::ids; use crate::wire::Envelope; use regex::Regex; +use serde::de::DeserializeOwned; +use serde::Deserialize; use serde_json::{json, Map, Value}; use std::collections::HashMap; use std::path::Path; @@ -37,6 +39,69 @@ const MISSING_TOOL_OUTPUT_ERROR: &str = "Tool output missing before turn ended"; /// Keep a single translator batch small even when the unread suffix is large. const CATCH_UP_BYTE_BUDGET: usize = 64 * 1024; +// Codex hooks are triggers for the rollout reader. These partial types cover +// the hook-owned correlation fields while preserving raw rollout JSON for the +// separately bounded transcript reducer. +#[derive(Deserialize)] +struct SessionStartHook { + #[serde(default)] + source: Option, + #[serde(default)] + permission_mode: Option, +} + +#[derive(Deserialize)] +struct CompactHook { + turn_id: String, + #[serde(default)] + trigger: Option, +} + +#[derive(Deserialize)] +struct SubagentStartHook { + agent_id: String, + #[serde(default)] + agent_type: Option, +} + +/// The stable outer shape of a rollout JSONL row. The payload deliberately +/// stays untyped here: Codex adds response-item fields regularly and several +/// handlers preserve those native objects in span input/output. +#[derive(Deserialize)] +struct RolloutRecord { + #[serde(default, deserialize_with = "deserialize_timestamp")] + timestamp_ms: Option, + #[serde(rename = "type", default)] + kind: Option, + #[serde(default)] + payload: Value, +} + +/// Event and response rows use a second discriminator inside their payload. +#[derive(Deserialize)] +struct PayloadKind { + #[serde(rename = "type", default)] + kind: Option, +} + +#[derive(Deserialize)] +struct WorkingDirectoryPayload { + #[serde(default)] + cwd: Option, +} + +#[derive(Deserialize)] +struct TurnContextPayload { + #[serde(default)] + model: Option, + #[serde(default)] + turn_id: Option, +} + +fn decode(value: &Value) -> Option { + serde_json::from_value(value.clone()).ok() +} + pub struct CodexTranslatorFactory { git: Arc, } @@ -200,11 +265,21 @@ impl AgentTranslator for CodexTranslator { // --- hook-specific side effects (before catch-up) --- match event.event.as_str() { "SessionStart" => { - self.session_source = str_field(payload, "source"); - self.permission_mode = str_field(payload, "permission_mode"); + if let Some(hook) = decode::(payload) { + self.session_source = hook.source; + self.permission_mode = hook.permission_mode; + } + } + "SubagentStart" => { + if let Some(hook) = decode::(payload) { + self.handle_subagent_start(event, hook); + } + } + "PreCompact" | "PostCompact" => { + if let Some(hook) = decode::(payload) { + self.record_compaction_trigger(hook, &mut ops); + } } - "SubagentStart" => self.handle_subagent_start(event), - "PreCompact" | "PostCompact" => self.record_compaction_trigger(payload, &mut ops), _ => {} } @@ -343,11 +418,9 @@ impl CodexTranslator { ); } - fn record_compaction_trigger(&mut self, payload: &Value, ops: &mut Vec) { - let Some(turn_id) = str_field(payload, "turn_id") else { - return; - }; - let trigger = str_field(payload, "trigger").unwrap_or_else(|| "manual".to_string()); + fn record_compaction_trigger(&mut self, hook: CompactHook, ops: &mut Vec) { + let turn_id = hook.turn_id; + let trigger = hook.trigger.unwrap_or_else(|| "manual".to_string()); self.compaction_trigger_by_turn .insert(turn_id.clone(), trigger.clone()); // Back-fill onto an already-built compaction span. @@ -383,14 +456,11 @@ impl CodexTranslator { } } - fn handle_subagent_start(&mut self, event: &Envelope) { - let payload = &event.payload; - let (Some(agent_id), Some(path)) = ( - str_field(payload, "agent_id"), - effective_transcript_path(event), - ) else { + fn handle_subagent_start(&mut self, event: &Envelope, hook: SubagentStartHook) { + let Some(path) = effective_transcript_path(event) else { return; }; + let agent_id = hook.agent_id; if self.scopes.contains_key(&path) { return; } @@ -401,7 +471,11 @@ impl CodexTranslator { let subagent_root = ids::span_id(&self.session_id, &format!("subagent:{agent_id}")); let mut scope = Scope::new(&path, ScopeKind::Subagent, subagent_root); scope.agent_id = Some(agent_id); - scope.agent_type = str_field(payload, "agent_type"); + scope.agent_type = hook + .agent_type + .as_ref() + .and_then(Value::as_str) + .map(str::to_owned); scope.spawning_turn_span_id = Some(parent); self.scopes.insert(path, scope); } @@ -475,7 +549,7 @@ impl CodexTranslator { CATCH_UP_BYTE_BUDGET, ); for line in read.lines { - if let Ok(rec) = serde_json::from_str::(&line) { + if let Ok(rec) = serde_json::from_str::(&line) { self.process_record(&mut scope, &rec, hook_ts, ops); } } @@ -486,81 +560,87 @@ impl CodexTranslator { fn process_record( &mut self, scope: &mut Scope, - rec: &Value, + rec: &RolloutRecord, hook_ts: i64, ops: &mut Vec, ) { let op_start = ops.len(); - let ts = parse_ts(rec).unwrap_or(hook_ts); - let ty = rec.get("type").and_then(Value::as_str).unwrap_or(""); - let payload = rec.get("payload").cloned().unwrap_or(Value::Null); - if matches!(ty, "session_meta" | "turn_context") { - if let Some(cwd) = str_field(&payload, "cwd") { + let ts = rec.timestamp_ms.unwrap_or(hook_ts); + let kind = rec.kind.as_deref().unwrap_or(""); + let payload = &rec.payload; + if matches!(kind, "session_meta" | "turn_context") { + if let Some(cwd) = decode::(payload).and_then(|v| v.cwd) { scope.current_cwd = Some(cwd); } } - match ty { - "session_meta" => self.open_root(scope, &payload, ts, ops), + match kind { + "session_meta" => self.open_root(scope, payload, ts, ops), "turn_context" => { - if let Some(m) = str_field(&payload, "model") { - let model_turn_id = str_field(&payload, "turn_id"); - scope.model = Some(m.clone()); - if scope.root_created { - let input = if scope.kind == ScopeKind::Main { - json!({ - "model": m, - "cwd": self.root_cwd, - "source": self.session_source, - }) - } else { - json!({ "model": m }) - }; - ops.push(SpanOp::Merge(SpanRow { - span_id: scope.turn_parent_span_id.clone(), - root_span_id: self.root_span_id.clone(), - input: Some(input), - metadata: Some(json!({ "model": m })), - ..Default::default() - })); - } - if let Some(turn) = model_turn_id.as_deref().and_then(|turn_id| { - scope.open_turns.iter().find(|turn| turn.turn_id == turn_id) - }) { - ops.push(SpanOp::Merge(SpanRow { - span_id: turn.span_id.clone(), - root_span_id: self.root_span_id.clone(), - metadata: Some(json!({ "model": m })), - ..Default::default() - })); + if let Some(context) = decode::(payload) { + if let Some(m) = context.model { + let model_turn_id = context.turn_id; + scope.model = Some(m.clone()); + if scope.root_created { + let input = if scope.kind == ScopeKind::Main { + json!({ + "model": m, + "cwd": self.root_cwd, + "source": self.session_source, + }) + } else { + json!({ "model": m }) + }; + ops.push(SpanOp::Merge(SpanRow { + span_id: scope.turn_parent_span_id.clone(), + root_span_id: self.root_span_id.clone(), + input: Some(input), + metadata: Some(json!({ "model": m })), + ..Default::default() + })); + } + if let Some(turn) = model_turn_id.as_deref().and_then(|turn_id| { + scope.open_turns.iter().find(|turn| turn.turn_id == turn_id) + }) { + ops.push(SpanOp::Merge(SpanRow { + span_id: turn.span_id.clone(), + root_span_id: self.root_span_id.clone(), + metadata: Some(json!({ "model": m })), + ..Default::default() + })); + } } } } "event_msg" => { - let sub = payload.get("type").and_then(Value::as_str).unwrap_or(""); - match sub { - "task_started" => self.open_turn(scope, &payload, ts, ops), - "user_message" => self.set_turn_input(scope, &payload, ops), - "token_count" => self.close_llm_with_tokens(scope, &payload, ts, ops), - "task_complete" => self.close_turn(scope, &payload, ts, ops), + let sub = decode::(payload) + .and_then(|v| v.kind) + .unwrap_or_default(); + match sub.as_str() { + "task_started" => self.open_turn(scope, payload, ts, ops), + "user_message" => self.set_turn_input(scope, payload, ops), + "token_count" => self.close_llm_with_tokens(scope, payload, ts, ops), + "task_complete" => self.close_turn(scope, payload, ts, ops), _ => {} } } "response_item" => { - let sub = payload.get("type").and_then(Value::as_str).unwrap_or(""); - match sub { - "message" => self.on_message(scope, &payload, ts, ops), - "reasoning" => self.on_reasoning(scope, &payload, ts, ops), + let sub = decode::(payload) + .and_then(|v| v.kind) + .unwrap_or_default(); + match sub.as_str() { + "message" => self.on_message(scope, payload, ts, ops), + "reasoning" => self.on_reasoning(scope, payload, ts, ops), "function_call" | "custom_tool_call" | "tool_search_call" => { - self.on_tool_call(scope, &payload, ts, ops) + self.on_tool_call(scope, payload, ts, ops) } "function_call_output" | "custom_tool_call_output" | "tool_search_output" => { - self.on_tool_output(scope, &payload, ts, ops) + self.on_tool_output(scope, payload, ts, ops) } _ => {} } } - "compacted" => self.on_compacted(scope, rec, &payload, ts, ops), + "compacted" => self.on_compacted(scope, payload, ts, ops), _ => {} } let cwd = scope.current_cwd.as_deref().or(self.root_cwd.as_deref()); @@ -1109,14 +1189,7 @@ impl CodexTranslator { self.scopes.insert(path, scope); } - fn on_compacted( - &mut self, - scope: &mut Scope, - _rec: &Value, - payload: &Value, - ts: i64, - ops: &mut Vec, - ) { + fn on_compacted(&mut self, scope: &mut Scope, payload: &Value, ts: i64, ops: &mut Vec) { let Some(turn) = scope.open_turns.last() else { return; }; @@ -1685,8 +1758,84 @@ fn compaction_history(payload: &Value, replacement: Option<&Vec>) -> Vec< } fn parse_ts(rec: &Value) -> Option { - let s = rec.get("timestamp").and_then(Value::as_str)?; - chrono::DateTime::parse_from_rfc3339(s) + rec.get("timestamp").and_then(normalize_timestamp) +} + +/// Deserialize Codex rollout timestamps into Unix milliseconds. Rollouts have +/// historically used RFC 3339 strings, but imported or newer producers may +/// supply numeric epoch seconds, milliseconds, microseconds, or nanoseconds. +/// Unknown values remain absent so the caller can use the hook timestamp. +fn deserialize_timestamp<'de, D>(deserializer: D) -> Result, D::Error> +where + D: serde::Deserializer<'de>, +{ + Ok(Option::::deserialize(deserializer)? + .as_ref() + .and_then(normalize_timestamp)) +} + +fn normalize_timestamp(value: &Value) -> Option { + match value { + Value::String(value) => { + parse_rfc3339_timestamp(value).or_else(|| normalize_numeric_timestamp(value)) + } + Value::Number(value) => { + if let Some(value) = value.as_i64() { + normalize_epoch_integer(value) + } else if let Some(value) = value.as_u64() { + i64::try_from(value).ok().and_then(normalize_epoch_integer) + } else { + value.as_f64().and_then(normalize_epoch_float) + } + } + _ => None, + } +} + +fn normalize_numeric_timestamp(value: &str) -> Option { + value + .parse::() + .ok() + .and_then(normalize_epoch_integer) + .or_else(|| value.parse::().ok().and_then(normalize_epoch_float)) +} + +fn normalize_epoch_integer(value: i64) -> Option { + let magnitude = value.unsigned_abs(); + if magnitude < 100_000_000_000 { + value.checked_mul(1_000) + } else if magnitude < 100_000_000_000_000 { + Some(value) + } else if magnitude < 100_000_000_000_000_000 { + Some(value / 1_000) + } else { + Some(value / 1_000_000) + } +} + +fn normalize_epoch_float(value: f64) -> Option { + if !value.is_finite() { + return None; + } + let magnitude = value.abs(); + let milliseconds = if magnitude < 100_000_000_000.0 { + value * 1_000.0 + } else if magnitude < 100_000_000_000_000.0 { + value + } else if magnitude < 100_000_000_000_000_000.0 { + value / 1_000.0 + } else { + value / 1_000_000.0 + }; + if milliseconds < i64::MIN as f64 || milliseconds > i64::MAX as f64 { + None + } else { + Some(milliseconds.trunc() as i64) + } +} + +fn parse_rfc3339_timestamp(value: &str) -> Option { + chrono::DateTime::parse_from_rfc3339(value) .ok() .map(|dt| dt.timestamp_millis()) } @@ -1848,7 +1997,8 @@ fn num_at(v: &Value, path: &str) -> Option { #[cfg(test)] mod tests { - use super::basename; + use super::{basename, normalize_timestamp, RolloutRecord}; + use serde_json::json; #[test] fn basename_accepts_unix_and_windows_paths() { @@ -1856,4 +2006,49 @@ mod tests { assert_eq!(basename(r"C:\Users\agent\project"), "project"); assert_eq!(basename(r"C:\Users\agent\project\\"), "project"); } + + #[test] + fn rollout_timestamps_normalize_supported_formats() { + let milliseconds = 1_704_067_202_000_i64; + assert_eq!( + normalize_timestamp(&json!(1_704_067_202_i64)), + Some(milliseconds) + ); + assert_eq!( + normalize_timestamp(&json!(milliseconds)), + Some(milliseconds) + ); + assert_eq!( + normalize_timestamp(&json!(1_704_067_202_000_000_i64)), + Some(milliseconds) + ); + assert_eq!( + normalize_timestamp(&json!(1_704_067_202_000_000_000_i64)), + Some(milliseconds) + ); + assert_eq!( + normalize_timestamp(&json!("2024-01-01T00:00:02Z")), + Some(milliseconds) + ); + assert_eq!( + normalize_timestamp(&json!("1704067202000")), + Some(milliseconds) + ); + assert_eq!( + normalize_timestamp(&json!(1_704_067_202.25_f64)), + Some(milliseconds + 250) + ); + assert_eq!(normalize_timestamp(&json!({ "seconds": 1 })), None); + } + + #[test] + fn rollout_record_treats_unknown_timestamp_formats_as_missing() { + let record: RolloutRecord = serde_json::from_value(json!({ + "timestamp": { "future": "format" }, + "type": "event_msg", + "payload": { "type": "task_started" }, + })) + .unwrap(); + assert_eq!(record.timestamp_ms, None); + } } diff --git a/bt-daemon/tests/codex_translator.rs b/bt-daemon/tests/codex_translator.rs index b0045b3..24184b4 100644 --- a/bt-daemon/tests/codex_translator.rs +++ b/bt-daemon/tests/codex_translator.rs @@ -747,13 +747,37 @@ fn tool_and_llm_payloads_preserve_original_contract() { assert_eq!(tool.tags, Some(vec!["permission-request".into()])); assert_eq!(tool.error.as_deref(), Some("boom")); - let mut llms: Vec<&SpanRow> = rows + let llms: Vec<&SpanRow> = rows .values() .filter(|row| row.span_type == SpanType::Llm) .collect(); - llms.sort_by_key(|row| row.start_ms); assert_eq!(llms.len(), 2); - assert!(llms[0] + let first = llms + .iter() + .copied() + .find(|row| { + row.output + .as_ref() + .and_then(|output| output.get("tool_calls")) + .is_some() + }) + .expect("LLM that emitted the tool call"); + let second = llms + .iter() + .copied() + .find(|row| { + row.input + .as_ref() + .and_then(Value::as_array) + .is_some_and(|input| { + input + .iter() + .any(|message| message.get("tool_calls").is_some()) + }) + }) + .expect("LLM that received the tool call result"); + assert_ne!(first.span_id, second.span_id); + assert!(first .input .as_ref() .unwrap() @@ -762,19 +786,19 @@ fn tool_and_llm_payloads_preserve_original_contract() { .iter() .all(|message| message.get("tool_calls").is_none())); assert_eq!( - llms[0].output.as_ref().unwrap()["tool_calls"][0]["function"]["arguments"], + first.output.as_ref().unwrap()["tool_calls"][0]["function"]["arguments"], json!("{\"cmd\":\"cat /tmp/review/SKILL.md\",\"sandbox_permissions\":\"require_escalated\",\"justification\":\"Need access\",\"prefix_rule\":[\"cat\"]}") ); - assert_eq!(llms[0].metrics.as_ref().unwrap()["cost"], json!(0.25)); + assert_eq!(first.metrics.as_ref().unwrap()["cost"], json!(0.25)); assert_eq!( - llms[0].metrics.as_ref().unwrap()["estimated_cost"], + first.metrics.as_ref().unwrap()["estimated_cost"], json!(0.25) ); - let second_input = llms[1].input.as_ref().unwrap().as_array().unwrap(); + let second_input = second.input.as_ref().unwrap().as_array().unwrap(); assert_eq!(second_input.last().unwrap()["role"], json!("tool")); assert_eq!(second_input.last().unwrap()["tool_call_id"], json!("c1")); assert_eq!( - llms[1].metadata.as_ref().unwrap()["usage_unavailable_reason"], + second.metadata.as_ref().unwrap()["usage_unavailable_reason"], json!("codex_token_count_missing_usage") ); } @@ -1025,7 +1049,7 @@ fn codex_untagged_replacement_history_preserves_native_order() { } #[test] -fn codex_subagent_nests_under_spawning_turn() { +fn codex_subagent_with_malformed_optional_type_nests_under_spawning_turn() { let tmp = tempfile::tempdir().unwrap(); let main_t = tmp.path().join("main.jsonl"); let sub_t = tmp.path().join("sub.jsonl"); @@ -1077,7 +1101,7 @@ fn codex_subagent_nests_under_spawning_turn() { json!({ "agent_id": "a1", "transcript_path": original_sub_p, - "agent_type": "reviewer", + "agent_type": { "name": "reviewer" }, "_bt_transcript_mirror": mirrored_subagent(), }), ), @@ -1152,7 +1176,7 @@ fn codex_subagent_nests_under_spawning_turn() { assert_eq!(subagent.parent_span_ids, vec![main_turn.span_id.clone()]); assert_eq!( subagent.metadata.as_ref().unwrap()["agent_type"], - json!("reviewer") + Value::Null ); assert!( subagent.end_ms.is_some(), @@ -1165,3 +1189,70 @@ fn codex_subagent_nests_under_spawning_turn() { let sub_llm = find(&rows, SpanType::Llm, "gpt-5.5-mini"); assert_eq!(sub_llm.parent_span_ids, vec![sub_turn.span_id.clone()]); } + +#[test] +fn codex_rollout_routing_ignores_future_fields() { + let tmp = tempfile::tempdir().unwrap(); + let transcript = tmp.path().join("rollout.jsonl"); + let records = [ + json!({ + "timestamp": "2026-01-01T00:00:01Z", + "type": "session_meta", + "future_record_field": { "preserve": true }, + "payload": { "id": "session-1", "cwd": "/work/app", "future_payload_field": [1, 2] }, + }), + json!({ + "timestamp": 1_704_067_202_000_i64, + "type": "event_msg", + "payload": { "type": "task_started", "turn_id": "t1", "future_event_field": {} }, + }), + json!({ + "timestamp": "2026-01-01T00:00:03Z", + "type": "response_item", + "payload": { + "type": "function_call", + "call_id": "c1", + "name": "shell", + "arguments": "{\"command\":\"pwd\"}", + "future_response_field": "new Codex field", + }, + }), + json!({ + "timestamp": "2026-01-01T00:00:04Z", + "type": "response_item", + "payload": { "type": "function_call_output", "call_id": "c1", "output": "/work/app" }, + }), + json!({ + "timestamp": "2026-01-01T00:00:05Z", + "type": "event_msg", + "payload": { "type": "task_complete", "last_agent_message": "done" }, + }), + ]; + let mut file = std::fs::File::create(&transcript).unwrap(); + for record in records { + writeln!(file, "{}", line(record)).unwrap(); + } + + let reg = Registry::default_agents(); + let mut translator = reg.create("codex", "sess-1"); + let ctx = SessionCtx { + session_id: "sess-1".into(), + config: None, + }; + let ops = translator + .handle( + &envelope( + "sess-1", + "SessionStart", + transcript.to_str().unwrap(), + json!({ "source": "startup" }), + ), + &ctx, + ) + .unwrap(); + + let rows = reduce(ops); + let tool = find(&rows, SpanType::Tool, "shell"); + assert_eq!(tool.input, Some(json!("{\"command\":\"pwd\"}"))); + assert_eq!(tool.output, Some(json!("/work/app"))); +}