1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
145
146
147
148
149
150
151
152
153
154
155
156
157
158
159
160
161
162
163
164
165
166
167
168
169
170
171
172
173
174
175
176
177
178
179
180
181
182
183
184
use std::pin::Pin;
use std::sync::Arc;
use std::sync::atomic::Ordering;
use futures::{Stream, StreamExt};
use crate::error::AgentError;
use crate::loop_::AgentEvent;
use crate::types::{AgentMessage, AgentResult, LlmMessage, StopReason, Usage};
use super::Agent;
impl Agent {
/// Collect a stream to completion, updating agent state along the way.
pub(super) async fn collect_stream(
&mut self,
mut stream: Pin<Box<dyn Stream<Item = AgentEvent> + Send>>,
) -> Result<AgentResult, AgentError> {
let mut all_messages: Vec<AgentMessage> = Vec::new();
let mut stop_reason = StopReason::Stop;
let mut usage = Usage::default();
let mut cost = crate::types::Cost::default();
let mut error: Option<String> = None;
let mut transfer_signal: Option<crate::transfer::TransferSignal> = None;
while let Some(event) = stream.next().await {
self.dispatch_event(&event);
self.update_state_from_event(&event);
match event {
AgentEvent::TransferInitiated { signal } => {
transfer_signal = Some(signal);
stop_reason = StopReason::Transfer;
}
AgentEvent::TurnEnd {
assistant_message,
tool_results,
..
} => {
// Preserve Transfer stop reason set by TransferInitiated event
if transfer_signal.is_none() {
stop_reason = assistant_message.stop_reason;
}
usage += assistant_message.usage.clone();
cost += assistant_message.cost.clone();
if let Some(ref err) = assistant_message.error_message {
error = Some(err.clone());
}
let assistant_llm = LlmMessage::Assistant(assistant_message);
// Write the completed turn back to observable state so
// that history is not lost if this future is dropped
// before completion (e.g. a `select!` timeout). The
// `AgentEnd` arm below replaces `state.messages`
// wholesale, so these incremental pushes never duplicate.
self.state
.messages
.push(AgentMessage::Llm(assistant_llm.clone()));
all_messages.push(AgentMessage::Llm(assistant_llm));
for tr in tool_results {
self.state
.messages
.push(AgentMessage::Llm(LlmMessage::ToolResult(tr.clone())));
all_messages.push(AgentMessage::Llm(LlmMessage::ToolResult(tr)));
}
}
AgentEvent::AgentEnd { messages } => match Arc::try_unwrap(messages) {
Ok(returned) => {
self.state.messages = returned;
}
Err(messages) => {
self.state.messages = clone_messages(messages.as_ref());
}
},
_ => {}
}
}
// If the stream ended without `AgentEnd` (channel closed early),
// `state.messages` already holds the pre-run snapshot plus every
// turn written back above — nothing to restore.
self.state.is_running = false;
self.loop_active.store(false, Ordering::Release);
self.pending_message_snapshot.clear();
self.loop_context_snapshot.clear();
self.state.error.clone_from(&error);
self.idle_notify.notify_waiters();
Ok(AgentResult {
messages: all_messages,
stop_reason,
usage,
cost,
error,
transfer_signal,
})
}
/// Processes a streaming event, updating [`Agent::state`] and notifying subscribers.
pub fn handle_stream_event(&mut self, event: &AgentEvent) {
self.dispatch_event(event);
self.update_state_from_event(event);
match event {
AgentEvent::TurnEnd {
assistant_message,
tool_results,
..
} => {
// Write the completed turn back to observable state so that
// dropping the stream before `AgentEnd` keeps every turn the
// host already processed. The `AgentEnd` arm below replaces
// `state.messages` wholesale, so these incremental pushes
// never duplicate.
self.state
.messages
.push(AgentMessage::Llm(LlmMessage::Assistant(
assistant_message.clone(),
)));
for tr in tool_results {
self.state
.messages
.push(AgentMessage::Llm(LlmMessage::ToolResult(tr.clone())));
}
// Capture terminal error so it survives through AgentEnd.
if let Some(ref err) = assistant_message.error_message {
self.state.error = Some(err.clone());
}
}
AgentEvent::AgentEnd { messages } => {
self.state.messages = clone_messages(messages.as_ref());
self.pending_message_snapshot.clear();
self.loop_context_snapshot.clear();
// Preserve terminal error — do not clear self.state.error.
self.idle_notify.notify_waiters();
}
_ => {}
}
}
fn update_state_from_event(&mut self, event: &AgentEvent) {
match event {
AgentEvent::MessageStart => {
self.state.stream_message = None;
}
AgentEvent::MessageEnd { message } => {
self.state.stream_message =
Some(AgentMessage::Llm(LlmMessage::Assistant(message.clone())));
}
AgentEvent::ToolExecutionStart { id, .. } => {
self.state.pending_tool_calls.insert(id.clone());
}
AgentEvent::TurnEnd { .. } => {
self.state.pending_tool_calls.clear();
self.state.stream_message = None;
}
AgentEvent::AgentEnd { .. } => {
self.state.is_running = false;
self.loop_active.store(false, Ordering::Release);
self.state.pending_tool_calls.clear();
self.state.stream_message = None;
}
_ => {}
}
}
}
fn clone_messages(messages: &[AgentMessage]) -> Vec<AgentMessage> {
messages
.iter()
.filter_map(|message| match message {
AgentMessage::Llm(llm) => Some(AgentMessage::Llm(llm.clone())),
AgentMessage::Custom(cm) => cm.clone_box().map_or_else(
|| {
tracing::warn!(
"CustomMessage {:?} does not support clone_box — dropped during state rebuild",
cm
);
None
},
|cloned| Some(AgentMessage::Custom(cloned)),
),
})
.collect()
}