Skip to main content

phi_agent/
event_log.rs

1//! Event log: serialize turn events to JSONL files.
2//!
3//! Provides event persistence integrated with SessionContext,
4//! shared across CLI, Web, and other consumers.
5//!
6//! Note: `save_turn_log` performs synchronous file I/O. Consumers should
7//! call it via `tokio::task::spawn_blocking` in async contexts to avoid
8//! blocking the runtime.
9
10use std::io::Write;
11
12use agent_base::{AgentResult, RuntimeEvent, UserEvent};
13
14use crate::session::SessionContext;
15
16/// Save all events from a turn to a JSONL file.
17///
18/// Performs synchronous file I/O. Callers should invoke this via
19/// `tokio::task::spawn_blocking`.
20pub 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
48/// Convert a RuntimeEvent to a JSON Value (shared by event_log and render).
49pub 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            // Progress on a call still being streamed in — recorded (not
59            // rendered) so the turn log can explain a silent window.
60            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
89/// Convert a RuntimeEvent to a JSONL line.
90pub 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    // ── event_to_value tests ──
106
107    #[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    // ── event_to_jsonl tests ──
197
198    #[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    // ── save_turn_log tests ──
247
248    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        // Line 0: metadata, Line 1: text_delta, Line 2: run_finished, Line 3: turn_end
300        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        // Should have two turn_end markers
347        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        // metadata header + turn_end marker
359        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    // ── agent_id serialization tests ──
366
367    #[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}