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 anyhow::Result;
13use agent_base::{RuntimeEvent, UserEvent};
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(
22    session_ctx: &SessionContext,
23    turn: u32,
24    events: &[RuntimeEvent],
25    user_input: &str,
26) -> Result<()> {
27    let turn_path = session_ctx.turn_path(turn as usize);
28    let mut file = std::fs::OpenOptions::new()
29        .create(true)
30        .append(true)
31        .open(&turn_path)?;
32
33    let meta = serde_json::json!({
34        "turn": turn,
35        "timestamp": chrono::Utc::now().to_rfc3339(),
36        "user_input": user_input,
37    });
38    writeln!(file, "{}", serde_json::to_string(&meta)?)?;
39
40    for event in events {
41        let line = event_to_jsonl(event);
42        writeln!(file, "{}", line)?;
43    }
44
45    writeln!(
46        file,
47        "{}",
48        serde_json::to_string(&serde_json::json!({"type": "turn_end", "turn": turn}))?
49    )?;
50
51    file.flush()?;
52
53    Ok(())
54}
55
56/// Convert a RuntimeEvent to a JSON Value (shared by event_log and render).
57pub fn event_to_value(event: &RuntimeEvent) -> serde_json::Value {
58    match event {
59        RuntimeEvent::ThoughtDelta { text, .. } => {
60            serde_json::json!({"type": "thought_delta", "text": text})
61        }
62        RuntimeEvent::TextDelta { text, .. } => {
63            serde_json::json!({"type": "text_delta", "text": text})
64        }
65        RuntimeEvent::ToolCallStarted { tool_name, args_json, .. } => {
66            let args: serde_json::Value =
67                serde_json::from_str(args_json).unwrap_or(serde_json::Value::Null);
68            serde_json::json!({"type": "tool_call_started", "tool": tool_name, "args": args})
69        }
70        RuntimeEvent::ToolCallFinished { tool_name, summary, .. } => {
71            serde_json::json!({"type": "tool_call_finished", "tool": tool_name, "summary": summary})
72        }
73        RuntimeEvent::AwaitingApproval { request, .. } => {
74            serde_json::json!({"type": "approval_request", "title": request.title, "message": request.message})
75        }
76        RuntimeEvent::PlanUpdated { explanation, plan, .. } => {
77            serde_json::json!({"type": "plan_updated", "explanation": explanation, "plan": plan})
78        }
79        RuntimeEvent::UserEvent {
80            event: UserEvent::Structured { event_type, data },
81            ..
82        } => {
83            serde_json::json!({"type": "user_event", "event_type": event_type, "data": data})
84        }
85        RuntimeEvent::UserEvent { .. } => serde_json::json!({"type": "other"}),
86        RuntimeEvent::RunCancelled { .. } => serde_json::json!({"type": "run_cancelled"}),
87        RuntimeEvent::RunFinished { .. } => serde_json::json!({"type": "run_finished"}),
88        RuntimeEvent::Checkpoint { .. } => serde_json::json!({"type": "checkpoint"}),
89    }
90}
91
92/// Convert a RuntimeEvent to a JSONL line.
93pub fn event_to_jsonl(event: &RuntimeEvent) -> String {
94    let value = event_to_value(event);
95    serde_json::to_string(&value).unwrap_or_else(|_| r#"{"type":"serialize_error"}"#.to_string())
96}