1use 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}
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}