phi_agent/render/
json_stream.rs1use std::io::{self, Write};
2
3use agent_base::{AgentResult, RuntimeEvent, UserEvent};
4use serde_json::{Value, json};
5
6use crate::render::EventRenderer;
7
8pub 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 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 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}