pub(in crate::tui) mod delivery;
use self::delivery::send_tui_event;
use super::TuiEvent;
use crate::{
agent::AgentOutputSink,
output::{
ActivityEvent, ActivityId, ActivityKind, ActivityMetadata, ActivitySender, ActivityStatus,
OutputEvent,
},
tui::activity::ReasoningSummaryGroup,
};
use crossbeam_channel::Sender;
const ASSISTANT_ACTIVITY_DELTA_BYTES: usize = 96;
pub(crate) struct TuiOutputSink {
sender: Sender<TuiEvent>,
activity_sender: ActivitySender,
activity_namespace: String,
reasoning_summary_group: ReasoningSummaryGroup,
assistant_turn: usize,
reasoning_summary_sequence: usize,
assistant_started: bool,
pending_assistant_activity: String,
}
impl TuiOutputSink {
pub(crate) fn new(
sender: Sender<TuiEvent>,
activity_sender: ActivitySender,
activity_namespace: String,
) -> Self {
Self {
sender,
activity_sender,
activity_namespace,
reasoning_summary_group: ReasoningSummaryGroup::default(),
assistant_turn: 0,
reasoning_summary_sequence: 0,
assistant_started: false,
pending_assistant_activity: String::new(),
}
}
fn emit_primary_reasoning_summary_lines(
&mut self,
text: &str,
item_id: Option<&str>,
turn_id: Option<&str>,
) -> anyhow::Result<()> {
let id = if let Some(item_id) = item_id {
ActivityId::new(format!(
"reasoning/{}/{item_id}",
turn_id.unwrap_or("current")
))
} else {
self.reasoning_summary_sequence = self.reasoning_summary_sequence.saturating_add(1);
ActivityId::new(format!(
"reasoning/{}/legacy-{}",
turn_id.unwrap_or(&self.assistant_turn.to_string()),
self.reasoning_summary_sequence
))
};
let event = self
.reasoning_summary_group
.append_identity_event(id, None, text);
if let Some(event) = event {
self.activity_event(event)?;
}
Ok(())
}
fn assistant_activity_id(&self) -> ActivityId {
ActivityId::new(format!("assistant/{}", self.assistant_turn))
}
fn flush_assistant_activity_delta(&mut self) -> anyhow::Result<()> {
if self.pending_assistant_activity.is_empty() {
return Ok(());
}
let id = self.assistant_activity_id();
let preview = std::mem::take(&mut self.pending_assistant_activity);
self.activity_event(ActivityEvent::Delta { id, preview })
}
fn emit_assistant_activity(&mut self, text: &str) -> anyhow::Result<()> {
if !self.assistant_started {
self.assistant_started = true;
let id = self.assistant_activity_id();
self.activity_event(ActivityEvent::Started {
id,
parent_id: None,
kind: ActivityKind::Assistant,
status: ActivityStatus::Running,
metadata: ActivityMetadata::new("assistant"),
})?;
}
self.pending_assistant_activity.push_str(text);
if self.pending_assistant_activity.len() >= ASSISTANT_ACTIVITY_DELTA_BYTES
|| self.pending_assistant_activity.contains('\n')
{
self.flush_assistant_activity_delta()?;
}
Ok(())
}
fn finish_assistant_activity(&mut self) -> anyhow::Result<()> {
if self.assistant_started {
self.flush_assistant_activity_delta()?;
self.assistant_started = false;
let id = self.assistant_activity_id();
self.assistant_turn += 1;
self.activity_event(ActivityEvent::Finished {
id,
status: ActivityStatus::Success,
metadata: None,
})?;
}
Ok(())
}
}
pub(super) fn namespace_activity_event(
event: ActivityEvent,
activity_namespace: &str,
) -> ActivityEvent {
match event {
ActivityEvent::Started {
id,
parent_id,
kind,
status,
metadata,
} => ActivityEvent::Started {
id: namespace_activity_id(id, activity_namespace),
parent_id: parent_id.map(|id| namespace_activity_id(id, activity_namespace)),
kind,
status,
metadata,
},
ActivityEvent::Delta { id, preview } => ActivityEvent::Delta {
id: namespace_activity_id(id, activity_namespace),
preview,
},
ActivityEvent::UsageUpdate {
id,
current_tokens,
max_tokens,
reasoning_tokens,
source,
request_sequence,
} => ActivityEvent::UsageUpdate {
id: namespace_activity_id(id, activity_namespace),
current_tokens,
max_tokens,
reasoning_tokens,
source,
request_sequence,
},
ActivityEvent::FinalPreview {
id,
preview,
metadata,
status,
} => ActivityEvent::FinalPreview {
id: namespace_activity_id(id, activity_namespace),
preview,
metadata,
status,
},
ActivityEvent::Finished {
id,
status,
metadata,
} => ActivityEvent::Finished {
id: namespace_activity_id(id, activity_namespace),
status,
metadata,
},
}
}
fn namespace_activity_id(id: ActivityId, activity_namespace: &str) -> ActivityId {
if activity_namespace.is_empty() || id.as_str().starts_with(activity_namespace) {
id
} else {
ActivityId::new(format!("{}{}", activity_namespace, id.as_str()))
}
}
fn namespace_output_event(event: OutputEvent, activity_namespace: &str) -> OutputEvent {
match event {
OutputEvent::ToolStarted { call, label } => {
let mut call = *call;
if !call.id.trim().is_empty() {
call.id = namespace_activity_id(ActivityId::new(call.id), activity_namespace).0;
}
OutputEvent::ToolStarted {
call: Box::new(call),
label,
}
}
OutputEvent::ToolResult {
call,
result,
summary,
} => {
let mut call = *call;
if !call.id.trim().is_empty() {
call.id = namespace_activity_id(ActivityId::new(call.id), activity_namespace).0;
}
OutputEvent::ToolResult {
call: Box::new(call),
result,
summary,
}
}
OutputEvent::ProviderContextInjection { metadata } => {
OutputEvent::ProviderContextInjection { metadata }
}
event => event,
}
}
impl AgentOutputSink for TuiOutputSink {
fn assistant_delta(&mut self, text: &str) -> anyhow::Result<()> {
self.output_event(OutputEvent::AssistantDelta {
text: text.to_string(),
})
}
fn output_event(&mut self, event: OutputEvent) -> anyhow::Result<()> {
match &event {
OutputEvent::ThinkingSummaryComplete { text } => {
self.emit_primary_reasoning_summary_lines(text, None, None)?
}
OutputEvent::ThinkingSummaryCompleteIdentified {
text,
item_id,
turn_id,
} => self.emit_primary_reasoning_summary_lines(
text,
item_id.as_deref(),
turn_id.as_deref(),
)?,
OutputEvent::AssistantDelta { text } => self.emit_assistant_activity(text)?,
OutputEvent::AssistantComplete { .. } => self.finish_assistant_activity()?,
_ => {}
}
let event = namespace_output_event(event, &self.activity_namespace);
send_tui_event(&self.sender, TuiEvent::Output(event))
}
fn activity_event(&mut self, event: ActivityEvent) -> anyhow::Result<()> {
let event = namespace_activity_event(event, &self.activity_namespace);
send_tui_event(&self.sender, TuiEvent::Activity(event))
}
fn activity_sender(&self) -> Option<ActivitySender> {
Some(self.activity_sender.clone())
}
fn current_parent_activity_id(&self) -> Option<ActivityId> {
None
}
fn tool_block(&mut self, _block: &str) -> anyhow::Result<()> {
Ok(())
}
}
#[cfg(test)]
mod tests {
use super::*;
use crate::{
output::{ActivityKind, ActivityMetadata, ActivityStatus, ToolDisplaySummary, ToolStatus},
providers::ToolCall,
tools::ToolResult,
};
#[test]
fn namespaces_activity_ids_per_prompt_run() {
let first = namespace_activity_event(
ActivityEvent::Started {
id: ActivityId::new("turn-0/tool-0-read"),
parent_id: None,
kind: ActivityKind::Tool,
status: ActivityStatus::Running,
metadata: ActivityMetadata::new("read"),
},
"run-1/",
);
let second = namespace_activity_event(
ActivityEvent::Started {
id: ActivityId::new("turn-0/tool-0-read"),
parent_id: None,
kind: ActivityKind::Tool,
status: ActivityStatus::Running,
metadata: ActivityMetadata::new("read"),
},
"run-2/",
);
assert!(
matches!(first, ActivityEvent::Started { ref id, .. } if id.as_str() == "run-1/turn-0/tool-0-read")
);
assert!(
matches!(second, ActivityEvent::Started { ref id, .. } if id.as_str() == "run-2/turn-0/tool-0-read")
);
}
#[test]
fn usage_update_activity_id_is_namespaced() {
let namespaced = namespace_activity_event(
ActivityEvent::UsageUpdate {
id: ActivityId::new("call_subagents/g1"),
current_tokens: 12,
max_tokens: 128,
reasoning_tokens: Some(3),
source: crate::output::ContextUsageSource::ProviderExact,
request_sequence: 1,
},
"run-9/",
);
assert!(matches!(
namespaced,
ActivityEvent::UsageUpdate {
ref id,
current_tokens: 12,
max_tokens: 128,
reasoning_tokens: Some(3),
source: crate::output::ContextUsageSource::ProviderExact,
request_sequence: 1,
} if id.as_str() == "run-9/call_subagents/g1"
));
}
#[test]
fn tool_result_display_id_is_namespaced_without_mutating_raw_call() {
let raw_call = ToolCall {
id: "call_1".to_string(),
name: "read".to_string(),
arguments: serde_json::json!({"path":"file.txt"}),
};
let event = namespace_output_event(
OutputEvent::ToolResult {
call: Box::new(raw_call.clone()),
result: Box::new(ToolResult {
tool_name: "read".to_string(),
success: true,
content: "hello".to_string(),
metadata: serde_json::json!({}),
display: crate::tools::ToolResultDisplay::default(),
}),
summary: Box::new(ToolDisplaySummary {
tool_name: "read".to_string(),
status: ToolStatus::Success,
unicode_mark: "✓",
ascii_mark: "OK",
label: "read file.txt".to_string(),
metadata: vec![],
}),
},
"run-7/",
);
assert_eq!(raw_call.id, "call_1");
assert!(
matches!(event, OutputEvent::ToolResult { ref call, .. } if call.id == "run-7/call_1")
);
}
#[test]
fn activity_final_preview_is_namespaced_and_sent_as_critical_event() {
let namespaced = namespace_activity_event(
ActivityEvent::FinalPreview {
id: ActivityId::new("tool-1"),
preview: "final".to_string(),
metadata: Some(ActivityMetadata::new("tool")),
status: Some(ActivityStatus::Failed),
},
"run-9/",
);
assert!(matches!(
namespaced,
ActivityEvent::FinalPreview { ref id, ref preview, status: Some(ActivityStatus::Failed), .. }
if id.as_str() == "run-9/tool-1" && preview == "final"
));
let (sender, receiver) = crossbeam_channel::bounded(1);
let activity_sender: ActivitySender = std::sync::Arc::new(|_| {});
let mut sink = TuiOutputSink::new(sender, activity_sender, "run-9/".to_string());
sink.activity_event(ActivityEvent::FinalPreview {
id: ActivityId::new("tool-2"),
preview: "final".to_string(),
metadata: None,
status: Some(ActivityStatus::Success),
})
.unwrap();
let received = receiver.try_recv().unwrap();
assert!(matches!(
received,
TuiEvent::Activity(ActivityEvent::FinalPreview { ref id, .. })
if id.as_str() == "run-9/tool-2"
));
}
#[test]
fn thinking_summary_complete_emits_primary_reasoning_activity_lines() {
let (sender, receiver) = crossbeam_channel::bounded(2);
let activity_sender: ActivitySender = std::sync::Arc::new(|_| {});
let mut sink = TuiOutputSink::new(sender, activity_sender, "run-1/".to_string());
sink.output_event(OutputEvent::ThinkingSummaryComplete {
text: " checked inputs\n\nplanned fix ".to_string(),
})
.unwrap();
let events = receiver.try_iter().collect::<Vec<_>>();
assert_eq!(events.len(), 2);
assert!(matches!(
events[0],
TuiEvent::Activity(ActivityEvent::Started {
ref id,
parent_id: None,
kind: ActivityKind::Assistant,
status: ActivityStatus::Success,
ref metadata,
}) if id.as_str() == "run-1/reasoning/0/legacy-1"
&& metadata.label == "reasoning summary"
&& metadata.detail.as_deref() == Some("1. checked inputs\n2. planned fix")
&& metadata.fields.is_empty()
));
assert!(matches!(
events[1],
TuiEvent::Output(OutputEvent::ThinkingSummaryComplete { ref text })
if text == " checked inputs\n\nplanned fix "
));
}
#[test]
fn repeated_thinking_summary_complete_updates_same_primary_reasoning_activity() {
let (sender, receiver) = crossbeam_channel::bounded(4);
let activity_sender: ActivitySender = std::sync::Arc::new(|_| {});
let mut sink = TuiOutputSink::new(sender, activity_sender, "run-1/".to_string());
sink.output_event(OutputEvent::ThinkingSummaryComplete {
text: "checked inputs".to_string(),
})
.unwrap();
sink.output_event(OutputEvent::ThinkingSummaryComplete {
text: "planned fix\nverified output".to_string(),
})
.unwrap();
let events = receiver.try_iter().collect::<Vec<_>>();
let activity_events = events
.iter()
.filter_map(|event| match event {
TuiEvent::Activity(ActivityEvent::Started { id, metadata, .. }) => Some((
id.as_str(),
metadata.label.as_str(),
metadata.detail.as_deref(),
)),
_ => None,
})
.collect::<Vec<_>>();
assert_eq!(
activity_events,
vec![
(
"run-1/reasoning/0/legacy-1",
"reasoning summary",
Some("1. checked inputs")
),
(
"run-1/reasoning/0/legacy-2",
"reasoning summary",
Some("1. planned fix\n2. verified output")
),
]
);
}
#[test]
fn whitespace_only_thinking_summary_complete_forwards_output_without_activity() {
let (sender, receiver) = crossbeam_channel::bounded(1);
let activity_sender: ActivitySender = std::sync::Arc::new(|_| {});
let mut sink = TuiOutputSink::new(sender, activity_sender, "run-1/".to_string());
sink.output_event(OutputEvent::ThinkingSummaryComplete {
text: " \n\t\n ".to_string(),
})
.unwrap();
let events = receiver.try_iter().collect::<Vec<_>>();
assert_eq!(events.len(), 1);
assert!(matches!(
events[0],
TuiEvent::Output(OutputEvent::ThinkingSummaryComplete { ref text }) if text == " \n\t\n "
));
}
#[test]
fn assistant_delta_emits_activity_started_delta_and_finished() {
let (sender, receiver) = crossbeam_channel::bounded(8);
let activity_sender: ActivitySender = std::sync::Arc::new(|_| {});
let mut sink = TuiOutputSink::new(sender, activity_sender, "run-1/".to_string());
sink.output_event(OutputEvent::AssistantDelta {
text: "hello".to_string(),
})
.unwrap();
sink.output_event(OutputEvent::AssistantDelta {
text: " world".to_string(),
})
.unwrap();
sink.output_event(OutputEvent::AssistantComplete {
text: "hello world".to_string(),
})
.unwrap();
let events = receiver.try_iter().collect::<Vec<_>>();
let started = events.iter().find(|e| {
matches!(e, TuiEvent::Activity(ActivityEvent::Started { id, kind, status, .. })
if id.as_str() == "run-1/assistant/0"
&& *kind == ActivityKind::Assistant
&& *status == ActivityStatus::Running)
});
assert!(
started.is_some(),
"expected ActivityEvent::Started for assistant/0"
);
let deltas = events
.iter()
.filter_map(|e| match e {
TuiEvent::Activity(ActivityEvent::Delta { id, preview })
if id.as_str() == "run-1/assistant/0" =>
{
Some(preview.as_str())
}
_ => None,
})
.collect::<Vec<_>>();
assert_eq!(deltas, vec!["hello world"]);
let finished = events.iter().find(|e| {
matches!(e, TuiEvent::Activity(ActivityEvent::Finished { id, status, .. })
if id.as_str() == "run-1/assistant/0"
&& *status == ActivityStatus::Success)
});
assert!(finished.is_some(), "expected Finished for assistant/0");
assert!(
events.iter().any(|e| {
matches!(e, TuiEvent::Output(OutputEvent::AssistantComplete { text })
if text == "hello world")
}),
"expected output event forwarded"
);
}
#[test]
fn assistant_activity_creates_new_node_per_turn() {
let (sender, receiver) = crossbeam_channel::bounded(16);
let activity_sender: ActivitySender = std::sync::Arc::new(|_| {});
let mut sink = TuiOutputSink::new(sender, activity_sender, "run-1/".to_string());
sink.output_event(OutputEvent::AssistantDelta {
text: "turn one".to_string(),
})
.unwrap();
sink.output_event(OutputEvent::AssistantComplete {
text: "turn one".to_string(),
})
.unwrap();
sink.output_event(OutputEvent::AssistantDelta {
text: "turn two".to_string(),
})
.unwrap();
sink.output_event(OutputEvent::AssistantComplete {
text: "turn two".to_string(),
})
.unwrap();
let events = receiver.try_iter().collect::<Vec<_>>();
assert!(
events.iter().any(|e| {
matches!(e, TuiEvent::Activity(ActivityEvent::Started { id, .. })
if id.as_str() == "run-1/assistant/0")
}),
"expected node for assistant/0"
);
assert!(
events.iter().any(|e| {
matches!(e, TuiEvent::Activity(ActivityEvent::Started { id, .. })
if id.as_str() == "run-1/assistant/1")
}),
"expected node for assistant/1"
);
}
}