1use 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#[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 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 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 write_json_file(&absolute_path, value)?;
101 Ok(RawPayloadRef {
102 raw_payload_id,
103 kind,
104 path: relative_path,
105 })
106 }
107
108 pub fn append(&self, payload: RawTraceEventPayload) -> Result<RawTraceEvent> {
110 self.append_with_context(RawTraceEventContext::default(), payload)
111 }
112
113 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 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}