1use std::io::Write;
11
12use agent_base::{AgentResult, RuntimeEvent, UserEvent};
13
14use crate::session::SessionContext;
15
16pub fn save_turn_log(
21 session_ctx: &SessionContext,
22 turn: u32,
23 events: &[RuntimeEvent],
24 user_input: &str,
25) -> AgentResult<()> {
26 let turn_path = session_ctx.turn_path(turn as usize);
27 let mut file = std::fs::OpenOptions::new().create(true).append(true).open(&turn_path)?;
28
29 let meta = serde_json::json!({
30 "turn": turn,
31 "timestamp": chrono::Utc::now().to_rfc3339(),
32 "user_input": user_input,
33 });
34 writeln!(file, "{}", serde_json::to_string(&meta)?)?;
35
36 for event in events {
37 let line = event_to_jsonl(event);
38 writeln!(file, "{}", line)?;
39 }
40
41 writeln!(file, "{}", serde_json::to_string(&serde_json::json!({"type": "turn_end", "turn": turn}))?)?;
42
43 file.flush()?;
44
45 Ok(())
46}
47
48pub fn event_to_value(event: &RuntimeEvent) -> serde_json::Value {
50 let mut value = match event {
51 RuntimeEvent::ThoughtDelta { text, .. } => {
52 serde_json::json!({"type": "thought_delta", "text": text})
53 },
54 RuntimeEvent::TextDelta { text, .. } => {
55 serde_json::json!({"type": "text_delta", "text": text})
56 },
57 RuntimeEvent::ToolCallDraft { name, args_len, count, .. } => {
58 serde_json::json!({"type": "tool_call_draft", "tool": name, "args_len": args_len, "count": count})
61 },
62 RuntimeEvent::ToolCallStarted { tool_name, args_json, .. } => {
63 let args: serde_json::Value = serde_json::from_str(args_json).unwrap_or(serde_json::Value::Null);
64 serde_json::json!({"type": "tool_call_started", "tool": tool_name, "args": args})
65 },
66 RuntimeEvent::ToolCallFinished { tool_name, summary, denied, .. } => {
67 serde_json::json!({"type": "tool_call_finished", "tool": tool_name, "summary": summary, "denied": denied})
68 },
69 RuntimeEvent::AwaitingApproval { request, .. } => {
70 serde_json::json!({"type": "approval_request", "title": request.title, "message": request.message})
71 },
72 RuntimeEvent::PlanUpdated { explanation, plan, .. } => {
73 serde_json::json!({"type": "plan_updated", "explanation": explanation, "plan": plan})
74 },
75 RuntimeEvent::UserEvent { event: UserEvent::Structured { event_type, data }, .. } => {
76 serde_json::json!({"type": "user_event", "event_type": event_type, "data": data})
77 },
78 RuntimeEvent::UserEvent { .. } => serde_json::json!({"type": "other"}),
79 RuntimeEvent::RunCancelled { .. } => serde_json::json!({"type": "run_cancelled"}),
80 RuntimeEvent::RunFinished { .. } => serde_json::json!({"type": "run_finished"}),
81 RuntimeEvent::Checkpoint { .. } => serde_json::json!({"type": "checkpoint"}),
82 };
83 if let Some(agent_id) = event.agent_id() {
84 value["agent_id"] = serde_json::json!(agent_id);
85 }
86 value
87}
88
89pub fn event_to_jsonl(event: &RuntimeEvent) -> String {
91 let value = event_to_value(event);
92 serde_json::to_string(&value).unwrap_or_else(|_| r#"{"type":"serialize_error"}"#.to_string())
93}
94
95#[cfg(test)]
96mod tests {
97 use super::*;
98 use agent_base::{SessionId, UserEvent};
99 use tempfile::TempDir;
100
101 fn session_id() -> SessionId {
102 SessionId { id: 1, external_id: None }
103 }
104
105 #[test]
108 fn test_event_to_value_text_delta() {
109 let v = event_to_value(&RuntimeEvent::TextDelta {
110 session_id: session_id(),
111 text: "hello".into(),
112 agent_id: None,
113 trace_id: None,
114 });
115 assert_eq!(v["type"], "text_delta");
116 assert_eq!(v["text"], "hello");
117 }
118
119 #[test]
120 fn test_event_to_value_thought_delta() {
121 let v = event_to_value(&RuntimeEvent::ThoughtDelta {
122 session_id: session_id(),
123 text: "thinking".into(),
124 agent_id: None,
125 trace_id: None,
126 });
127 assert_eq!(v["type"], "thought_delta");
128 }
129
130 #[test]
131 fn test_event_to_value_tool_call_started() {
132 let v = event_to_value(&RuntimeEvent::ToolCallStarted {
133 session_id: session_id(),
134 tool_name: "shell".into(),
135 args_json: r#"{"cmd":"ls"}"#.into(),
136 agent_id: None,
137 trace_id: None,
138 });
139 assert_eq!(v["type"], "tool_call_started");
140 assert_eq!(v["tool"], "shell");
141 assert_eq!(v["args"]["cmd"], "ls");
142 }
143
144 #[test]
145 fn test_event_to_value_tool_call_finished() {
146 let v = event_to_value(&RuntimeEvent::ToolCallFinished {
147 session_id: session_id(),
148 tool_name: "shell".into(),
149 summary: "done".into(),
150 agent_id: None,
151 trace_id: None,
152 denied: false,
153 details: None,
154 });
155 assert_eq!(v["type"], "tool_call_finished");
156 assert_eq!(v["summary"], "done");
157 assert_eq!(v["denied"], false);
158 }
159
160 #[test]
161 fn test_event_to_value_run_finished() {
162 let v = event_to_value(&RuntimeEvent::RunFinished { session_id: session_id(), agent_id: None, trace_id: None });
163 assert_eq!(v["type"], "run_finished");
164 }
165
166 #[test]
167 fn test_event_to_value_run_cancelled() {
168 let v =
169 event_to_value(&RuntimeEvent::RunCancelled { session_id: session_id(), agent_id: None, trace_id: None });
170 assert_eq!(v["type"], "run_cancelled");
171 }
172
173 #[test]
174 fn test_event_to_value_user_event_other() {
175 let v = event_to_value(&RuntimeEvent::UserEvent {
176 session_id: session_id(),
177 event: UserEvent::Progress { text: "loading".into() },
178 agent_id: None,
179 trace_id: None,
180 });
181 assert_eq!(v["type"], "other");
182 }
183
184 #[test]
185 fn test_event_to_value_user_event_structured() {
186 let v = event_to_value(&RuntimeEvent::UserEvent {
187 session_id: session_id(),
188 event: UserEvent::Structured { event_type: "custom".into(), data: serde_json::json!({"key": "value"}) },
189 agent_id: None,
190 trace_id: None,
191 });
192 assert_eq!(v["type"], "user_event");
193 assert_eq!(v["event_type"], "custom");
194 }
195
196 #[test]
199 fn test_event_to_jsonl_returns_valid_json() {
200 let line = event_to_jsonl(&RuntimeEvent::TextDelta {
201 session_id: session_id(),
202 text: "hello".into(),
203 agent_id: None,
204 trace_id: None,
205 });
206 let v: serde_json::Value = serde_json::from_str(&line).unwrap();
207 assert_eq!(v["type"], "text_delta");
208 }
209
210 #[test]
211 fn test_event_to_jsonl_all_variants_valid() {
212 let events: Vec<RuntimeEvent> = vec![
213 RuntimeEvent::TextDelta { session_id: session_id(), text: "t".into(), agent_id: None, trace_id: None },
214 RuntimeEvent::ThoughtDelta { session_id: session_id(), text: "t".into(), agent_id: None, trace_id: None },
215 RuntimeEvent::ToolCallStarted {
216 session_id: session_id(),
217 tool_name: "t".into(),
218 args_json: "{}".into(),
219 agent_id: None,
220 trace_id: None,
221 },
222 RuntimeEvent::ToolCallFinished {
223 session_id: session_id(),
224 tool_name: "t".into(),
225 summary: "s".into(),
226 agent_id: None,
227 trace_id: None,
228 denied: false,
229 details: None,
230 },
231 RuntimeEvent::RunFinished { session_id: session_id(), agent_id: None, trace_id: None },
232 RuntimeEvent::RunCancelled { session_id: session_id(), agent_id: None, trace_id: None },
233 RuntimeEvent::UserEvent {
234 session_id: session_id(),
235 event: UserEvent::Progress { text: "t".into() },
236 agent_id: None,
237 trace_id: None,
238 },
239 ];
240 for e in &events {
241 let line = event_to_jsonl(e);
242 assert!(serde_json::from_str::<serde_json::Value>(&line).is_ok(), "Failed to parse JSONL for event");
243 }
244 }
245
246 fn make_session_ctx(tmp: &TempDir) -> SessionContext {
249 crate::session::resolve_session(Some("test-session"), tmp.path()).unwrap()
250 }
251
252 #[test]
253 fn test_save_turn_log_creates_file() {
254 let tmp = TempDir::new().unwrap();
255 let ctx = make_session_ctx(&tmp);
256 let events = vec![RuntimeEvent::TextDelta {
257 session_id: session_id(),
258 text: "hello".into(),
259 agent_id: None,
260 trace_id: None,
261 }];
262 save_turn_log(&ctx, 1, &events, "test input").unwrap();
263 let path = ctx.turn_path(1);
264 assert!(path.exists());
265 let content = std::fs::read_to_string(&path).unwrap();
266 assert!(!content.is_empty());
267 }
268
269 #[test]
270 fn test_save_turn_log_contains_metadata() {
271 let tmp = TempDir::new().unwrap();
272 let ctx = make_session_ctx(&tmp);
273 let events = vec![RuntimeEvent::TextDelta {
274 session_id: session_id(),
275 text: "hello".into(),
276 agent_id: None,
277 trace_id: None,
278 }];
279 save_turn_log(&ctx, 1, &events, "my input").unwrap();
280 let content = std::fs::read_to_string(ctx.turn_path(1)).unwrap();
281 let lines: Vec<&str> = content.lines().collect();
282 let first: serde_json::Value = serde_json::from_str(lines[0]).unwrap();
283 assert_eq!(first["turn"], 1);
284 assert_eq!(first["user_input"], "my input");
285 assert!(first["timestamp"].as_str().is_some());
286 }
287
288 #[test]
289 fn test_save_turn_log_contains_events() {
290 let tmp = TempDir::new().unwrap();
291 let ctx = make_session_ctx(&tmp);
292 let events = vec![
293 RuntimeEvent::TextDelta { session_id: session_id(), text: "hello".into(), agent_id: None, trace_id: None },
294 RuntimeEvent::RunFinished { session_id: session_id(), agent_id: None, trace_id: None },
295 ];
296 save_turn_log(&ctx, 1, &events, "input").unwrap();
297 let content = std::fs::read_to_string(ctx.turn_path(1)).unwrap();
298 let lines: Vec<&str> = content.lines().collect();
299 assert!(lines.len() >= 4);
301 let line1: serde_json::Value = serde_json::from_str(lines[1]).unwrap();
302 assert_eq!(line1["type"], "text_delta");
303 let line2: serde_json::Value = serde_json::from_str(lines[2]).unwrap();
304 assert_eq!(line2["type"], "run_finished");
305 }
306
307 #[test]
308 fn test_save_turn_log_contains_turn_end_marker() {
309 let tmp = TempDir::new().unwrap();
310 let ctx = make_session_ctx(&tmp);
311 let events = vec![RuntimeEvent::TextDelta {
312 session_id: session_id(),
313 text: "hello".into(),
314 agent_id: None,
315 trace_id: None,
316 }];
317 save_turn_log(&ctx, 1, &events, "input").unwrap();
318 let content = std::fs::read_to_string(ctx.turn_path(1)).unwrap();
319 let last_line = content.lines().last().unwrap();
320 let v: serde_json::Value = serde_json::from_str(last_line).unwrap();
321 assert_eq!(v["type"], "turn_end");
322 assert_eq!(v["turn"], 1);
323 }
324
325 #[test]
326 fn test_save_turn_log_appends() {
327 let tmp = TempDir::new().unwrap();
328 let ctx = make_session_ctx(&tmp);
329 let events1 = vec![RuntimeEvent::TextDelta {
330 session_id: session_id(),
331 text: "turn1".into(),
332 agent_id: None,
333 trace_id: None,
334 }];
335 let events2 = vec![RuntimeEvent::TextDelta {
336 session_id: session_id(),
337 text: "turn2".into(),
338 agent_id: None,
339 trace_id: None,
340 }];
341 save_turn_log(&ctx, 1, &events1, "input1").unwrap();
342 save_turn_log(&ctx, 1, &events2, "input2").unwrap();
343 let content = std::fs::read_to_string(ctx.turn_path(1)).unwrap();
344 assert!(content.contains("turn1"));
345 assert!(content.contains("turn2"));
346 let turn_end_count = content.lines().filter(|l| l.contains("turn_end")).count();
348 assert_eq!(turn_end_count, 2);
349 }
350
351 #[test]
352 fn test_save_turn_log_empty_events() {
353 let tmp = TempDir::new().unwrap();
354 let ctx = make_session_ctx(&tmp);
355 save_turn_log(&ctx, 1, &[], "no events").unwrap();
356 let content = std::fs::read_to_string(ctx.turn_path(1)).unwrap();
357 let lines: Vec<&str> = content.lines().collect();
358 assert_eq!(lines.len(), 2);
360 let last: serde_json::Value = serde_json::from_str(lines[1]).unwrap();
361 assert_eq!(last["type"], "turn_end");
362 assert_eq!(last["turn"], 1);
363 }
364
365 #[test]
368 fn test_event_to_value_attaches_agent_id() {
369 let v = event_to_value(&RuntimeEvent::TextDelta {
370 session_id: session_id(),
371 text: "result".into(),
372 agent_id: Some("root/searcher".into()),
373 trace_id: None,
374 });
375 assert_eq!(v["type"], "text_delta");
376 assert_eq!(v["text"], "result");
377 assert_eq!(v["agent_id"], "root/searcher");
378 }
379
380 #[test]
381 fn test_event_to_value_omits_agent_id_when_none() {
382 let v = event_to_value(&RuntimeEvent::TextDelta {
383 session_id: session_id(),
384 text: "result".into(),
385 agent_id: None,
386 trace_id: None,
387 });
388 assert_eq!(v["type"], "text_delta");
389 assert!(v.get("agent_id").is_none());
390 }
391}