Skip to main content

phi_agent/render/
json_stream.rs

1use std::io::{self, Write};
2
3use agent_base::{AgentResult, RuntimeEvent, UserEvent};
4use serde_json::{Value, json};
5
6use crate::render::EventRenderer;
7
8/// JSON stream renderer: outputs one JSON line per event (JSONL format).
9/// Suitable for IDE integrations and programmatic consumers.
10pub struct JsonStreamRenderer {
11    writer: Box<dyn Write + Send>,
12    turn_start: Option<std::time::Instant>,
13    tool_call_count: u32,
14    last_assistant_text: String,
15}
16
17impl JsonStreamRenderer {
18    /// Create a new JSON stream renderer writing to the given writer.
19    pub fn new(writer: Box<dyn Write + Send>) -> Self {
20        Self { writer, turn_start: None, tool_call_count: 0, last_assistant_text: String::new() }
21    }
22
23    /// Create a renderer that writes to stdout.
24    pub fn stdout() -> Self {
25        Self::new(Box::new(io::stdout()))
26    }
27
28    fn emit(&mut self, value: &Value) -> AgentResult<()> {
29        let line = serde_json::to_string(value)
30            .map_err(|e| agent_base::AgentError::internal(format!("JSON serialize error: {e}")))?;
31        writeln!(self.writer, "{}", line).map_err(|e| agent_base::AgentError::internal(format!("write error: {e}")))?;
32        Ok(())
33    }
34}
35
36impl EventRenderer for JsonStreamRenderer {
37    fn render(&mut self, event: RuntimeEvent) -> AgentResult<()> {
38        if self.turn_start.is_none() {
39            self.turn_start = Some(std::time::Instant::now());
40        }
41
42        match &event {
43            RuntimeEvent::ThoughtDelta { text, .. } => {
44                self.emit(&json!({ "type": "thought_delta", "text": text }))?;
45            },
46            RuntimeEvent::TextDelta { text, .. } => {
47                self.last_assistant_text.push_str(text);
48                self.emit(&json!({ "type": "text_delta", "text": text }))?;
49            },
50            RuntimeEvent::ToolCallStarted { tool_name, args_json, .. } => {
51                self.tool_call_count += 1;
52                let args: Value = serde_json::from_str(args_json).unwrap_or(Value::Null);
53                self.emit(&json!({
54                    "type": "tool_call_started",
55                    "tool": tool_name,
56                    "args": args,
57                }))?;
58            },
59            RuntimeEvent::ToolCallFinished { tool_name, summary, .. } => {
60                self.emit(&json!({
61                    "type": "tool_call_finished",
62                    "tool": tool_name,
63                    "summary": summary,
64                }))?;
65            },
66            RuntimeEvent::AwaitingApproval { request, .. } => {
67                self.emit(&json!({
68                    "type": "approval_request",
69                    "title": request.title,
70                    "risk": format!("{:?}", request.risk_level),
71                    "message": request.message,
72                }))?;
73            },
74            RuntimeEvent::PlanUpdated { explanation, plan, .. } => {
75                self.emit(&json!({
76                    "type": "plan_updated",
77                    "explanation": explanation,
78                    "plan": plan,
79                }))?;
80            },
81            RuntimeEvent::UserEvent { event: UserEvent::Structured { event_type, data }, .. } => {
82                self.emit(&json!({
83                    "type": "user_event",
84                    "event_type": event_type,
85                    "data": data,
86                }))?;
87            },
88            RuntimeEvent::UserEvent { .. } => {},
89            RuntimeEvent::Checkpoint { .. } => {},
90            RuntimeEvent::RunFinished { .. } => {},
91            RuntimeEvent::RunCancelled { .. } => {
92                self.emit(&json!({ "type": "run_cancelled" }))?;
93            },
94        }
95
96        Ok(())
97    }
98
99    fn finish_turn(&mut self) -> AgentResult<()> {
100        let duration_ms = self.turn_start.map(|s| s.elapsed().as_millis() as u64).unwrap_or(0);
101
102        self.emit(&json!({
103            "type": "turn_finished",
104            "duration_ms": duration_ms,
105            "tool_call_count": self.tool_call_count,
106            "assistant_text": self.last_assistant_text.trim(),
107        }))?;
108
109        self.turn_start = None;
110        self.tool_call_count = 0;
111        self.last_assistant_text.clear();
112
113        Ok(())
114    }
115}