1use std::io::Write;
11
12use agent_base::{RuntimeEvent, UserEvent};
13use anyhow::Result;
14
15use crate::session::SessionContext;
16
17pub fn save_turn_log(session_ctx: &SessionContext, turn: u32, events: &[RuntimeEvent], user_input: &str) -> Result<()> {
22 let turn_path = session_ctx.turn_path(turn as usize);
23 let mut file = std::fs::OpenOptions::new().create(true).append(true).open(&turn_path)?;
24
25 let meta = serde_json::json!({
26 "turn": turn,
27 "timestamp": chrono::Utc::now().to_rfc3339(),
28 "user_input": user_input,
29 });
30 writeln!(file, "{}", serde_json::to_string(&meta)?)?;
31
32 for event in events {
33 let line = event_to_jsonl(event);
34 writeln!(file, "{}", line)?;
35 }
36
37 writeln!(file, "{}", serde_json::to_string(&serde_json::json!({"type": "turn_end", "turn": turn}))?)?;
38
39 file.flush()?;
40
41 Ok(())
42}
43
44pub fn event_to_value(event: &RuntimeEvent) -> serde_json::Value {
46 match event {
47 RuntimeEvent::ThoughtDelta { text, .. } => {
48 serde_json::json!({"type": "thought_delta", "text": text})
49 },
50 RuntimeEvent::TextDelta { text, .. } => {
51 serde_json::json!({"type": "text_delta", "text": text})
52 },
53 RuntimeEvent::ToolCallStarted { tool_name, args_json, .. } => {
54 let args: serde_json::Value = serde_json::from_str(args_json).unwrap_or(serde_json::Value::Null);
55 serde_json::json!({"type": "tool_call_started", "tool": tool_name, "args": args})
56 },
57 RuntimeEvent::ToolCallFinished { tool_name, summary, .. } => {
58 serde_json::json!({"type": "tool_call_finished", "tool": tool_name, "summary": summary})
59 },
60 RuntimeEvent::AwaitingApproval { request, .. } => {
61 serde_json::json!({"type": "approval_request", "title": request.title, "message": request.message})
62 },
63 RuntimeEvent::PlanUpdated { explanation, plan, .. } => {
64 serde_json::json!({"type": "plan_updated", "explanation": explanation, "plan": plan})
65 },
66 RuntimeEvent::UserEvent { event: UserEvent::Structured { event_type, data }, .. } => {
67 serde_json::json!({"type": "user_event", "event_type": event_type, "data": data})
68 },
69 RuntimeEvent::UserEvent { .. } => serde_json::json!({"type": "other"}),
70 RuntimeEvent::RunCancelled { .. } => serde_json::json!({"type": "run_cancelled"}),
71 RuntimeEvent::RunFinished { .. } => serde_json::json!({"type": "run_finished"}),
72 RuntimeEvent::Checkpoint { .. } => serde_json::json!({"type": "checkpoint"}),
73 }
74}
75
76pub fn event_to_jsonl(event: &RuntimeEvent) -> String {
78 let value = event_to_value(event);
79 serde_json::to_string(&value).unwrap_or_else(|_| r#"{"type":"serialize_error"}"#.to_string())
80}
81
82#[cfg(test)]
83mod tests {
84 use super::*;
85 use agent_base::{SessionId, UserEvent};
86 use tempfile::TempDir;
87
88 fn session_id() -> SessionId {
89 SessionId { id: 1, external_id: None }
90 }
91
92 #[test]
95 fn test_event_to_value_text_delta() {
96 let v = event_to_value(&RuntimeEvent::TextDelta {
97 session_id: session_id(),
98 text: "hello".into(),
99 });
100 assert_eq!(v["type"], "text_delta");
101 assert_eq!(v["text"], "hello");
102 }
103
104 #[test]
105 fn test_event_to_value_thought_delta() {
106 let v = event_to_value(&RuntimeEvent::ThoughtDelta {
107 session_id: session_id(),
108 text: "thinking".into(),
109 });
110 assert_eq!(v["type"], "thought_delta");
111 }
112
113 #[test]
114 fn test_event_to_value_tool_call_started() {
115 let v = event_to_value(&RuntimeEvent::ToolCallStarted {
116 session_id: session_id(),
117 tool_name: "shell".into(),
118 args_json: r#"{"cmd":"ls"}"#.into(),
119 });
120 assert_eq!(v["type"], "tool_call_started");
121 assert_eq!(v["tool"], "shell");
122 assert_eq!(v["args"]["cmd"], "ls");
123 }
124
125 #[test]
126 fn test_event_to_value_tool_call_finished() {
127 let v = event_to_value(&RuntimeEvent::ToolCallFinished {
128 session_id: session_id(),
129 tool_name: "shell".into(),
130 summary: "done".into(),
131 });
132 assert_eq!(v["type"], "tool_call_finished");
133 assert_eq!(v["summary"], "done");
134 }
135
136 #[test]
137 fn test_event_to_value_run_finished() {
138 let v = event_to_value(&RuntimeEvent::RunFinished {
139 session_id: session_id(),
140 });
141 assert_eq!(v["type"], "run_finished");
142 }
143
144 #[test]
145 fn test_event_to_value_run_cancelled() {
146 let v = event_to_value(&RuntimeEvent::RunCancelled {
147 session_id: session_id(),
148 });
149 assert_eq!(v["type"], "run_cancelled");
150 }
151
152 #[test]
153 fn test_event_to_value_user_event_other() {
154 let v = event_to_value(&RuntimeEvent::UserEvent {
155 session_id: session_id(),
156 event: UserEvent::Progress { text: "loading".into() },
157 });
158 assert_eq!(v["type"], "other");
159 }
160
161 #[test]
162 fn test_event_to_value_user_event_structured() {
163 let v = event_to_value(&RuntimeEvent::UserEvent {
164 session_id: session_id(),
165 event: UserEvent::Structured {
166 event_type: "custom".into(),
167 data: serde_json::json!({"key": "value"}),
168 },
169 });
170 assert_eq!(v["type"], "user_event");
171 assert_eq!(v["event_type"], "custom");
172 }
173
174 #[test]
177 fn test_event_to_jsonl_returns_valid_json() {
178 let line = event_to_jsonl(&RuntimeEvent::TextDelta {
179 session_id: session_id(),
180 text: "hello".into(),
181 });
182 let v: serde_json::Value = serde_json::from_str(&line).unwrap();
183 assert_eq!(v["type"], "text_delta");
184 }
185
186 #[test]
187 fn test_event_to_jsonl_all_variants_valid() {
188 let events: Vec<RuntimeEvent> = vec![
189 RuntimeEvent::TextDelta { session_id: session_id(), text: "t".into() },
190 RuntimeEvent::ThoughtDelta { session_id: session_id(), text: "t".into() },
191 RuntimeEvent::ToolCallStarted { session_id: session_id(), tool_name: "t".into(), args_json: "{}".into() },
192 RuntimeEvent::ToolCallFinished { session_id: session_id(), tool_name: "t".into(), summary: "s".into() },
193 RuntimeEvent::RunFinished { session_id: session_id() },
194 RuntimeEvent::RunCancelled { session_id: session_id() },
195 RuntimeEvent::UserEvent {
196 session_id: session_id(),
197 event: UserEvent::Progress { text: "t".into() },
198 },
199 ];
200 for e in &events {
201 let line = event_to_jsonl(e);
202 assert!(serde_json::from_str::<serde_json::Value>(&line).is_ok(),
203 "Failed to parse JSONL for event");
204 }
205 }
206
207 fn make_session_ctx(tmp: &TempDir) -> SessionContext {
210 crate::session::resolve_session(Some("test-session"), tmp.path()).unwrap()
211 }
212
213 #[test]
214 fn test_save_turn_log_creates_file() {
215 let tmp = TempDir::new().unwrap();
216 let ctx = make_session_ctx(&tmp);
217 let events = vec![
218 RuntimeEvent::TextDelta { session_id: session_id(), text: "hello".into() },
219 ];
220 save_turn_log(&ctx, 1, &events, "test input").unwrap();
221 let path = ctx.turn_path(1);
222 assert!(path.exists());
223 let content = std::fs::read_to_string(&path).unwrap();
224 assert!(!content.is_empty());
225 }
226
227 #[test]
228 fn test_save_turn_log_contains_metadata() {
229 let tmp = TempDir::new().unwrap();
230 let ctx = make_session_ctx(&tmp);
231 let events = vec![
232 RuntimeEvent::TextDelta { session_id: session_id(), text: "hello".into() },
233 ];
234 save_turn_log(&ctx, 1, &events, "my input").unwrap();
235 let content = std::fs::read_to_string(ctx.turn_path(1)).unwrap();
236 let lines: Vec<&str> = content.lines().collect();
237 let first: serde_json::Value = serde_json::from_str(lines[0]).unwrap();
238 assert_eq!(first["turn"], 1);
239 assert_eq!(first["user_input"], "my input");
240 assert!(first["timestamp"].as_str().is_some());
241 }
242
243 #[test]
244 fn test_save_turn_log_contains_events() {
245 let tmp = TempDir::new().unwrap();
246 let ctx = make_session_ctx(&tmp);
247 let events = vec![
248 RuntimeEvent::TextDelta { session_id: session_id(), text: "hello".into() },
249 RuntimeEvent::RunFinished { session_id: session_id() },
250 ];
251 save_turn_log(&ctx, 1, &events, "input").unwrap();
252 let content = std::fs::read_to_string(ctx.turn_path(1)).unwrap();
253 let lines: Vec<&str> = content.lines().collect();
254 assert!(lines.len() >= 4);
256 let line1: serde_json::Value = serde_json::from_str(lines[1]).unwrap();
257 assert_eq!(line1["type"], "text_delta");
258 let line2: serde_json::Value = serde_json::from_str(lines[2]).unwrap();
259 assert_eq!(line2["type"], "run_finished");
260 }
261
262 #[test]
263 fn test_save_turn_log_contains_turn_end_marker() {
264 let tmp = TempDir::new().unwrap();
265 let ctx = make_session_ctx(&tmp);
266 let events = vec![
267 RuntimeEvent::TextDelta { session_id: session_id(), text: "hello".into() },
268 ];
269 save_turn_log(&ctx, 1, &events, "input").unwrap();
270 let content = std::fs::read_to_string(ctx.turn_path(1)).unwrap();
271 let last_line = content.lines().last().unwrap();
272 let v: serde_json::Value = serde_json::from_str(last_line).unwrap();
273 assert_eq!(v["type"], "turn_end");
274 assert_eq!(v["turn"], 1);
275 }
276
277 #[test]
278 fn test_save_turn_log_appends() {
279 let tmp = TempDir::new().unwrap();
280 let ctx = make_session_ctx(&tmp);
281 let events1 = vec![
282 RuntimeEvent::TextDelta { session_id: session_id(), text: "turn1".into() },
283 ];
284 let events2 = vec![
285 RuntimeEvent::TextDelta { session_id: session_id(), text: "turn2".into() },
286 ];
287 save_turn_log(&ctx, 1, &events1, "input1").unwrap();
288 save_turn_log(&ctx, 1, &events2, "input2").unwrap();
289 let content = std::fs::read_to_string(ctx.turn_path(1)).unwrap();
290 assert!(content.contains("turn1"));
291 assert!(content.contains("turn2"));
292 let turn_end_count = content.lines().filter(|l| l.contains("turn_end")).count();
294 assert_eq!(turn_end_count, 2);
295 }
296
297 #[test]
298 fn test_save_turn_log_empty_events() {
299 let tmp = TempDir::new().unwrap();
300 let ctx = make_session_ctx(&tmp);
301 save_turn_log(&ctx, 1, &[], "no events").unwrap();
302 let content = std::fs::read_to_string(ctx.turn_path(1)).unwrap();
303 let lines: Vec<&str> = content.lines().collect();
304 assert_eq!(lines.len(), 2);
306 let last: serde_json::Value = serde_json::from_str(lines[1]).unwrap();
307 assert_eq!(last["type"], "turn_end");
308 assert_eq!(last["turn"], 1);
309 }
310}