1use crate::context::{
2 CompactionConfig, CompactionError, CompactionResult, Compactor, SessionUsageTracker, TokenTracker,
3};
4use crate::core::PromptCache;
5use crate::core::prompt_cache_key::derive_prompt_cache_key;
6use crate::core::queued_input::QueuedInput;
7pub use crate::core::retry_config::RetryConfig;
8use crate::core::tool_execution::{ToolAbortPolicy, ToolExecutionUpdate, ToolExecutions};
9use crate::events::{
10 AgentCommand, AgentEvent, AgentObserver, Command, CompactionId, CompactionOutcome, ContextEvent, LlmCallOutcome,
11 ModelEvent, StreamState, TaskOutcome, ToolEvent, TraceContext, TurnEvent, TurnOutcome, UserCommand,
12};
13use crate::mcp::McpHandle;
14use futures::Stream;
15use llm::{
16 AssistantReasoning, ChatMessage, Context, EncryptedReasoningContent, LlmCallPurpose, LlmError, LlmModel,
17 LlmResponse, MessageId, ModelIdentity, StopReason, StreamingModelProvider, TokenUsage, ToolCallError,
18 ToolCallRequest, ToolCallResult,
19};
20use mcp_utils::client::{CallToolError, CallToolOptions, ToolCallEvent};
21use std::collections::VecDeque;
22use std::pin::Pin;
23use std::sync::Arc;
24use std::time::Duration;
25use tokio::sync::mpsc;
26use tokio::time::sleep;
27use tokio_stream::StreamExt;
28use tokio_stream::StreamMap;
29use tokio_stream::wrappers::ReceiverStream;
30
31#[derive(Debug)]
33#[allow(clippy::large_enum_variant)]
34enum StreamEvent {
35 LlmRequestStarted { attempt: u32 },
36 Llm(Result<LlmResponse, LlmError>),
37 ToolExecution(ToolCallEvent),
38 Command(Command),
39 InputClosed,
40 Compaction(Result<CompactionResult, CompactionError>),
41}
42
43type EventStream = Pin<Box<dyn Stream<Item = StreamEvent> + Send>>;
44
45#[derive(Debug, Clone, PartialEq, Eq, Hash)]
48enum StreamKey {
49 Input,
50 Llm,
51 Compaction,
52 Tool(String),
53}
54
55pub(crate) struct AgentConfig {
56 pub llm: Arc<dyn StreamingModelProvider>,
57 pub context: Context,
58 pub mcp: Option<McpHandle>,
59 pub tool_timeout: Duration,
60 pub compaction_config: Option<CompactionConfig>,
61 pub auto_continue: AutoContinue,
62 pub retry_config: RetryConfig,
63 pub context_window: Option<u32>,
64 pub prompt_cache: PromptCache,
65 pub observers: Vec<Box<dyn AgentObserver>>,
66 pub session_usage: SessionUsageTracker,
67}
68
69pub struct Agent {
70 llm: Arc<dyn StreamingModelProvider>,
71 context: Context,
72 mcp: Option<McpHandle>,
73 message_tx: mpsc::Sender<AgentEvent>,
74 observers: Vec<Box<dyn AgentObserver>>,
75 streams: StreamMap<StreamKey, EventStream>,
76 tool_timeout: Duration,
77 token_tracker: TokenTracker,
78 compaction_config: Option<CompactionConfig>,
79 auto_continue: AutoContinue,
80 retry_config: RetryConfig,
81 tool_executions: ToolExecutions,
82 pending_inputs: VecDeque<QueuedInput>,
83 queued_inputs: VecDeque<QueuedInput>,
84 context_window: Option<u32>,
85 prompt_cache: PromptCache,
86 turn_active: bool,
87 llm_call_active: bool,
88 active_compaction: Option<CompactionId>,
89 active_model: Option<LlmModel>,
90 session_usage: SessionUsageTracker,
91}
92
93impl Agent {
94 pub(crate) fn new(
95 config: AgentConfig,
96 command_rx: mpsc::Receiver<Command>,
97 message_tx: mpsc::Sender<AgentEvent>,
98 ) -> Self {
99 let mut streams: StreamMap<StreamKey, EventStream> = StreamMap::new();
100 let input_stream = ReceiverStream::new(command_rx)
101 .map(StreamEvent::Command)
102 .chain(futures::stream::once(async { StreamEvent::InputClosed }));
103 streams.insert(StreamKey::Input, Box::pin(input_stream));
104
105 let context_limit = config.context_window.or_else(|| config.llm.context_window());
106
107 Self {
108 llm: config.llm,
109 context: config.context,
110 mcp: config.mcp,
111 message_tx,
112 observers: config.observers,
113 streams,
114 tool_timeout: config.tool_timeout,
115 token_tracker: TokenTracker::new(context_limit),
116 compaction_config: config.compaction_config,
117 auto_continue: config.auto_continue,
118 retry_config: config.retry_config,
119 tool_executions: ToolExecutions::default(),
120 pending_inputs: VecDeque::new(),
121 queued_inputs: VecDeque::new(),
122 context_window: config.context_window,
123 prompt_cache: config.prompt_cache,
124 turn_active: false,
125 llm_call_active: false,
126 active_compaction: None,
127 active_model: None,
128 session_usage: config.session_usage,
129 }
130 }
131
132 pub fn current_model_display_name(&self) -> String {
133 self.llm.display_name()
134 }
135
136 pub fn token_tracker(&self) -> &TokenTracker {
138 &self.token_tracker
139 }
140
141 pub async fn run(mut self) {
142 let mut state = IterationState::default();
143 let mut input_closed = false;
144 self.emit_tool_definitions().await;
145
146 while let Some((stream_key, event)) = self.streams.next().await {
147 match event {
148 StreamEvent::Command(Command::UserCommand(UserCommand::Cancel)) => {
149 self.on_user_cancel(&mut state).await;
150 }
151
152 StreamEvent::Command(Command::UserCommand(UserCommand::ClearContext)) => {
153 self.on_user_clear_context(&mut state).await;
154 }
155
156 StreamEvent::Command(Command::UserCommand(UserCommand::Text { message_id, content })) => {
157 if self.is_busy() {
158 self.queued_inputs.push_back(QueuedInput::User { message_id, content });
159 } else {
160 self.begin_turn(QueuedInput::User { message_id, content }, &mut state).await;
161 }
162 }
163
164 StreamEvent::Command(Command::AgentCommand(AgentCommand::SwitchModel(new_provider))) => {
165 self.on_switch_model(new_provider).await;
166 }
167
168 StreamEvent::Command(Command::AgentCommand(AgentCommand::UpdateTools(tools))) => {
169 self.context.set_tools(tools);
170 self.emit_tool_definitions().await;
171 }
172
173 StreamEvent::Command(Command::AgentCommand(AgentCommand::UpdateMcpInstructions { server, body })) => {
174 self.on_update_instruction(server, body).await;
175 }
176
177 StreamEvent::Command(Command::AgentCommand(AgentCommand::SetReasoningEffort(effort))) => {
178 self.context.set_reasoning_effort(effort.unwrap_or_default());
179 }
180
181 StreamEvent::Command(Command::AgentCommand(AgentCommand::ReplaceConversation(messages))) => {
182 self.on_replace_conversation(messages, &mut state).await;
183 }
184
185 StreamEvent::InputClosed => {
186 input_closed = true;
187 }
188
189 StreamEvent::LlmRequestStarted { attempt } => {
190 self.begin_chat_call(attempt).await;
191 }
192
193 StreamEvent::Llm(llm_event) => {
194 self.on_llm_event(llm_event, &mut state).await;
195 }
196
197 StreamEvent::ToolExecution(tool_event) => {
198 let StreamKey::Tool(tool_id) = stream_key else {
199 unreachable!("tool events must come from a tool stream")
200 };
201 self.on_tool_execution_event(tool_id, tool_event, &mut state).await;
202 }
203
204 StreamEvent::Compaction(result) => {
205 self.on_compaction_complete(result).await;
206 }
207 }
208
209 if state.is_complete(self.tool_executions.has_foreground())
210 && let Some(id) = state.current_message_id.take()
211 {
212 let iteration = std::mem::take(&mut state);
213 self.on_iteration_complete(id, iteration).await;
214 }
215
216 if input_closed && !self.turn_active && !self.is_busy() {
217 if self.tool_executions.is_empty() {
218 break;
219 }
220 self.abort_in_flight_work(ToolAbortPolicy::CancelAll).await;
221 }
222 }
223
224 tracing::debug!("Agent task shutting down - input channel closed");
225 }
226
227 async fn on_iteration_complete(&mut self, id: MessageId, iteration: IterationState) {
228 let IterationState {
229 message_content,
230 reasoning_summary_text,
231 encrypted_reasoning,
232 completed_tool_calls,
233 stop_reason,
234 ..
235 } = iteration;
236 let has_tool_calls = !completed_tool_calls.is_empty();
237 let has_content = !message_content.is_empty() || !reasoning_summary_text.is_empty() || has_tool_calls;
238 let should_auto_continue = self.auto_continue.should_continue(stop_reason.as_ref());
239
240 if has_content {
241 let reasoning = AssistantReasoning::from_parts(reasoning_summary_text.clone(), encrypted_reasoning);
242 self.context.push_assistant_turn(id.clone(), &message_content, reasoning, completed_tool_calls);
243
244 self.emit(AgentEvent::text(&id, &message_content, StreamState::Complete)).await;
245
246 if !reasoning_summary_text.is_empty() {
247 self.emit(AgentEvent::thought(&id, &reasoning_summary_text, StreamState::Complete)).await;
248 }
249 }
250
251 let has_queued_input = !self.queued_inputs.is_empty();
252 if has_queued_input || has_tool_calls {
253 self.auto_continue.reset();
254 self.start_next_turn().await;
255 } else if should_auto_continue {
256 self.auto_continue.advance();
257 tracing::info!(
258 "LLM stopped with {:?}, auto-continuing (attempt {}/{})",
259 stop_reason,
260 self.auto_continue.count,
261 self.auto_continue.max
262 );
263
264 self.inject_continuation_prompt(stop_reason.as_ref()).await;
265 self.start_next_turn().await;
266 } else {
267 tracing::debug!("LLM completed turn with stop reason: {:?}", stop_reason);
268 self.auto_continue.reset();
269 self.finish_turn(TurnOutcome::Completed).await;
270 }
271 }
272
273 async fn start_next_turn(&mut self) {
274 debug_assert!(self.pending_inputs.is_empty());
275 self.pending_inputs.append(&mut self.queued_inputs);
276 if self.compaction_needed() {
277 self.begin_compaction().await;
278 } else {
279 self.start_chat_turn().await;
280 }
281 }
282
283 async fn start_chat_turn(&mut self) {
284 self.commit_pending_inputs().await;
285 self.start_llm_stream(None, 0).await;
286 }
287
288 async fn on_user_cancel(&mut self, state: &mut IterationState) {
289 self.abort_in_flight_work(ToolAbortPolicy::PreserveBackgroundAcknowledgements).await;
290 self.commit_pending_inputs().await;
291 self.queued_inputs.retain(|input| matches!(input, QueuedInput::TaskOutcome(_)));
292 self.commit_queued_inputs().await;
293 *state = IterationState::default();
294 self.finish_turn(TurnOutcome::Cancelled).await;
295 }
296
297 async fn discard_in_flight_work(&mut self, state: &mut IterationState) {
298 self.abort_in_flight_work(ToolAbortPolicy::CancelAll).await;
299 self.pending_inputs.clear();
300 self.queued_inputs.clear();
301 self.auto_continue.reset();
302 *state = IterationState::default();
303 }
304
305 async fn on_user_clear_context(&mut self, state: &mut IterationState) {
306 self.discard_in_flight_work(state).await;
307 self.context.clear_conversation();
308 self.token_tracker.reset_current_usage();
309 self.emit(AgentEvent::Context(ContextEvent::Cleared)).await;
310 self.finish_turn(TurnOutcome::Cancelled).await;
311 }
312
313 async fn on_replace_conversation(&mut self, messages: Vec<ChatMessage>, state: &mut IterationState) {
314 self.discard_in_flight_work(state).await;
315 self.context.replace_conversation(messages);
316 self.emit(self.context_usage_message()).await;
317 self.finish_turn(TurnOutcome::Cancelled).await;
318 }
319
320 async fn begin_turn(&mut self, input: QueuedInput, state: &mut IterationState) {
321 *state = IterationState::default();
322 self.auto_continue.reset();
323 self.turn_active = true;
324 let content = input.content_blocks();
325 self.emit(AgentEvent::Turn(TurnEvent::Started { content })).await;
326 self.queued_inputs.push_back(input);
327 self.start_next_turn().await;
328 }
329
330 async fn enqueue_task_outcome(&mut self, outcome: TaskOutcome, state: &mut IterationState) {
331 let input = QueuedInput::TaskOutcome(Box::new(outcome));
332 if self.is_busy() {
333 self.queued_inputs.push_back(input);
334 } else {
335 self.begin_turn(input, state).await;
336 }
337 }
338
339 async fn on_update_instruction(&mut self, server: String, body: Option<String>) {
340 self.prompt_cache.update_mcp_instruction(server, body);
341 match self.prompt_cache.render().await {
342 Ok(content) => self.context.set_system_content(content),
343 Err(e) => tracing::warn!("Failed to rebuild system prompt after instructions update: {e}"),
344 }
345 }
346
347 async fn on_switch_model(&mut self, new_provider: Box<dyn StreamingModelProvider>) {
348 let previous = self.llm.display_name();
349 let new_context_limit = self.context_window.or_else(|| new_provider.context_window());
350 self.llm = Arc::from(new_provider);
351 self.token_tracker.reset_current_usage();
352 self.token_tracker.set_context_limit(new_context_limit);
353 let new = self.llm.display_name();
354 self.emit(AgentEvent::Model(ModelEvent::Switched { previous, new })).await;
355
356 self.emit(self.context_usage_message()).await;
357 }
358
359 async fn start_llm_stream(&mut self, delay: Option<Duration>, attempt: u32) {
360 self.refresh_prompt_cache_key();
361 self.streams.remove(&StreamKey::Llm);
362 let stream: EventStream = match delay {
363 None => {
364 self.begin_chat_call(attempt).await;
365 Box::pin(self.llm.stream_response(&self.context).map(StreamEvent::Llm))
366 }
367 Some(delay) => {
368 self.emit(AgentEvent::Turn(TurnEvent::RetryScheduled {
369 purpose: LlmCallPurpose::Chat,
370 attempt,
371 max_attempts: self.retry_config.max_attempts,
372 delay_ms: u64::try_from(delay.as_millis()).unwrap_or(u64::MAX),
373 }))
374 .await;
375 let llm = Arc::clone(&self.llm);
376 let context = self.context.clone();
377 Box::pin(async_stream::stream! {
378 sleep(delay).await;
379 yield StreamEvent::LlmRequestStarted { attempt };
380 let mut inner = llm.stream_response(&context);
381 while let Some(item) = inner.next().await {
382 yield StreamEvent::Llm(item);
383 }
384 })
385 }
386 };
387 self.streams.insert(StreamKey::Llm, stream);
388 }
389
390 async fn on_llm_error(&mut self, error: LlmError, state: &mut IterationState) {
391 let will_retry = error.is_retryable() && state.retry_attempt < self.retry_config.max_attempts;
392 let outcome = LlmCallOutcome::from_llm_error(&error, will_retry);
393 let error_message = error.to_string();
394 self.finish_chat_call(outcome).await;
395
396 if !will_retry {
397 self.finish_turn(TurnOutcome::Failed { error: error_message }).await;
398 return;
399 }
400
401 state.retry_attempt += 1;
402 let delay = self.retry_config.compute_delay(state.retry_attempt);
403
404 tracing::warn!(
405 attempt = state.retry_attempt,
406 max_attempts = self.retry_config.max_attempts,
407 delay_ms = u64::try_from(delay.as_millis()).unwrap_or(u64::MAX),
408 error = %error,
409 "Retrying LLM request after transient failure"
410 );
411
412 self.tool_executions.retire_foreground();
413 self.start_llm_stream(Some(delay), state.retry_attempt).await;
414 }
415
416 fn is_busy(&self) -> bool {
417 self.streams.contains_key(&StreamKey::Llm)
418 || self.streams.contains_key(&StreamKey::Compaction)
419 || self.tool_executions.has_foreground()
420 }
421
422 async fn abort_in_flight_work(&mut self, tool_policy: ToolAbortPolicy) {
423 if self.llm_call_active {
424 self.finish_chat_call(LlmCallOutcome::Cancelled).await;
425 }
426 if self.streams.remove(&StreamKey::Compaction).is_some() {
427 let compaction_id = self.active_compaction.take().expect("active compaction stream has an identity");
428 self.emit(AgentEvent::Turn(TurnEvent::LlmCallEnded {
429 purpose: LlmCallPurpose::Compaction,
430 outcome: LlmCallOutcome::Cancelled,
431 }))
432 .await;
433 self.emit(AgentEvent::Context(ContextEvent::CompactionEnded {
434 compaction_id,
435 outcome: CompactionOutcome::Cancelled,
436 }))
437 .await;
438 }
439 self.streams.remove(&StreamKey::Llm);
440 for tool_id in self.tool_executions.abort(&tool_policy) {
441 self.streams.remove(&StreamKey::Tool(tool_id));
442 }
443 }
444
445 async fn inject_continuation_prompt(&mut self, stop_reason: Option<&StopReason>) {
447 let reason = stop_reason.map_or_else(|| "Unknown".to_string(), |reason| format!("{reason:?}"));
448 let message_id = MessageId::new();
449 let content = vec![llm::ContentBlock::text(format!(
450 "<system-notification>The LLM API stopped with reason '{reason}'. Continue from where you left off and finish your task.</system-notification>"
451 ))];
452 self.context.add_message(ChatMessage::user_with_id(message_id.clone(), content.clone()));
453 self.emit(AgentEvent::Turn(TurnEvent::AutoContinue {
454 attempt: self.auto_continue.count,
455 max_attempts: self.auto_continue.max,
456 message_id,
457 content,
458 }))
459 .await;
460 }
461
462 async fn on_llm_event(&mut self, result: Result<LlmResponse, LlmError>, state: &mut IterationState) {
463 use LlmResponse::{
464 Done, EncryptedReasoning, Error, Reasoning, Start, Text, ToolRequestArg, ToolRequestComplete,
465 ToolRequestStart, Usage,
466 };
467
468 let response = match result {
469 Ok(response) => response,
470 Err(e) => {
471 self.on_llm_error(e, state).await;
472 return;
473 }
474 };
475
476 match response {
477 Start => state.on_llm_start(MessageId::new()),
478
479 Text { chunk } => {
480 self.handle_llm_text(chunk, state).await;
481 }
482
483 Reasoning { chunk } => {
484 state.reasoning_summary_text.push_str(&chunk);
485 if let Some(id) = state.current_message_id.clone() {
486 self.emit(AgentEvent::thought(&id, &chunk, StreamState::Partial)).await;
487 }
488 }
489
490 EncryptedReasoning { id, content } => {
491 if let Some(model) = self.active_model.clone() {
492 state.encrypted_reasoning = Some(EncryptedReasoningContent { id, model, content });
493 }
494 }
495
496 ToolRequestStart { id, name } => {
497 let request = ToolCallRequest { id, name, arguments: String::new() };
498 self.emit(AgentEvent::Tool(ToolEvent::Call { request })).await;
499 }
500
501 ToolRequestArg { id, chunk } => {
502 self.emit(AgentEvent::Tool(ToolEvent::CallUpdate { tool_call_id: id, chunk })).await;
503 }
504
505 ToolRequestComplete { tool_call } => {
506 self.handle_tool_completion(tool_call).await;
507 }
508
509 Done { stop_reason } => {
510 state.llm_done = true;
511 state.stop_reason = stop_reason;
512 self.finish_chat_call(LlmCallOutcome::Completed {
513 stop_reason: state.stop_reason.clone(),
514 usage: state.call_usage.take(),
515 })
516 .await;
517 }
518
519 Error { message } => {
520 self.finish_chat_call(LlmCallOutcome::failed(message.clone(), false)).await;
521 self.finish_turn(TurnOutcome::Failed { error: message }).await;
522 }
523
524 Usage { tokens: sample } => {
525 self.handle_llm_usage(sample, state).await;
526 }
527 }
528 }
529
530 async fn handle_llm_text(&mut self, chunk: String, state: &mut IterationState) {
531 state.message_content.push_str(&chunk);
532
533 if let Some(id) = state.current_message_id.clone() {
534 self.emit(AgentEvent::text(&id, &chunk, StreamState::Partial)).await;
535 }
536 }
537
538 async fn handle_tool_completion(&mut self, tool_call: ToolCallRequest) {
539 let cancel = self.tool_executions.start(tool_call.clone());
540
541 let tool_id = tool_call.id.clone();
542 tracing::debug!("Tool execution started: {} ({})", tool_call.name, tool_id);
543 self.emit(AgentEvent::Tool(ToolEvent::ExecutionStarted {
544 tool_id: tool_id.clone(),
545 tool_name: tool_call.name.clone(),
546 }))
547 .await;
548
549 let Some(mcp) = self.mcp.clone() else {
550 let stream = futures::stream::once(async {
551 StreamEvent::ToolExecution(ToolCallEvent::Complete(Err(CallToolError::Unavailable {
552 message: "MCP runtime is not available".to_string(),
553 })))
554 });
555 self.streams.insert(StreamKey::Tool(tool_id), Box::pin(stream));
556 return;
557 };
558
559 let trace_context = self.observers.iter().find_map(|observer| observer.tool_trace_context(&tool_id));
560 let options = CallToolOptions {
561 timeout: self.tool_timeout,
562 meta: trace_context.as_ref().map(TraceContext::to_meta),
563 cancel,
564 };
565 let stream =
566 mcp.call_model_visible(tool_call.name, &tool_call.arguments, options).map(StreamEvent::ToolExecution);
567 self.streams.insert(StreamKey::Tool(tool_id), Box::pin(stream));
568 }
569
570 async fn handle_llm_usage(&mut self, sample: TokenUsage, state: &mut IterationState) {
571 state.call_usage = Some(sample);
572 self.token_tracker.record_usage(sample);
573 let ratio_pct = self.token_tracker.usage_ratio().map(|r| r * 100.0);
574 let remaining = self.token_tracker.tokens_remaining();
575 tracing::debug!(?sample, ?ratio_pct, ?remaining, "Token usage");
576
577 self.emit(self.context_usage_message()).await;
578 self.emit_session_usage(LlmCallPurpose::Chat, sample).await;
579 }
580
581 async fn emit_session_usage(&mut self, purpose: LlmCallPurpose, tokens: TokenUsage) {
582 let model = ModelIdentity::of(self.active_model.as_ref());
583 let event = self.session_usage.record(purpose, model, tokens);
584 self.emit(AgentEvent::SessionUsage(event)).await;
585 }
586
587 fn context_usage_message(&self) -> AgentEvent {
588 AgentEvent::Context(ContextEvent::UsageUpdated { usage: self.token_tracker.snapshot().clone() })
589 }
590
591 fn compaction_needed(&self) -> bool {
592 self.compaction_config.as_ref().is_some_and(|config| {
593 self.token_tracker.needs_compaction(self.context.estimated_token_count(), config.threshold)
594 })
595 }
596
597 async fn begin_compaction(&mut self) {
598 tracing::info!("Starting context compaction - {} messages", self.context.message_count());
599 let compaction_id = CompactionId::new();
600 self.active_compaction = Some(compaction_id.clone());
601 self.emit(AgentEvent::Context(ContextEvent::CompactionStarted {
602 compaction_id,
603 message_count: self.context.message_count(),
604 }))
605 .await;
606 let started = self.begin_llm_call(LlmCallPurpose::Compaction, 0);
607 self.emit(started).await;
608
609 let compactor = Compactor::new(self.llm.clone());
610 let context = self.context.clone();
611 let stream: EventStream =
612 Box::pin(futures::stream::once(async move { StreamEvent::Compaction(compactor.compact(context).await) }));
613 self.streams.insert(StreamKey::Compaction, stream);
614 }
615
616 async fn on_compaction_complete(&mut self, result: Result<CompactionResult, CompactionError>) {
617 let compaction_id = self.active_compaction.take().expect("completed compaction has an identity");
618 if let Ok(result) = &result
619 && let Some(usage) = result.usage
620 {
621 self.emit_session_usage(LlmCallPurpose::Compaction, usage).await;
622 }
623 let outcome = match &result {
624 Ok(result) => LlmCallOutcome::Completed { stop_reason: None, usage: result.usage },
625 Err(e) => LlmCallOutcome::failed(e.to_string(), false),
626 };
627 self.emit(AgentEvent::Turn(TurnEvent::LlmCallEnded { purpose: LlmCallPurpose::Compaction, outcome })).await;
628
629 match result {
630 Ok(result) => {
631 tracing::info!("Context compacted: {} messages removed", result.messages_removed);
632 let message_id = MessageId::new();
633 self.context = self.context.with_compacted_summary(message_id.clone(), &result.summary);
634 self.token_tracker.reset_current_usage();
635 self.emit(AgentEvent::Context(ContextEvent::CompactionResult {
636 compaction_id: compaction_id.clone(),
637 message_id,
638 summary: result.summary,
639 messages_removed: result.messages_removed,
640 }))
641 .await;
642 self.emit(AgentEvent::Context(ContextEvent::CompactionEnded {
643 compaction_id,
644 outcome: CompactionOutcome::Completed,
645 }))
646 .await;
647 }
648 Err(e) => {
649 tracing::warn!("Context compaction failed: {e}");
650 self.emit(AgentEvent::Context(ContextEvent::CompactionEnded {
651 compaction_id,
652 outcome: CompactionOutcome::Failed { error: e.to_string() },
653 }))
654 .await;
655 }
656 }
657
658 self.start_chat_turn().await;
659 }
660
661 async fn on_tool_execution_event(&mut self, tool_id: String, event: ToolCallEvent, state: &mut IterationState) {
662 match self.tool_executions.on_event(&tool_id, event) {
663 ToolExecutionUpdate::Event(event) => {
664 if let ToolEvent::SubAgentProgress { payload, .. } = &event
665 && let AgentEvent::SessionUsage(child) = &payload.event
666 {
667 let folded = self.session_usage.record_child(&payload.task_id, child.clone());
668 self.emit(AgentEvent::SessionUsage(folded)).await;
669 }
670 self.emit(AgentEvent::Tool(event)).await;
671 }
672 ToolExecutionUpdate::Completed { result, event } => {
673 self.streams.remove(&StreamKey::Tool(tool_id));
674 state.completed_tool_calls.push(result);
675 self.emit(AgentEvent::Tool(event)).await;
676 }
677 ToolExecutionUpdate::TaskCreated { result, event } => {
678 state.completed_tool_calls.push(Ok(result));
679 self.emit(AgentEvent::Tool(event)).await;
680 }
681 ToolExecutionUpdate::TaskCompleted(outcome) => {
682 self.streams.remove(&StreamKey::Tool(tool_id));
683 self.enqueue_task_outcome(outcome, state).await;
684 }
685 ToolExecutionUpdate::TaskCancelled(outcome) => {
686 self.streams.remove(&StreamKey::Tool(tool_id));
687 self.record_task_outcome(outcome).await;
688 }
689 ToolExecutionUpdate::Retired => {
690 self.streams.remove(&StreamKey::Tool(tool_id));
691 }
692 ToolExecutionUpdate::Ignored => {
693 tracing::debug!(%tool_id, "Ignoring unexpected tool execution event");
694 }
695 }
696 }
697
698 async fn record_task_outcome(&mut self, outcome: TaskOutcome) {
699 self.context.add_message(outcome.context_message());
700 self.emit(AgentEvent::Tool(outcome.into())).await;
701 }
702
703 fn refresh_prompt_cache_key(&mut self) {
704 let key = derive_prompt_cache_key(self.llm.as_ref(), &self.context);
705 self.context.set_prompt_cache_key(Some(key));
706 }
707
708 async fn commit_pending_inputs(&mut self) {
709 let inputs = std::mem::take(&mut self.pending_inputs);
710 self.commit_inputs(inputs).await;
711 }
712
713 async fn commit_queued_inputs(&mut self) {
714 let inputs = std::mem::take(&mut self.queued_inputs);
715 self.commit_inputs(inputs).await;
716 }
717
718 async fn commit_inputs(&mut self, inputs: VecDeque<QueuedInput>) {
719 for input in inputs {
720 match input {
721 QueuedInput::User { message_id, content } => {
722 self.context.add_message(ChatMessage::user_with_id(message_id, content));
723 }
724 QueuedInput::TaskOutcome(outcome) => self.record_task_outcome(*outcome).await,
725 }
726 }
727 }
728
729 async fn emit_tool_definitions(&mut self) {
730 let tools = self.context.tools().clone();
731 if !tools.is_empty() {
732 self.emit(AgentEvent::Tool(ToolEvent::DefinitionsUpdated { tools })).await;
733 }
734 }
735
736 async fn emit(&mut self, message: AgentEvent) {
737 for observer in &mut self.observers {
738 observer.on_event(&message);
739 }
740
741 if let Err(e) = self.message_tx.send(message).await {
742 tracing::warn!("Failed to send agent message: {e:?}");
743 }
744 }
745
746 async fn finish_turn(&mut self, outcome: TurnOutcome) {
747 if std::mem::take(&mut self.turn_active) {
748 self.emit(AgentEvent::turn_ended(outcome)).await;
749 }
750 }
751
752 async fn begin_chat_call(&mut self, attempt: u32) {
753 self.llm_call_active = true;
754 let started = self.begin_llm_call(LlmCallPurpose::Chat, attempt);
755 if let Some(system_prompt) = self.context.system_content() {
756 for observer in &mut self.observers {
757 observer.on_system_prompt(system_prompt);
758 }
759 }
760 self.emit(started).await;
761 }
762
763 async fn finish_chat_call(&mut self, outcome: LlmCallOutcome) {
764 if std::mem::take(&mut self.llm_call_active) {
765 self.emit(AgentEvent::Turn(TurnEvent::LlmCallEnded { purpose: LlmCallPurpose::Chat, outcome })).await;
766 }
767 }
768
769 fn begin_llm_call(&mut self, purpose: LlmCallPurpose, attempt: u32) -> AgentEvent {
770 self.active_model = self.llm.model();
771 AgentEvent::Turn(TurnEvent::LlmCallStarted {
772 purpose,
773 model: ModelIdentity::of(self.active_model.as_ref()),
774 display_name: self.llm.display_name(),
775 attempt,
776 max_attempts: self.retry_config.max_attempts,
777 })
778 }
779}
780
781pub(crate) struct AutoContinue {
782 max: u32,
783 count: u32,
784}
785
786impl AutoContinue {
787 pub(crate) fn new(max: u32) -> Self {
788 Self { max, count: 0 }
789 }
790
791 fn reset(&mut self) {
792 self.count = 0;
793 }
794
795 fn should_continue(&self, stop_reason: Option<&StopReason>) -> bool {
796 matches!(stop_reason, Some(StopReason::Length)) && self.count < self.max
797 }
798
799 fn advance(&mut self) {
800 self.count += 1;
801 }
802}
803
804#[derive(Debug, Default)]
805struct IterationState {
806 current_message_id: Option<MessageId>,
807 message_content: String,
808 reasoning_summary_text: String,
809 encrypted_reasoning: Option<EncryptedReasoningContent>,
810 completed_tool_calls: Vec<Result<ToolCallResult, ToolCallError>>,
811 llm_done: bool,
812 stop_reason: Option<StopReason>,
813 retry_attempt: u32,
814 call_usage: Option<TokenUsage>,
815}
816
817impl IterationState {
818 fn on_llm_start(&mut self, message_id: MessageId) {
819 self.current_message_id = Some(message_id);
820 self.message_content.clear();
821 self.reasoning_summary_text.clear();
822 self.encrypted_reasoning = None;
823 self.stop_reason = None;
824 self.call_usage = None;
825 }
826
827 fn is_complete(&self, has_foreground_tools: bool) -> bool {
828 self.llm_done && !has_foreground_tools
829 }
830}