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::{RuntimeEvent, UserEvent};
13use anyhow::Result;
14
15use crate::session::SessionContext;
16
17/// Save all events from a turn to a JSONL file.
18///
19/// Performs synchronous file I/O. Callers should invoke this via
20/// `tokio::task::spawn_blocking`.
21pub 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
44/// Convert a RuntimeEvent to a JSON Value (shared by event_log and render).
45pub 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
76/// Convert a RuntimeEvent to a JSONL line.
77pub 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    // ── event_to_value tests ──
93
94    #[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    // ── event_to_jsonl tests ──
175
176    #[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    // ── save_turn_log tests ──
208
209    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        // Line 0: metadata, Line 1: text_delta, Line 2: run_finished, Line 3: turn_end
255        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        // Should have two turn_end markers
293        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        // metadata header + turn_end marker
305        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}