Skip to main content

codex_rollout_trace/
writer.rs

1//! Hot-path trace bundle writer.
2
3use std::fs::File;
4use std::fs::OpenOptions;
5use std::io::BufWriter;
6use std::io::Write;
7use std::path::Path;
8use std::path::PathBuf;
9use std::sync::Mutex;
10use std::sync::MutexGuard;
11use std::sync::PoisonError;
12use std::time::SystemTime;
13use std::time::UNIX_EPOCH;
14
15use anyhow::Context;
16use anyhow::Result;
17use serde::Serialize;
18
19use crate::bundle::MANIFEST_FILE_NAME;
20use crate::bundle::PAYLOADS_DIR_NAME;
21use crate::bundle::RAW_EVENT_LOG_FILE_NAME;
22use crate::bundle::TraceBundleManifest;
23use crate::model::AgentThreadId;
24use crate::payload::RawPayloadKind;
25use crate::payload::RawPayloadRef;
26use crate::raw_event::RAW_TRACE_EVENT_SCHEMA_VERSION;
27use crate::raw_event::RawTraceEvent;
28use crate::raw_event::RawTraceEventContext;
29use crate::raw_event::RawTraceEventPayload;
30
31/// Local trace bundle writer.
32///
33/// The writer appends raw events and writes payload files. It does not keep a
34/// reduced `RolloutTrace` in memory; replay is owned by the reducer.
35#[derive(Debug)]
36pub struct TraceWriter {
37    inner: Mutex<TraceWriterInner>,
38}
39
40#[derive(Debug)]
41struct TraceWriterInner {
42    manifest: TraceBundleManifest,
43    payloads_dir: PathBuf,
44    event_log: BufWriter<File>,
45    next_seq: u64,
46    next_payload_ordinal: u64,
47}
48
49impl TraceWriter {
50    /// Creates a trace bundle directory and writes its manifest.
51    pub fn create(
52        bundle_dir: impl AsRef<Path>,
53        trace_id: String,
54        rollout_id: String,
55        root_thread_id: AgentThreadId,
56    ) -> Result<Self> {
57        let bundle_dir = bundle_dir.as_ref().to_path_buf();
58        let payloads_dir = bundle_dir.join(PAYLOADS_DIR_NAME);
59        std::fs::create_dir_all(&payloads_dir)
60            .with_context(|| format!("create trace payload dir {}", payloads_dir.display()))?;
61
62        let started_at_unix_ms = unix_time_ms();
63        let manifest =
64            TraceBundleManifest::new(trace_id, rollout_id, root_thread_id, started_at_unix_ms);
65        write_json_file(&bundle_dir.join(MANIFEST_FILE_NAME), &manifest)?;
66
67        let event_log_path = bundle_dir.join(RAW_EVENT_LOG_FILE_NAME);
68        let event_log = OpenOptions::new()
69            .create(true)
70            .append(true)
71            .open(&event_log_path)
72            .with_context(|| format!("open trace event log {}", event_log_path.display()))?;
73
74        Ok(Self {
75            inner: Mutex::new(TraceWriterInner {
76                manifest,
77                payloads_dir,
78                event_log: BufWriter::new(event_log),
79                next_seq: 1,
80                next_payload_ordinal: 1,
81            }),
82        })
83    }
84
85    /// Writes a JSON payload file and returns its reduced-state reference.
86    pub fn write_json_payload(
87        &self,
88        kind: RawPayloadKind,
89        value: &impl Serialize,
90    ) -> Result<RawPayloadRef> {
91        let mut inner = self.lock_inner();
92        let ordinal = inner.next_payload_ordinal;
93        inner.next_payload_ordinal += 1;
94        let raw_payload_id = format!("raw_payload:{ordinal}");
95        let relative_path = format!("{PAYLOADS_DIR_NAME}/{ordinal}.json");
96        let absolute_path = inner.payloads_dir.join(format!("{ordinal}.json"));
97        // Payload files are created before the event that references them. A
98        // replay interrupted after an event is appended should never point at a
99        // payload file that the writer planned but had not written yet.
100        write_json_file(&absolute_path, value)?;
101        Ok(RawPayloadRef {
102            raw_payload_id,
103            kind,
104            path: relative_path,
105        })
106    }
107
108    /// Appends one raw event with no extra envelope context.
109    pub fn append(&self, payload: RawTraceEventPayload) -> Result<RawTraceEvent> {
110        self.append_with_context(RawTraceEventContext::default(), payload)
111    }
112
113    /// Appends one raw event with explicit thread/turn context.
114    pub fn append_with_context(
115        &self,
116        context: RawTraceEventContext,
117        payload: RawTraceEventPayload,
118    ) -> Result<RawTraceEvent> {
119        let mut inner = self.lock_inner();
120        let event = RawTraceEvent {
121            schema_version: RAW_TRACE_EVENT_SCHEMA_VERSION,
122            seq: inner.next_seq,
123            wall_time_unix_ms: unix_time_ms(),
124            rollout_id: inner.manifest.rollout_id.clone(),
125            thread_id: context.thread_id,
126            codex_turn_id: context.codex_turn_id,
127            payload,
128        };
129        inner.next_seq += 1;
130        serde_json::to_writer(&mut inner.event_log, &event)?;
131        inner.event_log.write_all(b"\n")?;
132        inner.event_log.flush()?;
133        Ok(event)
134    }
135
136    fn lock_inner(&self) -> MutexGuard<'_, TraceWriterInner> {
137        // Preserve the event log after a panic in tracing code. Dropping the
138        // writer would lose subsequent diagnostic events in exactly the session
139        // we are trying to debug.
140        self.inner.lock().unwrap_or_else(PoisonError::into_inner)
141    }
142}
143
144fn write_json_file(path: &Path, value: &impl Serialize) -> Result<()> {
145    let file = File::create(path).with_context(|| format!("create {}", path.display()))?;
146    serde_json::to_writer_pretty(file, value)
147        .with_context(|| format!("write JSON {}", path.display()))
148}
149
150pub(crate) fn unix_time_ms() -> i64 {
151    let duration = SystemTime::now()
152        .duration_since(UNIX_EPOCH)
153        .unwrap_or_default();
154    i64::try_from(duration.as_millis()).unwrap_or(i64::MAX)
155}
156
157#[cfg(test)]
158mod tests {
159    use pretty_assertions::assert_eq;
160    use serde_json::json;
161    use tempfile::TempDir;
162
163    use crate::model::ExecutionStatus;
164    use crate::model::RolloutStatus;
165    use crate::payload::RawPayloadKind;
166    use crate::raw_event::RawTraceEventPayload;
167    use crate::replay_bundle;
168    use crate::writer::TraceWriter;
169
170    #[test]
171    fn writer_records_payload_refs_and_replays_rollout_status() -> anyhow::Result<()> {
172        let temp = TempDir::new()?;
173        let writer = TraceWriter::create(
174            temp.path(),
175            "trace-1".to_string(),
176            "rollout-1".to_string(),
177            "thread-root".to_string(),
178        )?;
179
180        writer.append(RawTraceEventPayload::RolloutStarted {
181            trace_id: "trace-1".to_string(),
182            root_thread_id: "thread-root".to_string(),
183        })?;
184        let metadata_payload = writer.write_json_payload(
185            RawPayloadKind::ProtocolEvent,
186            &json!({
187                "source": "test",
188                "model": "gpt-test",
189            }),
190        )?;
191        writer.append(RawTraceEventPayload::ThreadStarted {
192            thread_id: "thread-root".to_string(),
193            agent_path: "/root".to_string(),
194            metadata_payload: Some(metadata_payload.clone()),
195        })?;
196        writer.append(RawTraceEventPayload::CodexTurnStarted {
197            codex_turn_id: "turn-1".to_string(),
198            thread_id: "thread-root".to_string(),
199        })?;
200        let inference_request = writer.write_json_payload(
201            RawPayloadKind::InferenceRequest,
202            &json!({
203                "model": "gpt-test",
204                "input": [{
205                    "type": "message",
206                    "role": "user",
207                    "content": [{"type": "input_text", "text": "hello"}]
208                }],
209            }),
210        )?;
211        writer.append(RawTraceEventPayload::InferenceStarted {
212            inference_call_id: "inference-1".to_string(),
213            thread_id: "thread-root".to_string(),
214            codex_turn_id: "turn-1".to_string(),
215            model: "gpt-test".to_string(),
216            provider_name: "test-provider".to_string(),
217            request_payload: inference_request.clone(),
218        })?;
219        let inference_response = writer.write_json_payload(
220            RawPayloadKind::InferenceResponse,
221            &json!({
222                "response_id": "resp-1",
223                "output_items": [],
224            }),
225        )?;
226        writer.append(RawTraceEventPayload::InferenceCompleted {
227            inference_call_id: "inference-1".to_string(),
228            response_id: Some("resp-1".to_string()),
229            upstream_request_id: Some("req-1".to_string()),
230            response_payload: inference_response.clone(),
231        })?;
232        writer.append(RawTraceEventPayload::CodexTurnEnded {
233            codex_turn_id: "turn-1".to_string(),
234            status: ExecutionStatus::Completed,
235        })?;
236        writer.append(RawTraceEventPayload::RolloutEnded {
237            status: RolloutStatus::Completed,
238        })?;
239
240        let rollout = replay_bundle(temp.path())?;
241
242        assert_eq!(rollout.status, RolloutStatus::Completed);
243        assert_eq!(rollout.root_thread_id, "thread-root");
244        assert_eq!(rollout.threads["thread-root"].agent_path, "/root");
245        assert_eq!(rollout.codex_turns["turn-1"].thread_id, "thread-root");
246        assert_eq!(
247            rollout.codex_turns["turn-1"].execution.status,
248            ExecutionStatus::Completed,
249        );
250        assert_eq!(
251            rollout.inference_calls["inference-1"].raw_request_payload_id,
252            inference_request.raw_payload_id,
253        );
254        assert_eq!(
255            rollout.inference_calls["inference-1"].raw_response_payload_id,
256            Some(inference_response.raw_payload_id),
257        );
258        assert_eq!(
259            rollout.raw_payloads[&metadata_payload.raw_payload_id].path,
260            "payloads/1.json"
261        );
262
263        Ok(())
264    }
265}