Skip to main content

aether_core/core/
agent.rs

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/// Internal event type for merging LLM and tool result streams
32#[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/// Keys for the merged stream map. Tool-call IDs come from providers, so the
46/// typed key keeps them from colliding with reserved streams.
47#[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    /// Get a reference to the token tracker
137    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_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    /// Inject a continuation prompt when the LLM stops due to a resumable reason.
446    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(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}