Skip to main content

pi/modes/print/
text.rs

1//! Text renderer: drain an [`AgentSessionEvent`] stream to stdout/stderr text.
2//!
3//! Mirrors the text-mode branch of `.references/pi/packages/coding-agent/src/
4//! modes/print-mode.ts`: the session runs to completion, then the final
5//! assistant message's text blocks are written to stdout (one line each). An
6//! `error` or `aborted` stop reason writes the failure text to stderr and
7//! yields a nonzero exit code. Every other event variant (thinking deltas,
8//! tool-execution lifecycle, compaction, retry, queue updates, …) is consumed
9//! losslessly so fragmented streaming updates coalesce into the final message.
10//!
11//! The renderer is split into a synchronous, allocation-light state machine
12//! ([`TextRenderer`]) and an async drain ([`render_text`]) so tests can feed
13//! fragmented events directly and assert output without touching real I/O.
14
15use 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/// Outcome of draining an event stream through the text renderer.
26#[derive(Clone, Copy, Debug, Eq, PartialEq)]
27pub struct TextOutcome {
28    /// Process exit code (`0` on success, `1` on error / abort).
29    pub exit_code: i32,
30}
31
32impl TextOutcome {
33    /// Success exit code.
34    pub const SUCCESS: i32 = 0;
35    /// Failure exit code (error or aborted stop reason).
36    pub const FAILURE: i32 = 1;
37}
38
39/// Stateful text-mode renderer.
40///
41/// Tracks the latest assistant message observed across the whole stream so that
42/// multiple prompts, retries, and compaction rounds collapse to a single final
43/// answer. No output is produced until [`finish`](Self::finish); intermediate
44/// streaming events only update internal state.
45#[derive(Default)]
46pub struct TextRenderer {
47    last_assistant: Option<AssistantMessage>,
48}
49
50impl TextRenderer {
51    /// Create an empty renderer.
52    #[must_use]
53    pub fn new() -> Self {
54        Self::default()
55    }
56
57    /// Observe one event.
58    ///
59    /// Every [`AgentSessionEvent`] variant is handled. Authoritative assistant
60    /// snapshots come from turn-end, message-end, message-update, and agent-end
61    /// so fragmented streaming and multi-turn runs coalesce into a single final
62    /// answer. Other variants are consumed without side effects so the renderer
63    /// never panics on thinking, tool, compaction, retry, or queue events.
64    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            // Lifecycle / tool / session events update no text state.
78            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    /// Borrow the tracked final assistant message, if any.
107    #[must_use]
108    pub fn last_assistant(&self) -> Option<&AssistantMessage> {
109        self.last_assistant.as_ref()
110    }
111
112    /// Emit the final text (or error) and return the exit code.
113    ///
114    /// # Errors
115    ///
116    /// Propagates sink write failures.
117    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            // No assistant message observed (empty run / no prompt).
123            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
152/// Drain `events` through a fresh [`TextRenderer`] and emit the final output.
153///
154/// # Errors
155///
156/// Propagates sink write failures from [`TextRenderer::finish`].
157pub 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
170/// Return the assistant message inside a transcript message, if any.
171fn 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
181/// Wire string for a stop reason, matching the TypeScript literal.
182fn 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        // An agent_end with will_retry=true should not prematurely finalize;
334        // a later agent_end with the real answer wins.
335        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        // Smoke: every variant must be handleable without panic.
369        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}