1use std::io;
16
17use futures::Stream;
18use futures::StreamExt;
19use pi_agent::AgentMessage;
20use pi_ai::{AssistantMessage, Message, StopReason};
21
22use super::PrintSink;
23use crate::core::agent_session::AgentSessionEvent;
24
25#[derive(Clone, Copy, Debug, Eq, PartialEq)]
27pub struct TextOutcome {
28 pub exit_code: i32,
30}
31
32impl TextOutcome {
33 pub const SUCCESS: i32 = 0;
35 pub const FAILURE: i32 = 1;
37}
38
39#[derive(Default)]
46pub struct TextRenderer {
47 last_assistant: Option<AssistantMessage>,
48}
49
50impl TextRenderer {
51 #[must_use]
53 pub fn new() -> Self {
54 Self::default()
55 }
56
57 pub fn handle(&mut self, event: &AgentSessionEvent) {
65 match event {
66 AgentSessionEvent::TurnEnd { message, .. }
67 | AgentSessionEvent::MessageStart { message }
68 | AgentSessionEvent::MessageUpdate { message, .. }
69 | AgentSessionEvent::MessageEnd { message } => {
70 self.record_assistant(message);
71 }
72 AgentSessionEvent::AgentEnd { messages, .. } => {
73 for message in messages {
74 self.record_assistant(message);
75 }
76 }
77 AgentSessionEvent::AgentStart
79 | AgentSessionEvent::SessionBeforeSwitch { .. }
80 | AgentSessionEvent::SessionBeforeFork { .. }
81 | AgentSessionEvent::SessionStart { .. }
82 | AgentSessionEvent::SessionShutdown { .. }
83 | AgentSessionEvent::ModelSelect { .. }
84 | AgentSessionEvent::TurnStart
85 | AgentSessionEvent::ToolExecutionStart { .. }
86 | AgentSessionEvent::ToolExecutionUpdate { .. }
87 | AgentSessionEvent::ToolExecutionEnd { .. }
88 | AgentSessionEvent::AgentSettled
89 | AgentSessionEvent::QueueUpdate { .. }
90 | AgentSessionEvent::CompactionStart { .. }
91 | AgentSessionEvent::CompactionEnd { .. }
92 | AgentSessionEvent::EntryAppended { .. }
93 | AgentSessionEvent::SessionInfoChanged { .. }
94 | AgentSessionEvent::ThinkingLevelChanged { .. }
95 | AgentSessionEvent::AutoRetryStart { .. }
96 | AgentSessionEvent::AutoRetryEnd { .. } => {}
97 }
98 }
99
100 fn record_assistant(&mut self, message: &AgentMessage) {
101 if let Some(assistant) = assistant_of(message) {
102 self.last_assistant = Some(assistant.clone());
103 }
104 }
105
106 #[must_use]
108 pub fn last_assistant(&self) -> Option<&AssistantMessage> {
109 self.last_assistant.as_ref()
110 }
111
112 pub async fn finish<K>(&self, sink: &K) -> io::Result<i32>
118 where
119 K: PrintSink,
120 {
121 let Some(assistant) = self.last_assistant.as_ref() else {
122 sink.flush().await?;
124 return Ok(TextOutcome::SUCCESS);
125 };
126
127 match assistant.stop_reason {
128 StopReason::Error | StopReason::Aborted => {
129 let message = assistant.error_message.clone().unwrap_or_else(|| {
130 format!("Request {}", stop_reason_wire(assistant.stop_reason))
131 });
132 sink.write_stderr(&message).await?;
133 sink.write_stderr("\n").await?;
134 sink.flush().await?;
135 Ok(TextOutcome::FAILURE)
136 }
137 StopReason::Stop | StopReason::Length | StopReason::ToolUse => {
138 for content in &assistant.content {
139 if let pi_ai::AssistantContent::Text(text_block) = content {
140 let text = text_block.text.to_string();
141 sink.write_stdout(&text).await?;
142 sink.write_stdout("\n").await?;
143 }
144 }
145 sink.flush().await?;
146 Ok(TextOutcome::SUCCESS)
147 }
148 }
149 }
150}
151
152pub async fn render_text<S, K>(events: S, sink: &K) -> io::Result<i32>
158where
159 S: Stream<Item = AgentSessionEvent> + Send + Unpin,
160 K: PrintSink,
161{
162 let mut renderer = TextRenderer::new();
163 let mut events = events;
164 while let Some(event) = events.next().await {
165 renderer.handle(&event);
166 }
167 renderer.finish(sink).await
168}
169
170fn assistant_of(message: &AgentMessage) -> Option<&AssistantMessage> {
172 match message {
173 AgentMessage::Llm(boxed) => match boxed.as_ref() {
174 Message::Assistant(assistant) => Some(assistant),
175 Message::User(_) | Message::ToolResult(_) => None,
176 },
177 AgentMessage::Custom(_) => None,
178 }
179}
180
181fn stop_reason_wire(reason: StopReason) -> &'static str {
183 match reason {
184 StopReason::Stop => "stop",
185 StopReason::Length => "length",
186 StopReason::ToolUse => "toolUse",
187 StopReason::Error => "error",
188 StopReason::Aborted => "aborted",
189 }
190}
191
192#[cfg(test)]
193mod tests {
194 use super::*;
195 use crate::modes::print::BufferSink;
196 use futures::stream;
197 use pi_agent::user_text;
198 use pi_ai::{AssistantContent, TextContent};
199
200 type TestResult = Result<(), Box<dyn std::error::Error>>;
201
202 fn assistant_with(text: &str, reason: StopReason) -> AgentMessage {
203 let mut msg = AssistantMessage::new("api", "provider", "model", 2);
204 if !text.is_empty() {
205 msg.content
206 .push(AssistantContent::Text(TextContent::new(text)));
207 }
208 msg.stop_reason = reason;
209 AgentMessage::Llm(Box::new(Message::Assistant(msg)))
210 }
211
212 fn assistant_error(message: &str) -> AgentMessage {
213 let mut msg = AssistantMessage::new("api", "provider", "model", 2);
214 msg.stop_reason = StopReason::Error;
215 msg.error_message = Some(message.to_owned());
216 AgentMessage::Llm(Box::new(Message::Assistant(msg)))
217 }
218
219 #[tokio::test]
220 async fn text_renders_final_assistant_text() -> TestResult {
221 let events = vec![
222 AgentSessionEvent::AgentStart,
223 AgentSessionEvent::MessageEnd {
224 message: assistant_with("Hello\nWorld", StopReason::Stop),
225 },
226 AgentSessionEvent::AgentEnd {
227 messages: vec![assistant_with("Hello\nWorld", StopReason::Stop)],
228 will_retry: false,
229 },
230 ];
231 let sink = BufferSink::default();
232 let code = render_text(stream::iter(events), &sink).await?;
233 assert_eq!(code, 0);
234 assert_eq!(sink.stdout_string(), "Hello\nWorld\n");
235 assert!(sink.stderr_string().is_empty());
236 Ok(())
237 }
238
239 #[tokio::test]
240 async fn text_error_stop_reason_to_stderr_exit_one() -> TestResult {
241 let events = vec![AgentSessionEvent::AgentEnd {
242 messages: vec![assistant_error("rate limited")],
243 will_retry: false,
244 }];
245 let sink = BufferSink::default();
246 let code = render_text(stream::iter(events), &sink).await?;
247 assert_eq!(code, 1);
248 assert_eq!(sink.stderr_string(), "rate limited\n");
249 assert!(sink.stdout_string().is_empty());
250 Ok(())
251 }
252
253 #[tokio::test]
254 async fn text_aborted_without_error_message_uses_request_prefix() -> TestResult {
255 let mut msg = AssistantMessage::new("api", "provider", "model", 2);
256 msg.stop_reason = StopReason::Aborted;
257 let message = AgentMessage::Llm(Box::new(Message::Assistant(msg)));
258 let events = vec![AgentSessionEvent::AgentEnd {
259 messages: vec![message],
260 will_retry: false,
261 }];
262 let sink = BufferSink::default();
263 let code = render_text(stream::iter(events), &sink).await?;
264 assert_eq!(code, 1);
265 assert_eq!(sink.stderr_string(), "Request aborted\n");
266 Ok(())
267 }
268
269 #[tokio::test]
270 async fn text_fragmented_updates_coalesce() -> TestResult {
271 let mut partial = AssistantMessage::new("api", "provider", "model", 2);
272 partial
273 .content
274 .push(AssistantContent::Text(TextContent::new("Hel")));
275 partial
276 .content
277 .push(AssistantContent::Text(TextContent::new("lo")));
278 let final_msg = assistant_with("Hello", StopReason::Stop);
279 let events = vec![
280 AgentSessionEvent::MessageStart {
281 message: AgentMessage::Llm(Box::new(Message::Assistant(partial))),
282 },
283 AgentSessionEvent::MessageEnd {
284 message: final_msg.clone(),
285 },
286 AgentSessionEvent::AgentEnd {
287 messages: vec![final_msg],
288 will_retry: false,
289 },
290 ];
291 let sink = BufferSink::default();
292 let code = render_text(stream::iter(events), &sink).await?;
293 assert_eq!(code, 0);
294 assert_eq!(sink.stdout_string(), "Hello\n");
295 Ok(())
296 }
297
298 #[tokio::test]
299 async fn text_consumes_tool_and_thinking_events_without_output() -> TestResult {
300 let tool_start = AgentSessionEvent::ToolExecutionStart {
301 tool_call_id: "tc1".into(),
302 tool_name: "bash".into(),
303 args: serde_json::Map::new(),
304 };
305 let final_msg = assistant_with("done", StopReason::Stop);
306 let events = vec![
307 AgentSessionEvent::AgentStart,
308 tool_start,
309 AgentSessionEvent::AgentEnd {
310 messages: vec![final_msg],
311 will_retry: false,
312 },
313 ];
314 let sink = BufferSink::default();
315 let code = render_text(stream::iter(events), &sink).await?;
316 assert_eq!(code, 0);
317 assert_eq!(sink.stdout_string(), "done\n");
318 Ok(())
319 }
320
321 #[tokio::test]
322 async fn text_empty_stream_no_output_exit_zero() -> TestResult {
323 let sink = BufferSink::default();
324 let code = render_text(stream::empty::<AgentSessionEvent>(), &sink).await?;
325 assert_eq!(code, 0);
326 assert!(sink.stdout_string().is_empty());
327 assert!(sink.stderr_string().is_empty());
328 Ok(())
329 }
330
331 #[tokio::test]
332 async fn text_ignores_will_retry_agent_end_until_final() -> TestResult {
333 let retry_msg = assistant_error("transient");
336 let final_msg = assistant_with("recovered", StopReason::Stop);
337 let events = vec![
338 AgentSessionEvent::AgentEnd {
339 messages: vec![retry_msg],
340 will_retry: true,
341 },
342 AgentSessionEvent::AgentEnd {
343 messages: vec![final_msg],
344 will_retry: false,
345 },
346 ];
347 let sink = BufferSink::default();
348 let code = render_text(stream::iter(events), &sink).await?;
349 assert_eq!(code, 0);
350 assert_eq!(sink.stdout_string(), "recovered\n");
351 Ok(())
352 }
353
354 #[tokio::test]
355 async fn text_ignores_non_assistant_messages() -> TestResult {
356 let events = vec![AgentSessionEvent::MessageEnd {
357 message: user_text("hi", std::iter::empty()),
358 }];
359 let sink = BufferSink::default();
360 let code = render_text(stream::iter(events), &sink).await?;
361 assert_eq!(code, 0);
362 assert!(sink.stdout_string().is_empty());
363 Ok(())
364 }
365
366 #[tokio::test]
367 async fn text_handles_all_session_event_variants() -> TestResult {
368 let mut renderer = TextRenderer::new();
370 renderer.handle(&AgentSessionEvent::TurnStart);
371 renderer.handle(&AgentSessionEvent::TurnEnd {
372 message: assistant_with("x", StopReason::Stop),
373 tool_results: Vec::new(),
374 });
375 renderer.handle(&AgentSessionEvent::QueueUpdate {
376 steering: vec!["s".into()],
377 follow_up: vec!["f".into()],
378 });
379 renderer.handle(&AgentSessionEvent::CompactionStart {
380 reason: crate::core::agent_session::CompactionReason::Manual,
381 });
382 renderer.handle(&AgentSessionEvent::CompactionEnd {
383 reason: crate::core::agent_session::CompactionReason::Manual,
384 result: None,
385 aborted: false,
386 will_retry: false,
387 error_message: None,
388 });
389 renderer.handle(&AgentSessionEvent::AgentSettled);
390 renderer.handle(&AgentSessionEvent::AutoRetryStart {
391 attempt: 1,
392 max_attempts: 3,
393 delay_ms: 100,
394 error_message: "e".into(),
395 });
396 renderer.handle(&AgentSessionEvent::AutoRetryEnd {
397 success: true,
398 attempt: 1,
399 final_error: None,
400 });
401 renderer.handle(&AgentSessionEvent::SessionStart {
402 reason: crate::core::agent_session::SessionStartReason::Startup,
403 previous_session_file: None,
404 });
405 renderer.handle(&AgentSessionEvent::SessionShutdown {
406 reason: crate::core::agent_session::SessionShutdownReason::Quit,
407 target_session_file: None,
408 });
409 assert!(renderer.last_assistant.is_some());
410 Ok(())
411 }
412}