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}
116
117#[cfg(test)]
118mod tests {
119    use super::*;
120    use agent_base::{ApprovalRequest, PlanItem, PlanStepStatus, RiskLevel, SessionId, UserEvent};
121    use std::io::Write;
122    use std::sync::{Arc, Mutex};
123
124    struct SharedWriter {
125        inner: Arc<Mutex<Vec<u8>>>,
126    }
127
128    impl Write for SharedWriter {
129        fn write(&mut self, data: &[u8]) -> std::io::Result<usize> {
130            self.inner.lock().unwrap().extend_from_slice(data);
131            Ok(data.len())
132        }
133        fn flush(&mut self) -> std::io::Result<()> { Ok(()) }
134    }
135
136    impl SharedWriter {
137        fn new() -> (Self, Arc<Mutex<Vec<u8>>>) {
138            let inner = Arc::new(Mutex::new(Vec::new()));
139            (Self { inner: inner.clone() }, inner)
140        }
141    }
142
143    fn session_id() -> SessionId {
144        SessionId { id: 1, external_id: None }
145    }
146
147    fn render_one(event: RuntimeEvent) -> Vec<String> {
148        let (writer, buf) = SharedWriter::new();
149        let mut r = JsonStreamRenderer::new(Box::new(writer));
150        r.render(event).unwrap();
151        drop(r);
152        let text = String::from_utf8(buf.lock().unwrap().clone()).unwrap();
153        if text.is_empty() { vec![] } else { text.lines().map(|l| l.to_string()).collect() }
154    }
155
156    fn render_and_finish(events: &[RuntimeEvent]) -> Vec<String> {
157        let (writer, buf) = SharedWriter::new();
158        let mut r = JsonStreamRenderer::new(Box::new(writer));
159        for e in events {
160            r.render(e.clone()).unwrap();
161        }
162        r.finish_turn().unwrap();
163        drop(r);
164        let text = String::from_utf8(buf.lock().unwrap().clone()).unwrap();
165        text.lines().map(|l| l.to_string()).collect()
166    }
167
168    #[test]
169    fn test_text_delta_produces_valid_json() {
170        let lines = render_one(RuntimeEvent::TextDelta {
171            session_id: session_id(), text: "hello".into(),
172        });
173        assert_eq!(lines.len(), 1);
174        let v: serde_json::Value = serde_json::from_str(&lines[0]).unwrap();
175        assert_eq!(v["type"], "text_delta");
176        assert_eq!(v["text"], "hello");
177    }
178
179    #[test]
180    fn test_thought_delta_produces_valid_json() {
181        let lines = render_one(RuntimeEvent::ThoughtDelta {
182            session_id: session_id(), text: "thinking...".into(),
183        });
184        let v: serde_json::Value = serde_json::from_str(&lines[0]).unwrap();
185        assert_eq!(v["type"], "thought_delta");
186    }
187
188    #[test]
189    fn test_tool_call_started_parses_args() {
190        let lines = render_one(RuntimeEvent::ToolCallStarted {
191            session_id: session_id(),
192            tool_name: "shell".into(),
193            args_json: r#"{"cmd":"ls"}"#.into(),
194        });
195        let v: serde_json::Value = serde_json::from_str(&lines[0]).unwrap();
196        assert_eq!(v["type"], "tool_call_started");
197        assert_eq!(v["tool"], "shell");
198        assert_eq!(v["args"]["cmd"], "ls");
199    }
200
201    #[test]
202    fn test_tool_call_finished_produces_valid_json() {
203        let lines = render_one(RuntimeEvent::ToolCallFinished {
204            session_id: session_id(),
205            tool_name: "shell".into(),
206            summary: "done".into(),
207        });
208        let v: serde_json::Value = serde_json::from_str(&lines[0]).unwrap();
209        assert_eq!(v["type"], "tool_call_finished");
210        assert_eq!(v["tool"], "shell");
211        assert_eq!(v["summary"], "done");
212    }
213
214    #[test]
215    fn test_awaiting_approval_produces_valid_json() {
216        let lines = render_one(RuntimeEvent::AwaitingApproval {
217            session_id: session_id(),
218            request: ApprovalRequest {
219                title: "Delete".into(),
220                message: "Dangerous".into(),
221                action_key: None,
222                risk_level: RiskLevel::Destructive,
223                raw: None,
224            },
225        });
226        let v: serde_json::Value = serde_json::from_str(&lines[0]).unwrap();
227        assert_eq!(v["type"], "approval_request");
228        assert_eq!(v["title"], "Delete");
229    }
230
231    #[test]
232    fn test_plan_updated_produces_valid_json() {
233        let lines = render_one(RuntimeEvent::PlanUpdated {
234            session_id: session_id(),
235            objective: "test".into(),
236            explanation: Some("step 1 done".into()),
237            plan: vec![PlanItem { step: "Step 1".into(), status: PlanStepStatus::Completed }],
238        });
239        let v: serde_json::Value = serde_json::from_str(&lines[0]).unwrap();
240        assert_eq!(v["type"], "plan_updated");
241    }
242
243    #[test]
244    fn test_user_event_structured() {
245        let lines = render_one(RuntimeEvent::UserEvent {
246            session_id: session_id(),
247            event: UserEvent::Structured {
248                event_type: "custom".into(),
249                data: serde_json::json!({"key": "value"}),
250            },
251        });
252        let v: serde_json::Value = serde_json::from_str(&lines[0]).unwrap();
253        assert_eq!(v["type"], "user_event");
254        assert_eq!(v["event_type"], "custom");
255        assert_eq!(v["data"]["key"], "value");
256    }
257
258    #[test]
259    fn test_user_event_progress_ignored() {
260        let lines = render_one(RuntimeEvent::UserEvent {
261            session_id: session_id(),
262            event: UserEvent::Progress { text: "loading...".into() },
263        });
264        assert!(lines.is_empty());
265    }
266
267    #[test]
268    fn test_run_finished_no_output() {
269        let lines = render_one(RuntimeEvent::RunFinished { session_id: session_id() });
270        assert!(lines.is_empty());
271    }
272
273    #[test]
274    fn test_run_cancelled_produces_valid_json() {
275        let lines = render_one(RuntimeEvent::RunCancelled { session_id: session_id() });
276        let v: serde_json::Value = serde_json::from_str(&lines[0]).unwrap();
277        assert_eq!(v["type"], "run_cancelled");
278    }
279
280    #[test]
281    fn test_finish_turn_emits_summary() {
282        let lines = render_and_finish(&[
283            RuntimeEvent::TextDelta { session_id: session_id(), text: "hello".into() },
284        ]);
285        let last: serde_json::Value = serde_json::from_str(lines.last().unwrap()).unwrap();
286        assert_eq!(last["type"], "turn_finished");
287        assert!(last["duration_ms"].as_u64().is_some());
288        assert_eq!(last["tool_call_count"], 0);
289    }
290
291    #[test]
292    fn test_tool_call_count_incremented() {
293        let lines = render_and_finish(&[
294            RuntimeEvent::ToolCallStarted {
295                session_id: session_id(), tool_name: "a".into(), args_json: "{}".into(),
296            },
297            RuntimeEvent::ToolCallStarted {
298                session_id: session_id(), tool_name: "b".into(), args_json: "{}".into(),
299            },
300        ]);
301        let last: serde_json::Value = serde_json::from_str(lines.last().unwrap()).unwrap();
302        assert_eq!(last["tool_call_count"], 2);
303    }
304
305    #[test]
306    fn test_assistant_text_accumulated() {
307        let lines = render_and_finish(&[
308            RuntimeEvent::TextDelta { session_id: session_id(), text: "Hello ".into() },
309            RuntimeEvent::TextDelta { session_id: session_id(), text: "World".into() },
310        ]);
311        let last: serde_json::Value = serde_json::from_str(lines.last().unwrap()).unwrap();
312        assert_eq!(last["assistant_text"], "Hello World");
313    }
314}