Skip to main content

everruns_engine/execution/
reason.rs

1//! ReasonAtom - Atom for LLM reasoning (model call)
2//!
3//! This atom handles:
4//! 1. Emitting reason.started event
5//! 2. Context preparation (loading message history, adding system message)
6//! 3. Fixing invalid context (e.g., missing tool_results for dangling tool calls)
7//! 4. LLM call with streaming support
8//! 5. Storing the assistant response
9//! 6. Emitting reason.completed event
10//! 7. Returning the result with tool calls (if any)
11//!
12//! NOTES from Python spec:
13//! - Context preparation includes loading message history, adding system message, editing context if needed
14//! - Before LLM call, invalid context (e.g. missing tool_results) should be fixed
15//! - LLM call should emit start/end events
16//! - Failure of the LLM call should be "normal" result, should user message that LLM call failed
17//! - Reason should be cancellable, cancellation should stop LLM call and exit with message
18
19use futures::StreamExt;
20use serde::{Deserialize, Serialize};
21use std::collections::{HashMap, HashSet};
22use std::sync::Arc;
23use std::time::Instant;
24use uuid::Uuid;
25
26fn add_compaction_cost(usage: &mut TokenUsage, compaction_cost: f64) {
27    let generation_cost = usage.effective_cost_usd();
28    usage.effective_cost_usd = Some(generation_cost.unwrap_or(0.0) + compaction_cost);
29    if let Some(actual_cost) = usage.actual_cost_usd.as_mut() {
30        *actual_cost += compaction_cost;
31    }
32}
33
34use super::ExecutionContext;
35use crate::annotation_hook::{collect_annotations, verify_annotations};
36use crate::capabilities::CapabilityRegistry;
37use crate::driver_registry::{LlmStreamEvent, Message, MessageContent, MessageRole};
38use crate::error::{AgentLoopError, Result};
39use crate::events::{
40    EventContext, EventRequest, LlmCompactionInfo, LlmGenerationData, LlmRetryInfo,
41    OutputMessageCompletedData, OutputMessageDeltaData, OutputMessageReplacedData,
42    OutputMessageStartedData, ReasonCompletedData, ReasonItemData, ReasonRecoveredData,
43    ReasonStartedData, ReasonThinkingCompletedData, ReasonThinkingDeltaData,
44    ReasonThinkingStartedData, RecoveryMode, TokenUsage, ToolCompletedData, ToolDefinitionSummary,
45};
46use crate::llm_retry::{
47    LlmRetryConfig, RetryMetadata, is_transient_error_message, remaining_retry_time,
48    reserve_retry_wait,
49};
50use crate::message::{ContentPart, RuntimeMessage, RuntimeMessageRole};
51use crate::message_retriever::MessageRetriever;
52use crate::output_guardrail::{
53    ArmedGuardrail, OutputGuardrailContext, PostGenerationOutputContext, evaluate_guardrails,
54    evaluate_post_generation_guardrails,
55};
56use crate::phase_effects::{PhaseEffectEmitter, PhaseEffectSink};
57use crate::runtime_context::{AssembledTurnContext, TurnContextRequest, TurnContextResolver};
58use crate::tool_types::{ToolCall, ToolDefinition};
59use crate::typed_id::{AgentId, HarnessId, MessageId, SessionId};
60use crate::{ErrorDisclosure, UserFacingError, UserFacingErrorContext};
61use crate::{
62    durability::DurableToolResultStore,
63    durability::PartialStreamState,
64    durability::PartialStreamStore,
65    file_services::{FileResolver, ResolvedFile},
66    image_services::ImageResolver,
67    image_services::ResolvedImage,
68};
69use everruns_provider::reasoning::{ReasoningContentPart, ReasoningText};
70
71mod compaction;
72mod error_policy;
73mod facts;
74mod finalized_calls;
75mod observability;
76mod output_hooks;
77mod provider_managed_compaction;
78mod reasoning_updates;
79mod request_controls;
80mod stream_state;
81mod transcript;
82
83use compaction::{
84    ProactiveCompactionContext, ReactiveCompactionContext, apply_proactive_compaction,
85    apply_reactive_compaction,
86};
87use error_policy::{filter_response_text, is_error_placeholder_message};
88#[cfg(test)]
89use observability::capability_usage_snapshot_records;
90use observability::{build_request_options, emit_capability_usage_snapshot};
91use output_hooks::{client_visible_guardrail_text, collect_output_hooks};
92use request_controls::resolve_request_controls;
93use stream_state::{
94    StreamReplayState, StreamTermination, advances_stall_deadline, append_guarded_thinking_delta,
95    inspect_guarded_reasoning_item, merge_retry_metadata,
96};
97use transcript::repair_dangling_tool_calls;
98
99fn unix_now_secs() -> u64 {
100    std::time::SystemTime::now()
101        .duration_since(std::time::UNIX_EPOCH)
102        .unwrap_or_default()
103        .as_secs()
104}
105
106/// Input for ReasonAtom
107#[derive(Debug, Clone, Serialize, Deserialize)]
108pub struct ReasonInput {
109    /// Atom execution context
110    pub context: ExecutionContext,
111    /// Harness ID for loading base configuration
112    pub harness_id: HarnessId,
113    /// Agent ID for loading configuration (optional)
114    #[serde(skip_serializing_if = "Option::is_none")]
115    pub agent_id: Option<AgentId>,
116    /// Organization ID for multi-tenancy tracking
117    #[serde(default)]
118    pub org_id: i64,
119    /// MCP tool definitions from agent's MCP capabilities (pre-resolved)
120    /// These are passed from the control-plane since MCP capabilities
121    /// are not in the CapabilityRegistry.
122    #[serde(default)]
123    pub mcp_tool_definitions: Vec<ToolDefinition>,
124    /// Previous LLM response ID for stateful continuation.
125    /// Enables server-side context caching across reason iterations.
126    #[serde(skip_serializing_if = "Option::is_none")]
127    pub previous_response_id: Option<String>,
128    /// Current iteration number within this turn (1-based).
129    /// Used for output.message.started events so UI can show progress.
130    #[serde(default = "default_iteration")]
131    pub iteration: u32,
132}
133
134fn default_iteration() -> u32 {
135    1
136}
137
138/// Internal continuations accounted as part of one scheduled Reason activity.
139#[derive(Debug, Clone, Serialize, Deserialize)]
140pub struct NativeExecutionCounts {
141    pub llm_calls: u32,
142    pub tool_calls: u32,
143}
144
145/// Result of the ReasonAtom
146#[derive(Debug, Clone, Default, Serialize, Deserialize)]
147pub struct ReasonResult {
148    #[serde(default, skip_serializing_if = "Option::is_none")]
149    pub native_counts: Option<NativeExecutionCounts>,
150    /// Whether the LLM call succeeded
151    pub success: bool,
152    /// Text response from the model
153    pub text: String,
154    /// Tool calls requested by the model
155    #[serde(default)]
156    pub tool_calls: Vec<ToolCall>,
157    /// Whether tool execution is needed
158    pub has_tool_calls: bool,
159    /// Tool definitions from applied capabilities (for tool execution)
160    #[serde(default)]
161    pub tool_definitions: Vec<ToolDefinition>,
162    /// Maximum iterations configured for the agent
163    #[serde(default = "default_max_iterations")]
164    pub max_iterations: usize,
165    /// Error message if the call failed
166    #[serde(skip_serializing_if = "Option::is_none")]
167    pub error: Option<String>,
168    /// Disclosed user-facing decision of the failure, already filtered
169    /// through the resolved error-disclosure mode. Hosts must prefer this over
170    /// re-classifying `error`/`text` strings so disclosure stays consistent.
171    #[serde(default, skip_serializing_if = "Option::is_none")]
172    pub user_facing_error: Option<UserFacingError>,
173    /// Error-disclosure mode that was applied to `user_facing_error`.
174    #[serde(default, skip_serializing_if = "Option::is_none")]
175    pub error_disclosure: Option<ErrorDisclosure>,
176    /// Token usage from the LLM call
177    #[serde(skip_serializing_if = "Option::is_none")]
178    pub usage: Option<TokenUsage>,
179    /// Assistant message emitted by `output.message.completed` for this generation.
180    #[serde(skip_serializing_if = "Option::is_none")]
181    pub output_message_id: Option<MessageId>,
182    /// Streaming latency for this LLM call, when available.
183    #[serde(skip_serializing_if = "Option::is_none")]
184    pub time_to_first_token_ms: Option<u64>,
185    /// LLM provider's response ID for chaining with `previous_response_id`
186    #[serde(skip_serializing_if = "Option::is_none")]
187    pub response_id: Option<String>,
188    /// Raw provider finish reason for this generation.
189    #[serde(default, skip_serializing_if = "Option::is_none")]
190    pub finish_reason: Option<String>,
191    /// Resolved locale used for this turn's prompt and backend-authored strings.
192    #[serde(skip_serializing_if = "Option::is_none")]
193    pub locale: Option<String>,
194    /// Merged network access list for URL filtering in tools.
195    #[serde(default, skip_serializing_if = "Option::is_none")]
196    pub network_access: Option<crate::network_access::NetworkAccessList>,
197    /// Request-level parallel tool calling preference (EVE-598), carried from
198    /// the resolved agent config into `ActInput` so the act scheduler can honor
199    /// `Some(false)` (force serialize). `None` preserves the default schedule.
200    #[serde(default, skip_serializing_if = "Option::is_none")]
201    pub parallel_tool_calls: Option<bool>,
202}
203
204fn default_max_iterations() -> usize {
205    500
206}
207
208// ============================================================================
209// ReasonAtom
210// ============================================================================
211
212/// Atom that calls the LLM model for reasoning
213///
214/// This atom:
215/// 1. Emits reason.started event
216/// 2. Retrieves agent and session configuration from stores
217/// 3. Resolves model using priority: controls.model_id > session.model_id > agent.default_model_id
218/// 4. Builds configuration with capabilities applied
219/// 5. Loads messages from the store
220/// 6. Patches dangling tool calls
221/// 7. Resolves image_file content parts to actual image data (if ImageResolver provided)
222/// 8. Calls the LLM with the messages
223/// 9. Stores the assistant response
224/// 10. Emits reason.completed event
225/// 11. Returns the result with tool calls (if any)
226pub struct ReasonAtom {
227    native_async: Option<Arc<tokio::sync::Mutex<crate::native_async::NativeAsyncCoordinator>>>,
228    context_resolver: Arc<dyn TurnContextResolver>,
229    message_retriever: Arc<dyn MessageRetriever>,
230    capability_registry: CapabilityRegistry,
231    event_emitter: PhaseEffectEmitter<dyn PhaseEffectSink>,
232    /// Optional image resolver for resolving image_file content parts
233    image_resolver: Option<Arc<dyn ImageResolver>>,
234    file_resolver: Option<Arc<dyn FileResolver>>,
235    /// Optional heartbeater for stream-liveness signalling (EVE-531).
236    stream_heartbeater: Option<Arc<dyn crate::durability::StreamHeartbeater>>,
237    /// Optional provider stall timeout (EVE-531). Default: 120s.
238    provider_stall_timeout: Option<std::time::Duration>,
239    /// Shared attempt/backoff/time budget for automatic provider recovery.
240    provider_retry_config: LlmRetryConfig,
241    /// Optional durable tool result store for transcript repair (EVE-533).
242    durable_tool_result_store: Option<Arc<dyn DurableToolResultStore>>,
243    /// Optional partial-stream store for ContinuePartial recovery (EVE-532).
244    partial_stream_store: Option<Arc<dyn PartialStreamStore>>,
245    /// Optional live reasoning-effort handle (EVE-595). When set and holding a
246    /// value, it overrides the message-derived effort on every LLM step, so a
247    /// tool can change effort mid-turn and have subsequent steps observe it.
248    reasoning_effort_handle: Option<crate::tool_context::ReasoningEffortHandle>,
249    /// Optional utility LLM service (EVE-573). Powers model-backed
250    /// end-of-message output guardrails (e.g. moderation). When absent, those
251    /// guardrails fail open and the seam is a no-op.
252    utility_llm_service: Option<Arc<dyn crate::UtilityLlmService>>,
253    /// Optional decisions. Powers guardrail checks that ask for a
254    /// calibrated number rather than text to parse. When absent, those checks
255    /// fail open exactly like the utility-model-backed ones.
256    decisions: Option<Arc<dyn crate::DecisionsService>>,
257    /// Optional session schedule store. Used by the `usage_limit_auto_continue`
258    /// capability to schedule a one-shot continuation after a provider usage
259    /// limit resets. When absent, the capability degrades to a no-op (no
260    /// continuation is scheduled and the error copy makes no auto-resume
261    /// promise).
262    schedule_store: Option<Arc<dyn crate::session_services::SessionScheduleStore>>,
263    /// Optional durable store for replacement context checkpoints.
264    compaction_checkpoint_store: Option<Arc<dyn crate::CompactionCheckpointStore>>,
265}
266
267impl ReasonAtom {
268    pub fn with_native_async(
269        mut self,
270        coordinator: Arc<tokio::sync::Mutex<crate::native_async::NativeAsyncCoordinator>>,
271    ) -> Self {
272        self.native_async = Some(coordinator);
273        self
274    }
275
276    /// Create a new ReasonAtom
277    pub fn new(
278        context_resolver: impl TurnContextResolver + 'static,
279        message_retriever: impl MessageRetriever + 'static,
280        capability_registry: CapabilityRegistry,
281        event_emitter: impl PhaseEffectSink + 'static,
282    ) -> Self {
283        Self {
284            native_async: None,
285            context_resolver: Arc::new(context_resolver),
286            message_retriever: Arc::new(message_retriever),
287            capability_registry,
288            event_emitter: PhaseEffectEmitter::new(Arc::new(event_emitter)),
289            image_resolver: None,
290            file_resolver: None,
291            stream_heartbeater: None,
292            provider_stall_timeout: None,
293            provider_retry_config: LlmRetryConfig::default(),
294            durable_tool_result_store: None,
295            partial_stream_store: None,
296            reasoning_effort_handle: None,
297            utility_llm_service: None,
298            decisions: None,
299            schedule_store: None,
300            compaction_checkpoint_store: None,
301        }
302    }
303
304    /// Set the session schedule store used by `usage_limit_auto_continue` to
305    /// schedule a continuation after a provider usage limit resets.
306    pub fn with_schedule_store(
307        mut self,
308        store: Arc<dyn crate::session_services::SessionScheduleStore>,
309    ) -> Self {
310        self.schedule_store = Some(store);
311        self
312    }
313
314    pub fn with_compaction_checkpoint_store(
315        mut self,
316        store: Arc<dyn crate::CompactionCheckpointStore>,
317    ) -> Self {
318        self.compaction_checkpoint_store = Some(store);
319        self
320    }
321
322    /// Collect the [`LlmErrorHook`]s contributed by the active capabilities,
323    /// paired with each capability's per-agent config. Hooks are invoked
324    /// generically on the terminal-error path; the reason atom has no knowledge
325    /// of any specific capability's behavior. Capabilities that contribute no
326    /// hook — the common case — are skipped at zero allocation cost.
327    fn collect_llm_error_hooks(
328        &self,
329        resolved_capability_configs: &[crate::CapabilityRef],
330    ) -> Vec<(
331        Arc<dyn crate::llm_error_hook::LlmErrorHook>,
332        serde_json::Value,
333    )> {
334        resolved_capability_configs
335            .iter()
336            .filter_map(|cfg| {
337                let cap = self.capability_registry.get(cfg.capability_id())?;
338                let hook = cap.llm_error_hook()?;
339                Some((hook, cfg.config_value().clone()))
340            })
341            .collect()
342    }
343
344    /// Set the image resolver for resolving image_file content parts
345    ///
346    /// When set, image_file references in messages will be resolved to actual
347    /// image data before being sent to the LLM. This is required for multimodal
348    /// conversations that include image attachments.
349    ///
350    /// # Example
351    ///
352    /// ```ignore
353    /// let resolver = Arc::new(GrpcImageResolver::new(client));
354    /// let atom = ReasonAtom::new(/* ... */).with_image_resolver(resolver);
355    /// ```
356    pub fn with_image_resolver(mut self, resolver: Arc<dyn ImageResolver>) -> Self {
357        self.image_resolver = Some(resolver);
358        self
359    }
360
361    pub fn with_file_resolver(mut self, resolver: Arc<dyn FileResolver>) -> Self {
362        self.file_resolver = Some(resolver);
363        self
364    }
365
366    /// Set the stream heartbeater for liveness signalling during LLM streaming.
367    pub fn with_stream_heartbeater(
368        mut self,
369        heartbeater: Arc<dyn crate::durability::StreamHeartbeater>,
370    ) -> Self {
371        self.stream_heartbeater = Some(heartbeater);
372        self
373    }
374
375    /// Set the provider stall timeout. If no token arrives within this window,
376    /// the stream is aborted and the activity fails with a retryable error.
377    pub fn with_provider_stall_timeout(mut self, timeout: std::time::Duration) -> Self {
378        self.provider_stall_timeout = Some(timeout);
379        self
380    }
381
382    /// Set the bounded provider recovery policy.
383    pub fn with_provider_retry_config(mut self, config: LlmRetryConfig) -> Self {
384        self.provider_retry_config = config;
385        self
386    }
387
388    /// Set the durable tool result store for transcript repair (EVE-533).
389    ///
390    /// When provided, transcript repair consults this store to replay settled tool
391    /// results or synthesize appropriate interrupted placeholders rather than always
392    /// emitting a generic "cancelled" message.
393    pub fn with_durable_tool_result_store(
394        mut self,
395        store: Arc<dyn DurableToolResultStore>,
396    ) -> Self {
397        self.durable_tool_result_store = Some(store);
398        self
399    }
400
401    /// Set the partial-stream store for ContinuePartial recovery (EVE-532).
402    pub fn with_partial_stream_store(mut self, store: Arc<dyn PartialStreamStore>) -> Self {
403        self.partial_stream_store = Some(store);
404        self
405    }
406
407    /// Set the live reasoning-effort handle (EVE-595).
408    ///
409    /// When set and holding a value, the effort it carries overrides the
410    /// message-derived effort for every LLM step. Because the handle is shared
411    /// and re-read on each step, a tool that mutates it mid-turn causes
412    /// subsequent steps in the same turn to use the new effort.
413    pub fn with_reasoning_effort_handle(
414        mut self,
415        handle: crate::tool_context::ReasoningEffortHandle,
416    ) -> Self {
417        self.reasoning_effort_handle = Some(handle);
418        self
419    }
420
421    /// Set the utility LLM service used by model-backed end-of-message output
422    /// guardrails (EVE-573). When unset, those guardrails fail open.
423    pub fn with_utility_llm_service(mut self, service: Arc<dyn crate::UtilityLlmService>) -> Self {
424        self.utility_llm_service = Some(service);
425        self
426    }
427
428    /// Set the decisions used by guardrail checks that ask for typed
429    /// answers. When unset, those checks fail open.
430    pub fn with_decisions(mut self, service: Arc<dyn crate::DecisionsService>) -> Self {
431        self.decisions = Some(service);
432        self
433    }
434}
435
436impl ReasonAtom {
437    /// Stable phase name used by logs and durable activity adapters.
438    pub fn name(&self) -> &'static str {
439        "reason"
440    }
441
442    /// Execute a reason phase using host-injected portable contracts.
443    pub async fn execute(&self, input: ReasonInput) -> Result<ReasonResult> {
444        self.execute_inner(input, None).await
445    }
446}
447
448impl ReasonAtom {
449    /// Execute using a pre-assembled turn context.
450    ///
451    /// Hosts that already assembled turn context for the current reason phase can
452    /// pass it through here to avoid reloading messages and rebuilding the agent.
453    pub async fn execute_with_assembled_context(
454        &self,
455        input: ReasonInput,
456        assembled: AssembledTurnContext,
457    ) -> Result<ReasonResult> {
458        self.execute_inner(input, Some(assembled)).await
459    }
460
461    async fn execute_inner(
462        &self,
463        input: ReasonInput,
464        assembled: Option<AssembledTurnContext>,
465    ) -> Result<ReasonResult> {
466        let ReasonInput {
467            context,
468            harness_id,
469            agent_id,
470            org_id,
471            mcp_tool_definitions,
472            previous_response_id,
473            iteration,
474        } = input;
475
476        tracing::info!(
477            session_id = %context.session_id,
478            turn_id = %context.turn_id,
479            exec_id = %context.exec_id,
480            harness_id = %harness_id,
481            agent_id = ?agent_id,
482            mcp_tools_count = %mcp_tool_definitions.len(),
483            "ReasonAtom: starting LLM call"
484        );
485
486        // Generate OTel-style span IDs for hierarchical tracing
487        // trace_id: groups all events in this turn
488        // span_id: unique identifier for this reason span (shared by started/completed)
489        // parent_span_id: links to turn as parent
490        //
491        // NOTE: TurnId::to_string() returns prefixed format (e.g., "turn_abc123")
492        // matching the format used by turn.started/completed events in Braintrust.
493        let trace_id = context.turn_id.to_string();
494        let reason_span_id = Uuid::now_v7().to_string();
495        let parent_span_id = trace_id.clone(); // Parent is the turn
496
497        // Create event context from atom context with span info
498        let event_context = EventContext::from_execution_context(&context).with_span(
499            trace_id.clone(),
500            reason_span_id.clone(),
501            Some(parent_span_id.clone()),
502        );
503
504        // Track reason phase timing for Braintrust observability
505        let reason_start = Instant::now();
506
507        // Emit reason.started event
508        if let Err(e) = self
509            .event_emitter
510            .emit(EventRequest::new(
511                context.session_id,
512                event_context.clone(),
513                ReasonStartedData {
514                    harness_id,
515                    agent_id,
516                    metadata: None, // Will be populated after model resolution
517                },
518            ))
519            .await
520        {
521            tracing::warn!(
522                session_id = %context.session_id,
523                error = %e,
524                "ReasonAtom: failed to emit reason.started event"
525            );
526        }
527
528        // Assemble the turn context up-front so the error path below knows
529        // the resolved provider/model and the error-disclosure mode even when
530        // the LLM call (or the assembly itself) fails.
531        let assembled = match assembled {
532            Some(assembled) => Ok(assembled),
533            None => {
534                self.context_resolver
535                    .resolve_turn_context(TurnContextRequest {
536                        session_id: context.session_id,
537                        harness_id,
538                        agent_id,
539                        mcp_tool_definitions: mcp_tool_definitions.clone(),
540                        allow_provider_managed_reduction: true,
541                    })
542                    .await
543            }
544        };
545
546        let (error_disclosure, error_context, error_hooks, call_result) = match assembled {
547            Ok(assembled) => {
548                let outcome = provider_managed_compaction::execute_with_fallback(
549                    provider_managed_compaction::Call {
550                        atom: self,
551                        session_id: context.session_id,
552                        harness_id,
553                        agent_id,
554                        org_id,
555                        context: &context,
556                        trace_id: &trace_id,
557                        reason_span_id: &reason_span_id,
558                        previous_response_id,
559                        iteration,
560                        mcp_tool_definitions: &mcp_tool_definitions,
561                        assembled,
562                    },
563                )
564                .await;
565                (
566                    outcome.disclosure,
567                    outcome.error_context,
568                    outcome.error_hooks,
569                    outcome.result,
570                )
571            }
572            Err(error) => (
573                ErrorDisclosure::default(),
574                UserFacingErrorContext::default(),
575                Vec::new(),
576                Err(error),
577            ),
578        };
579
580        // Handle LLM call errors gracefully
581        let result = match call_result {
582            Ok(result) => {
583                // Calculate reason phase duration
584                let reason_duration_ms = reason_start.elapsed().as_millis() as u64;
585
586                // Emit reason.completed event (same span as reason.started, parent is turn)
587                let completed_context = EventContext::from_execution_context(&context).with_span(
588                    trace_id.clone(),
589                    reason_span_id.clone(), // Same span_id as started
590                    Some(parent_span_id.clone()),
591                );
592                if let Err(e) = self
593                    .event_emitter
594                    .emit(EventRequest::new(
595                        context.session_id,
596                        completed_context,
597                        ReasonCompletedData::success(
598                            &result.text,
599                            result.has_tool_calls,
600                            result.tool_calls.len() as u32,
601                            Some(reason_duration_ms),
602                            result.usage.clone(),
603                        ),
604                    ))
605                    .await
606                {
607                    tracing::warn!(
608                        session_id = %context.session_id,
609                        error = %e,
610                        "ReasonAtom: failed to emit reason.completed event"
611                    );
612                }
613                result
614            }
615            Err(e) => {
616                // Calculate reason phase duration even for failures
617                let reason_duration_ms = reason_start.elapsed().as_millis() as u64;
618
619                // LLM call failure is a "normal" result per the spec
620                // Return a result indicating failure with the error message
621                tracing::warn!(
622                    session_id = %context.session_id,
623                    turn_id = %context.turn_id,
624                    error = %e,
625                    "ReasonAtom: LLM call failed"
626                );
627
628                let error_msg = e.to_string();
629                let mut source_error = e.user_facing_error(error_context);
630
631                // Only emit user-facing error events for non-transient errors.
632                // Transient errors (server errors, rate limits, timeouts) will be
633                // retried by the durable task engine. Emitting error events on each
634                // retry attempt causes duplicate error messages in the UI.
635                // The durable worker emits a single error event when all retries
636                // are exhausted (DLQ).
637                let is_transient = e.is_transient_llm_error()
638                    || (e.llm_error_kind().is_none() && is_transient_error_message(&error_msg));
639
640                // Capability error-hook seam: on the terminal (non-retried)
641                // error path, let active capabilities react — perform a side
642                // effect and/or augment the user-facing error fields — before the
643                // message is built. The atom stays behavior-agnostic; each hook
644                // (e.g. `usage_limit_auto_continue`) owns its own logic.
645                if !is_transient && !error_hooks.is_empty() {
646                    let services = crate::llm_error_hook::LlmErrorHookServices {
647                        schedule_store: self.schedule_store.clone(),
648                    };
649                    for (hook, config) in &error_hooks {
650                        let outcome = {
651                            let ctx = crate::llm_error_hook::LlmErrorContext {
652                                session_id: context.session_id,
653                                error_code: &source_error.code,
654                                error_fields: &source_error.fields,
655                                config,
656                                services: &services,
657                            };
658                            hook.on_llm_error(&ctx).await
659                        };
660                        for (key, value) in outcome.extra_error_fields {
661                            source_error = source_error.with_field(key, value);
662                        }
663                    }
664                }
665
666                let user_error = source_error.apply_disclosure(error_disclosure, Some(&error_msg));
667                let user_error_text = user_error.fallback_message();
668
669                let mut output_message_id = None;
670
671                if !is_transient {
672                    // Create error message for the user to see
673                    let mut error_message = RuntimeMessage::assistant(&user_error_text);
674                    let mut metadata = std::collections::HashMap::new();
675                    user_error.apply_to_message_metadata(&mut metadata);
676                    UserFacingError::apply_disclosure_to_message_metadata(
677                        &mut metadata,
678                        error_disclosure,
679                        &source_error.code,
680                    );
681                    error_message.metadata = Some(metadata);
682
683                    output_message_id = Some(error_message.id);
684
685                    // Emit output.message.completed event (stores message as event with proper turn context)
686                    // output.message.completed is child of reason span
687                    let error_msg_context = EventContext::from_execution_context(&context)
688                        .with_span(
689                            trace_id.clone(),
690                            Uuid::now_v7().to_string(),   // Own span_id
691                            Some(reason_span_id.clone()), // Parent is reason span
692                        );
693                    if let Err(emit_err) = self
694                        .event_emitter
695                        .emit(EventRequest::new(
696                            context.session_id,
697                            error_msg_context,
698                            OutputMessageCompletedData::new(error_message)
699                                .with_user_facing_error(&user_error)
700                                .with_error_disclosure(error_disclosure),
701                        ))
702                        .await
703                    {
704                        tracing::warn!(
705                            session_id = %context.session_id,
706                            error = %emit_err,
707                            "ReasonAtom: failed to emit output.message.completed event for error"
708                        );
709                    }
710                } else {
711                    tracing::info!(
712                        session_id = %context.session_id,
713                        "ReasonAtom: skipping error event for transient LLM error (will be retried)"
714                    );
715                }
716
717                // Emit reason.completed event for failure (same span as started, parent is turn)
718                let completed_context = EventContext::from_execution_context(&context).with_span(
719                    trace_id.clone(),
720                    reason_span_id.clone(), // Same span_id as started
721                    Some(parent_span_id.clone()),
722                );
723                if let Err(emit_err) = self
724                    .event_emitter
725                    .emit(EventRequest::new(
726                        context.session_id,
727                        completed_context,
728                        ReasonCompletedData::failure(error_msg.clone(), Some(reason_duration_ms)),
729                    ))
730                    .await
731                {
732                    tracing::warn!(
733                        session_id = %context.session_id,
734                        error = %emit_err,
735                        "ReasonAtom: failed to emit reason.completed event"
736                    );
737                }
738
739                ReasonResult {
740                    native_counts: None,
741                    success: false,
742                    text: user_error_text,
743                    tool_calls: vec![],
744                    has_tool_calls: false,
745                    tool_definitions: vec![],
746                    max_iterations: default_max_iterations(),
747                    error: Some(error_msg.clone()),
748                    user_facing_error: Some(user_error),
749                    error_disclosure: Some(error_disclosure),
750                    usage: None,
751                    output_message_id,
752                    time_to_first_token_ms: None,
753                    response_id: None,
754                    finish_reason: error_msg
755                        .to_ascii_lowercase()
756                        .contains("model refused")
757                        .then(|| "refusal".to_string()),
758                    locale: None,
759                    network_access: None,
760                    parallel_tool_calls: None,
761                }
762            }
763        };
764
765        Ok(result)
766    }
767
768    /// Execute the actual LLM call
769    #[allow(clippy::too_many_arguments)]
770    async fn execute_llm_call(
771        &self,
772        session_id: SessionId,
773        harness_id: HarnessId,
774        agent_id: Option<AgentId>,
775        org_id: i64,
776        context: &ExecutionContext,
777        trace_id: &str,
778        reason_span_id: &str,
779        previous_response_id: Option<String>,
780        iteration: u32,
781        assembled: AssembledTurnContext,
782    ) -> Result<ReasonResult> {
783        let prior_usage = assembled.cumulative_usage();
784        let mut messages = transcript::order_native_results(assembled.messages);
785        let mut message_source_sequence = assembled.message_source_sequence;
786        let model_with_provider = assembled.model;
787        let provider_managed = model_with_provider
788            .provider_managed_reduction_option
789            .is_some();
790        let supports_clear_at = facts::supports_clear_at(
791            &model_with_provider.provider_type,
792            &model_with_provider.model,
793        );
794        let resolved_model_id = assembled.resolved_model_id;
795        let resolved_locale = assembled.resolved_locale;
796        let compaction_policy = assembled.compaction_policy;
797        let resolved_capability_configs = assembled.resolved_capability_configs;
798        let runtime_agent = assembled.runtime_agent;
799        let embedder_metadata = assembled.embedder_metadata;
800
801        emit_capability_usage_snapshot(
802            self.event_emitter.as_ref(),
803            &self.capability_registry,
804            session_id,
805            context,
806            &resolved_capability_configs,
807            &runtime_agent.tools,
808        )
809        .await;
810
811        let output_hooks =
812            collect_output_hooks(&self.capability_registry, &resolved_capability_configs);
813        let guardrail_providers = output_hooks.streaming;
814        let post_output_providers = output_hooks.post_generation;
815        let annotation_providers = output_hooks.annotations;
816        let citation_verifiers = output_hooks.citation_verifiers;
817
818        // 7. Create LLM driver using factory
819        let chat_driver = Arc::clone(&model_with_provider.driver);
820        let stateful_response_continuation =
821            previous_response_id.is_some() && chat_driver.supports_stateful_responses();
822        let mut restored_checkpoint: Option<crate::CompactionCheckpoint> = None;
823        let mut checkpoint_suffix_message_count = 0usize;
824        let native_reasoning_compaction = compaction_policy.as_ref().is_none_or(|policy| {
825            matches!(
826                policy.settings().strategy,
827                crate::compaction_policy::CompactionStrategy::Native
828                    | crate::compaction_policy::CompactionStrategy::Auto
829            ) && chat_driver.supports_compact()
830        });
831
832        let checkpoint_format = super::provider_checkpoint::format_version(provider_managed);
833        if (compaction_policy.is_some() || provider_managed)
834            && let Some(store) = self.compaction_checkpoint_store.as_ref()
835            && let Some(checkpoint) = store
836                .get_latest_format(
837                    session_id,
838                    model_with_provider.provider_type.as_str(),
839                    &model_with_provider.model,
840                    checkpoint_format,
841                )
842                .await?
843            && super::provider_checkpoint::is_restorable(
844                &checkpoint,
845                chat_driver.as_ref(),
846                model_with_provider.provider_type.as_str(),
847                &model_with_provider.model,
848                provider_managed,
849                native_reasoning_compaction,
850            )
851        {
852            let filters = crate::capabilities::collect_message_filters_only_with_context(
853                &resolved_capability_configs,
854                &self.capability_registry,
855                provider_managed,
856            );
857            let mut query =
858                crate::MessageQuery::new(session_id).after_sequence(checkpoint.source_sequence);
859            filters.apply_message_filters(&mut query);
860            let history = self.message_retriever.load_filtered_history(query).await?;
861            messages = history.messages;
862            checkpoint_suffix_message_count = messages.len();
863            filters.apply_post_load_filters(&mut messages);
864            if let crate::CompactionCheckpointPayload::Summary { text } = &checkpoint.payload {
865                messages.insert(
866                    0,
867                    RuntimeMessage::system(format!(
868                        "[CONVERSATION_SUMMARY]\n{text}\n[/CONVERSATION_SUMMARY]"
869                    )),
870                );
871            }
872            message_source_sequence = history.source_sequence.or(message_source_sequence);
873            restored_checkpoint = Some(checkpoint);
874        }
875
876        let controls = resolve_request_controls(
877            &messages,
878            self.reasoning_effort_handle.as_ref(),
879            &model_with_provider.provider_type,
880            &model_with_provider.model,
881        );
882        let reasoning_effort = controls.reasoning_effort;
883        let speed = controls.speed;
884        let verbosity = controls.verbosity;
885        let checkpoint_reasoning =
886            super::provider_checkpoint::reasoning_state(restored_checkpoint.as_ref());
887        let mut reasoning_replay = reasoning_updates::prepare(
888            &messages,
889            model_with_provider.provider_type.as_str(),
890            &model_with_provider.model,
891            reasoning_effort,
892            self.reasoning_effort_handle
893                .as_ref()
894                .and_then(crate::tool_context::ReasoningEffortHandle::get),
895            checkpoint_reasoning,
896        )
897        .filter(|_| native_reasoning_compaction);
898
899        // 9. Check for an in-flight partial assistant stream from a previous worker (EVE-532).
900        // If found, apply the ContinuePartial recovery policy: finalize from accumulated
901        // text (if non-empty) or restart clean (if empty/usable partial only).
902        if let Some(ref store) = self.partial_stream_store {
903            let turn_id_str = context.turn_id.to_string();
904            match store.get_partial_stream(session_id, &turn_id_str).await {
905                Ok(Some(partial)) if !partial.accumulated.is_empty() => {
906                    // Finalize: emit completed from persisted accumulated text.
907                    return self
908                        .finalize_partial_stream(
909                            session_id,
910                            context,
911                            partial,
912                            iteration,
913                            &runtime_agent,
914                            &resolved_capability_configs,
915                        )
916                        .await;
917                }
918                Ok(Some(partial)) => {
919                    if let (Some(replay), Some(mut saved)) =
920                        (reasoning_replay.as_mut(), partial.reasoning_state)
921                    {
922                        // The old worker persisted the effective live override
923                        // before sending. Its process-local handle is gone.
924                        saved.pending = saved.effective;
925                        replay.state = saved;
926                    }
927                    // Empty accumulated: restart clean — fall through to normal LLM call.
928                    // Emit reason.recovered { mode: Restart } for observability.
929                    let recovery_ctx = EventContext::from_execution_context(context);
930                    let _ = self
931                        .event_emitter
932                        .emit(EventRequest::new(
933                            session_id,
934                            recovery_ctx,
935                            ReasonRecoveredData {
936                                turn_id: context.turn_id,
937                                mode: RecoveryMode::Restart,
938                                accumulated_len: 0,
939                            },
940                        ))
941                        .await;
942                    tracing::info!(
943                        session_id = %session_id,
944                        turn_id = %context.turn_id,
945                        "ReasonAtom: partial stream detected with empty accumulated; restarting clean"
946                    );
947                }
948                Ok(None) => {} // No partial; normal first-run execution.
949                Err(e) => {
950                    if reasoning_replay.is_some() {
951                        return Err(e);
952                    }
953                    // Best-effort: log and continue with normal execution.
954                    tracing::warn!(
955                        session_id = %session_id,
956                        turn_id = %context.turn_id,
957                        error = %e,
958                        "ReasonAtom: partial-stream store error; proceeding with normal execution"
959                    );
960                }
961            }
962        }
963
964        // 10. Repair dangling tool calls (EVE-533): ensure every assistant tool_call
965        // has a matching ToolResult before the LLM call. Consults durable_tool_results
966        // when available to replay settled results or synthesize interrupted placeholders.
967        let repair_event_context = EventContext::from_execution_context(context);
968        let patched_messages = if self.native_async.is_some() {
969            messages.clone()
970        } else {
971            repair_dangling_tool_calls(
972                &messages,
973                self.durable_tool_result_store.as_deref(),
974                self.event_emitter.as_ref(),
975                session_id,
976                &repair_event_context,
977                &context.turn_id.to_string(),
978            )
979            .await
980        };
981        let raw_tool_result_bytes = compaction_policy
982            .as_ref()
983            .map(|policy| policy.total_tool_result_bytes(&patched_messages))
984            .unwrap_or(0);
985
986        // 9b. Let enabled capabilities build a prompt-facing model view from
987        // lossless stored messages. Storage remains unchanged.
988        let model_view_providers = crate::capabilities::collect_model_view_providers(
989            &resolved_capability_configs,
990            &self.capability_registry,
991            Some(model_with_provider.model.as_str()),
992        );
993        let model_view_context = crate::capabilities::ModelViewContext {
994            session_id,
995            prior_usage: prior_usage.as_ref(),
996            provider_managed_reduction: provider_managed,
997        };
998        let mut context_messages =
999            model_view_providers.apply_model_view(patched_messages, &model_view_context);
1000        context_messages = crate::tool_call_integrity::retain_complete_message_tool_exchanges(
1001            &context_messages,
1002            stateful_response_continuation || restored_checkpoint.is_some(),
1003        );
1004
1005        // 9c. Dynamic facts after each answered input, rendered as of that input.
1006        // Capable Anthropic models use turn-scoped system messages; others keep
1007        // the user-message fallback (see `facts`).
1008        let render_facts = |at| {
1009            crate::capabilities::render_facts_block(&crate::capabilities::collect_dynamic_facts(
1010                &resolved_capability_configs,
1011                &self.capability_registry,
1012                Some(model_with_provider.model.as_str()),
1013                &crate::capabilities::FactsContext::new(session_id).at(at),
1014            ))
1015        };
1016        let (context_messages, volatile_suffix_len) =
1017            facts::interleave_facts(context_messages, render_facts, supports_clear_at);
1018        let mut context_messages = context_messages;
1019
1020        // 9d. Prepend conversation context (e.g. hierarchical AGENTS.md) as
1021        // the leading user-role message.
1022        //
1023        // Project instructions ride here: model-visible on every turn and
1024        // re-resolved alongside the system prompt, but never folded into the
1025        // cached system prompt. Untrusted workspace content must stay below
1026        // harness safety instructions in the instruction hierarchy, and file
1027        // edits must not invalidate the cache-stable system prefix.
1028        if let Some(context) = runtime_agent.conversation_context.as_ref()
1029            && !context.is_empty()
1030        {
1031            context_messages.insert(0, RuntimeMessage::user(context.clone()));
1032        }
1033
1034        // 10. Resolve images from image_file references (if any)
1035        //
1036        // Image resolution converts image_file content parts (which only contain UUIDs)
1037        // into actual base64-encoded image data that can be sent to LLMs.
1038        let resolved_images = self.resolve_images(&context_messages).await;
1039        let resolved_files = self.resolve_files(&context_messages).await;
1040
1041        // 11. Build LLM messages
1042        let mut llm_messages = Vec::new();
1043
1044        // Add system prompt
1045        let has_system_prompt = !runtime_agent.system_prompt.is_empty();
1046        if has_system_prompt {
1047            llm_messages.push(Message {
1048                native_tool_calls: Vec::new(),
1049                role: MessageRole::System,
1050                content: MessageContent::Text(runtime_agent.system_prompt.clone()),
1051                tool_calls: None,
1052                tool_call_id: None,
1053                phase: None,
1054                reasoning: Vec::new(),
1055                configuration_update: None,
1056            });
1057        }
1058
1059        // Build messages for llm.generation event (includes system message)
1060        let messages_for_event: Vec<RuntimeMessage> = if has_system_prompt {
1061            std::iter::once(RuntimeMessage::system(&runtime_agent.system_prompt))
1062                .chain(context_messages.iter().cloned())
1063                .collect()
1064        } else {
1065            context_messages.clone()
1066        };
1067
1068        // Add conversation messages with resolved images.
1069        // For user messages with an external_actor, prefix the first text part
1070        // with the actor's display label so the LLM knows who is speaking.
1071        // Skip error placeholder messages from prior failed turns — they add
1072        // no conversational value and inflate the request.
1073        let mut stripped_error_count = 0u32;
1074        for msg in &context_messages {
1075            if is_error_placeholder_message(msg) {
1076                stripped_error_count += 1;
1077                continue;
1078            }
1079            let mut llm_msg = crate::llm_conversions::llm_message_from_message_with_attachments(
1080                msg,
1081                &resolved_images,
1082                &resolved_files,
1083            );
1084            llm_msg.configuration_update = reasoning_replay
1085                .as_ref()
1086                .and_then(|replay| replay.transitions.get(&msg.id).copied());
1087            if msg.role == RuntimeMessageRole::User
1088                && let Some(ref actor) = msg.external_actor
1089            {
1090                llm_msg.prepend_text_prefix(&format!("[{}] ", actor.display_label()));
1091            }
1092            facts::mark_turn_scoped(&mut llm_msg, msg, supports_clear_at);
1093            llm_messages.push(llm_msg);
1094        }
1095        if stripped_error_count > 0 {
1096            tracing::info!(
1097                session_id = %session_id,
1098                stripped_error_count,
1099                "ReasonAtom: stripped error placeholder messages from LLM input"
1100            );
1101        }
1102
1103        // Context reducers operate on prompt-facing copies and may select only
1104        // one side of a tool exchange at a window boundary. Stateless requests
1105        // must be self-contained; stateful Responses requests may retain
1106        // result-only deltas whose calls live behind `previous_response_id`.
1107        llm_messages = crate::tool_call_integrity::retain_complete_llm_tool_exchanges_for_request(
1108            llm_messages,
1109            stateful_response_continuation || restored_checkpoint.is_some(),
1110        );
1111
1112        // 12. Build LLM call config with reasoning effort and metadata
1113        let mut llm_config_builder =
1114            crate::llm_conversions::llm_call_config_builder_from_agent(&runtime_agent);
1115        if let Some(effort) = reasoning_effort {
1116            llm_config_builder = llm_config_builder.reasoning_effort(effort);
1117        }
1118        if let Some(speed) = speed {
1119            llm_config_builder = llm_config_builder.speed(speed);
1120        }
1121        if let Some(verbosity) = verbosity {
1122            llm_config_builder = llm_config_builder.verbosity(verbosity);
1123        }
1124
1125        // Inject embedder metadata first; system keys added below take precedence
1126        for (k, v) in &embedder_metadata {
1127            llm_config_builder = llm_config_builder.with_metadata(k, v.clone());
1128        }
1129
1130        // Add metadata for API tracking and debugging
1131        // These IDs help correlate API requests with Everruns entities
1132        // TypedId::to_string() produces prefixed format (e.g., "session_abc123")
1133        llm_config_builder = llm_config_builder
1134            .with_metadata("session_id", session_id.to_string())
1135            .with_metadata("harness_id", harness_id.to_string())
1136            .with_metadata("turn_id", context.turn_id.to_string())
1137            .with_metadata("exec_id", context.exec_id.to_string())
1138            .with_metadata("org_id", format!("org_{:032x}", org_id));
1139        if let Some(agent_id) = agent_id {
1140            llm_config_builder = llm_config_builder.with_metadata("agent_id", agent_id.to_string());
1141        }
1142
1143        // Add model_id if we have one (not available for system default model)
1144        if let Some(model_id) = &resolved_model_id {
1145            llm_config_builder = llm_config_builder.with_metadata("model_id", model_id.to_string());
1146        }
1147
1148        let mut llm_config = llm_config_builder
1149            .previous_response_id(previous_response_id.clone())
1150            .volatile_suffix_len(volatile_suffix_len)
1151            .build();
1152        if let Some(replay) = &reasoning_replay {
1153            llm_config.reasoning_effort = replay.state.baseline;
1154            llm_config.reasoning_state = Some(replay.state.clone());
1155            if replay.reset_continuation {
1156                llm_config.previous_response_id = None;
1157            }
1158        } else if messages
1159            .iter()
1160            .rev()
1161            .find(|message| {
1162                message.role == RuntimeMessageRole::Agent && !is_error_placeholder_message(message)
1163            })
1164            .and_then(|message| message.metadata.as_ref())
1165            .is_some_and(|metadata| metadata.contains_key(reasoning_updates::STATE_KEY))
1166        {
1167            // Leaving Astra's configuration-update mode starts a fresh provider
1168            // chain. Never inherit its updates in another model or protocol.
1169            llm_config.previous_response_id = None;
1170        }
1171        if let Some(checkpoint) = restored_checkpoint.as_ref()
1172            && let crate::CompactionCheckpointPayload::ProviderOpaque { context } =
1173                &checkpoint.payload
1174        {
1175            llm_config.previous_response_id = None;
1176            llm_config.provider_opaque_context = Some(context.clone());
1177        }
1178        provider_managed_compaction::decorate_request(
1179            &mut llm_config,
1180            provider_managed,
1181            restored_checkpoint.is_some(),
1182        );
1183
1184        tracing::debug!(
1185            session_id = %session_id,
1186            turn_id = %context.turn_id,
1187            model = %runtime_agent.model,
1188            message_count = %llm_messages.len(),
1189            "ReasonAtom: calling LLM"
1190        );
1191
1192        // 13. Emit output.message.started event BEFORE starting LLM call
1193        // This allows UI to show a thinking indicator immediately
1194        let streaming_event_context = EventContext::from_execution_context(context);
1195
1196        // Arm output guardrails for this stream. Each guardrail sees the
1197        // assembled system prompt and its own per-capability config (already
1198        // borrowed in `guardrail_providers` above, so no second scan over
1199        // `resolved_capability_configs`). Guardrails that decline to arm —
1200        // e.g. the canary couldn't extract a long-enough sentence — are
1201        // skipped, leaving the streaming hot path entirely free of work.
1202        let mut armed_guardrails: Vec<ArmedGuardrail> = Vec::new();
1203        for (cap_id, cfg, provider) in &guardrail_providers {
1204            let ctx = OutputGuardrailContext {
1205                system_prompt: &runtime_agent.system_prompt,
1206                config: cfg,
1207            };
1208            let guardrail_id = provider.id().to_string();
1209            if let Some(run) = provider.arm(&ctx) {
1210                armed_guardrails.push(ArmedGuardrail {
1211                    capability_id: cap_id.clone(),
1212                    guardrail_id,
1213                    run,
1214                });
1215            }
1216        }
1217        // Blocking post-generation guardrails need the full assistant message
1218        // before they can decide. When active, withhold text deltas until the
1219        // seam allows the finalized output so blocked tokens are never emitted
1220        // or persisted as output.message.delta events.
1221        let buffer_output_deltas = !post_output_providers.is_empty();
1222        // Allocate the public message id before the first lifecycle event so
1223        // started/delta/replaced/completed can be grouped without turn-level
1224        // heuristics. Each reasoning iteration reaches this point separately.
1225        let output_message_id = MessageId::new();
1226        tracing::info!(
1227            session_id = %session_id,
1228            turn_id = %context.turn_id,
1229            "ReasonAtom: emitting output.message.started event"
1230        );
1231        if let Err(e) = self
1232            .event_emitter
1233            .emit(EventRequest::new(
1234                session_id,
1235                streaming_event_context.clone(),
1236                OutputMessageStartedData {
1237                    reasoning_state: llm_config.reasoning_state.clone(),
1238                    turn_id: context.turn_id,
1239                    message_id: output_message_id,
1240                    model: Some(runtime_agent.model.clone()),
1241                    iteration: Some(iteration),
1242                    // Emitted before the LLM call — phase is not yet known, so the
1243                    // streamed hint starts `None` (treat as assistant text).
1244                    phase: None,
1245                },
1246            ))
1247            .await
1248        {
1249            if llm_config.reasoning_state.is_some() {
1250                return Err(e);
1251            }
1252            tracing::warn!(
1253                session_id = %session_id,
1254                error = %e,
1255                "ReasonAtom: failed to emit output.message.started event"
1256            );
1257        } else {
1258            tracing::info!(
1259                session_id = %session_id,
1260                "ReasonAtom: output.message.started event emitted successfully"
1261            );
1262        }
1263
1264        // Also emit reason.thinking.started if extended thinking is enabled
1265        let thinking_enabled = reasoning_effort.is_some();
1266        if thinking_enabled {
1267            tracing::info!(
1268                session_id = %session_id,
1269                turn_id = %context.turn_id,
1270                "ReasonAtom: emitting reason.thinking.started event"
1271            );
1272            if let Err(e) = self
1273                .event_emitter
1274                .emit(EventRequest::new(
1275                    session_id,
1276                    streaming_event_context.clone(),
1277                    ReasonThinkingStartedData {
1278                        turn_id: context.turn_id,
1279                        model: Some(runtime_agent.model.clone()),
1280                    },
1281                ))
1282                .await
1283            {
1284                tracing::warn!(
1285                    session_id = %session_id,
1286                    error = %e,
1287                    "ReasonAtom: failed to emit reason.thinking.started event"
1288                );
1289            } else {
1290                tracing::info!(
1291                    session_id = %session_id,
1292                    "ReasonAtom: reason.thinking.started event emitted successfully"
1293                );
1294            }
1295        }
1296
1297        // Track LLM call timing
1298        let llm_start = Instant::now();
1299
1300        // Try LLM call with automatic compaction on RequestTooLarge.
1301        // Transient errors (429, 5xx) are retried at the driver level.
1302        // Stream-level errors are not retried here to avoid duplicate user-visible messages.
1303        let mut compaction_info: Option<LlmCompactionInfo> = None;
1304        let mut llm_messages_for_call = llm_messages.clone();
1305        let compaction_lifecycle = provider_managed_compaction::Lifecycle::new(
1306            self.event_emitter.as_ref(),
1307            session_id,
1308            &streaming_event_context,
1309            &model_with_provider.model,
1310            model_with_provider.provider_type.as_str(),
1311            llm_messages_for_call.len(),
1312            message_source_sequence,
1313            &llm_config,
1314        );
1315
1316        if !provider_managed && let Some(policy) = compaction_policy.as_deref() {
1317            compaction_info = apply_proactive_compaction(
1318                ProactiveCompactionContext {
1319                    chat_driver: chat_driver.as_ref(),
1320                    policy,
1321                    checkpoint_store: self.compaction_checkpoint_store.as_ref(),
1322                    event_emitter: self.event_emitter.as_ref(),
1323                    event_context: &streaming_event_context,
1324                    session_id,
1325                    message_source_sequence,
1326                    provider_type: model_with_provider.provider_type.as_str(),
1327                    model: &model_with_provider.model,
1328                    system_prompt: has_system_prompt
1329                        .then_some(runtime_agent.system_prompt.as_str()),
1330                    stateful_response_continuation,
1331                    checkpoint_restored: restored_checkpoint.is_some(),
1332                    checkpoint_suffix_message_count,
1333                    raw_tool_result_bytes,
1334                    prior_usage: prior_usage.as_ref(),
1335                },
1336                &mut llm_messages_for_call,
1337                &mut llm_config,
1338            )
1339            .await?;
1340        }
1341
1342        // 14. Process stream with batched output.message.delta emissions
1343        // Batch deltas every 100ms to reduce event volume while providing real-time feedback
1344        const DELTA_BATCH_INTERVAL_MS: u64 = 100;
1345        let retry_config = self.provider_retry_config.clone();
1346        // OpenRouter server tools execute inside the provider and therefore do
1347        // not surface as agent ToolCalls. Reissuing their request can duplicate
1348        // side effects even when the stream has emitted only reasoning. The
1349        // routing payload shape is owned by the OpenRouter driver crate
1350        // (`everruns_openrouter::options`); the engine only sniffs the opaque
1351        // `driver_options` data so the `engine -> core/provider/capability`
1352        // dependency direction holds.
1353        let has_provider_executed_tools = llm_config
1354            .driver_options
1355            .get("openrouter/routing")
1356            .and_then(|raw| raw.get("server_tools"))
1357            .and_then(|tools| tools.as_array())
1358            .is_some_and(|tools| !tools.is_empty());
1359        let mut stream_retry_metadata = RetryMetadata::default();
1360        let mut retry_started_at = None;
1361        // Best-effort streamed phase hint (EVE-774). Starts `None` ("not yet
1362        // classified — treat as assistant text") and is refined monotonically
1363        // once a provider reveals a native phase mid-stream. Declared outside the
1364        // retry loop so it is available to the post-loop guarded delta emission.
1365        let mut streamed_phase: Option<everruns_provider::ExecutionPhase> = None;
1366        let mut native_calls = std::collections::BTreeMap::new();
1367        let mut compaction_started_at: Option<Instant> = None;
1368        let (
1369            text,
1370            thinking,
1371            reasoning,
1372            tool_calls,
1373            completion_metadata,
1374            time_to_first_token_ms,
1375            pending_delta,
1376            mut tripped,
1377        ) = 'stream_attempt: loop {
1378            let stream_result = if let Some(remaining) =
1379                remaining_retry_time(&retry_config, retry_started_at)
1380            {
1381                match tokio::time::timeout(
1382                    remaining,
1383                    chat_driver.chat_completion_stream(
1384                        &crate::ProviderEndpoint::default(),
1385                        llm_messages_for_call.clone(),
1386                        &llm_config,
1387                    ),
1388                )
1389                .await
1390                {
1391                    Ok(result) => result,
1392                    Err(_) => {
1393                        compaction_lifecycle.fail_if(provider_managed).await;
1394                        return Err(AgentLoopError::llm_kind(
1395                            crate::error::LlmErrorKind::Unavailable,
1396                            format!(
1397                                "provider retry time budget exhausted after {} retries over {:.1}s; the turn is safe to resume",
1398                                stream_retry_metadata.attempts,
1399                                retry_config.max_retry_elapsed.as_secs_f64()
1400                            ),
1401                        )
1402                        .with_retry_metadata(&stream_retry_metadata));
1403                    }
1404                }
1405            } else {
1406                chat_driver
1407                    .chat_completion_stream(
1408                        &crate::ProviderEndpoint::default(),
1409                        llm_messages_for_call.clone(),
1410                        &llm_config,
1411                    )
1412                    .await
1413            };
1414            let mut stream = match stream_result {
1415                Ok(stream) => stream,
1416                Err(e) if e.is_request_too_large() => {
1417                    compaction_lifecycle.fail_if(provider_managed).await;
1418                    if provider_managed {
1419                        return Err(e);
1420                    }
1421                    let Some(policy) = compaction_policy.as_deref() else {
1422                        tracing::warn!(
1423                            session_id = %session_id,
1424                            turn_id = %context.turn_id,
1425                            "ReasonAtom: context too large and compaction capability is not enabled"
1426                        );
1427                        return Err(e);
1428                    };
1429                    let outcome = apply_reactive_compaction(
1430                        ReactiveCompactionContext {
1431                            chat_driver: chat_driver.as_ref(),
1432                            policy,
1433                            checkpoint_store: self.compaction_checkpoint_store.as_ref(),
1434                            event_emitter: self.event_emitter.as_ref(),
1435                            event_context: &streaming_event_context,
1436                            session_id,
1437                            message_source_sequence,
1438                            provider_type: model_with_provider.provider_type.as_str(),
1439                            model: &model_with_provider.model,
1440                            summarization_model_fallback: &runtime_agent.model,
1441                            system_prompt: has_system_prompt
1442                                .then_some(runtime_agent.system_prompt.as_str()),
1443                            stateful_response_continuation,
1444                        },
1445                        &mut llm_messages_for_call,
1446                        &mut llm_config,
1447                    )
1448                    .await?;
1449                    let Some(outcome) = outcome else {
1450                        return Err(e);
1451                    };
1452                    if outcome.generation_info.is_some() {
1453                        compaction_info = outcome.generation_info;
1454                    }
1455
1456                    chat_driver
1457                        .chat_completion_stream(
1458                            &crate::ProviderEndpoint::default(),
1459                            llm_messages_for_call.clone(),
1460                            &llm_config,
1461                        )
1462                        .await?
1463                }
1464                Err(e)
1465                    if e.is_transient_llm_error()
1466                        && !e.llm_retry_handled()
1467                        && !has_provider_executed_tools
1468                        && stream_retry_metadata.attempts < retry_config.max_retries =>
1469                {
1470                    let proposed_wait =
1471                        retry_config.calculate_backoff(stream_retry_metadata.attempts);
1472                    let Some(wait_duration) =
1473                        reserve_retry_wait(&retry_config, &mut retry_started_at, proposed_wait)
1474                    else {
1475                        compaction_lifecycle.fail_if(provider_managed).await;
1476                        return Err(AgentLoopError::llm_kind(
1477                            e.llm_error_kind()
1478                                .unwrap_or(crate::error::LlmErrorKind::Unavailable),
1479                            format!(
1480                                "{e}; automatic recovery time budget exhausted after {} retries; the turn is safe to resume",
1481                                stream_retry_metadata.attempts
1482                            ),
1483                        )
1484                        .with_retry_metadata(&stream_retry_metadata));
1485                    };
1486                    tracing::warn!(
1487                        session_id = %session_id,
1488                        turn_id = %context.turn_id,
1489                        attempt = stream_retry_metadata.attempts + 1,
1490                        max_retries = retry_config.max_retries,
1491                        wait_secs = wait_duration.as_secs_f64(),
1492                        error = %e,
1493                        "ReasonAtom: transient provider failure before stream, retrying"
1494                    );
1495                    stream_retry_metadata.record_retry(wait_duration, None);
1496                    tokio::time::sleep(wait_duration).await;
1497                    continue 'stream_attempt;
1498                }
1499                Err(error) => {
1500                    return compaction_lifecycle
1501                        .fail_start(error, provider_managed)
1502                        .await;
1503                }
1504            };
1505
1506            if let Some(coordinator) = &self.native_async {
1507                coordinator
1508                    .lock()
1509                    .await
1510                    .begin_transcript_response(output_message_id.to_string())
1511                    .await?;
1512                let coordinator = coordinator.clone();
1513                stream = Box::pin(futures::stream::unfold(
1514                    Some((coordinator, stream)),
1515                    |state| async move {
1516                        let (coordinator, mut source) = state?;
1517                        let event = coordinator
1518                            .lock()
1519                            .await
1520                            .next_response_event(&mut source)
1521                            .await;
1522                        let finished = matches!(&event, Ok(LlmStreamEvent::Done(_)) | Err(_));
1523                        Some((event, (!finished).then_some((coordinator, source))))
1524                    },
1525                ));
1526            }
1527            let mut text = String::new();
1528            // Reasoning artifacts in emission order. One entry per provider
1529            // block, each keeping its own signature/id, so interleaved thinking
1530            // and per-call thought signatures survive replay.
1531            let mut reasoning: Vec<ReasoningContentPart> = Vec::new();
1532            // Live-render buffer only; the durable text lives on the artifacts.
1533            let mut thinking = String::new();
1534            let mut tool_calls = Vec::new();
1535            let mut termination = StreamTermination::Exhausted;
1536            let mut replay_state = StreamReplayState::for_request(has_provider_executed_tools);
1537            let mut pending_delta = String::new();
1538            let mut pending_thinking_delta = String::new();
1539            let mut last_delta_emit = Instant::now();
1540            let mut last_thinking_delta_emit = Instant::now();
1541            let mut time_to_first_token_ms: Option<u64> = None;
1542
1543            // EVE-531: stall timeout + keepalive heartbeat for stream-liveness
1544            let stall_timeout = self
1545                .provider_stall_timeout
1546                .unwrap_or(std::time::Duration::from_secs(120));
1547            let initial_stall_timeout = remaining_retry_time(&retry_config, retry_started_at)
1548                .map_or(stall_timeout, |remaining| remaining.min(stall_timeout));
1549            let mut stall_sleep = Box::pin(tokio::time::sleep(initial_stall_timeout));
1550            let mut keepalive_ticker = tokio::time::interval(std::time::Duration::from_secs(12));
1551            keepalive_ticker.set_missed_tick_behavior(tokio::time::MissedTickBehavior::Skip);
1552            keepalive_ticker.tick().await; // consume immediate first tick
1553            let mut last_stream_heartbeat = Instant::now();
1554            // Tracks the wall-clock time of the last actual token received.
1555            // Updated only on content events; keepalive heartbeats use this
1556            // so the control plane can distinguish "alive/slow" from "making
1557            // progress" without conflating keepalive pings with real tokens.
1558            let mut last_token_at_unix: u64 = unix_now_secs();
1559
1560            loop {
1561                let event = tokio::select! {
1562                    biased;
1563                    next = stream.next() => match next {
1564                        Some(e) => e,
1565                        None => break,
1566                    },
1567                    _ = &mut stall_sleep => {
1568                        // EVE-806: a stream that produced no tokens within the
1569                        // liveness window is equivalent to a dropped connection.
1570                        // Route it through the same bounded transient-retry path
1571                        // as an in-stream provider error (everruns-provider
1572                        // classifies this message as transient) instead of
1573                        // failing the turn immediately. Retrying re-issues the
1574                        // same request with no artificial history messages; a
1575                        // stall after partial output is not retried, and repeated
1576                        // stalls stay bounded by retry_config.max_retries.
1577                        let stall_error =
1578                            crate::driver_registry::LlmStreamError::new(format!(
1579                                "provider stream stall: no tokens for {}s",
1580                                stall_timeout.as_secs()
1581                            ));
1582                        tracing::warn!(
1583                            session_id = %session_id,
1584                            turn_id = %context.turn_id,
1585                            stall_secs = stall_timeout.as_secs(),
1586                            "ReasonAtom: provider stream stall timeout"
1587                        );
1588                        if replay_state.should_retry(
1589                            &stall_error,
1590                            stream_retry_metadata.attempts,
1591                            retry_config.max_retries,
1592                        ) {
1593                            let proposed_wait = retry_config
1594                                .calculate_backoff(stream_retry_metadata.attempts);
1595                            let Some(wait_duration) = reserve_retry_wait(
1596                                &retry_config,
1597                                &mut retry_started_at,
1598                                proposed_wait,
1599                            ) else {
1600                                compaction_lifecycle
1601                                    .fail_if(provider_managed)
1602                                    .await;
1603                                return Err(AgentLoopError::llm_kind(
1604                                    crate::error::LlmErrorKind::Unavailable,
1605                                    format!(
1606                                        "{}; automatic recovery time budget exhausted after {} retries; the turn is safe to resume",
1607                                        stall_error.message,
1608                                        stream_retry_metadata.attempts
1609                                    ),
1610                                )
1611                                .with_retry_metadata(&stream_retry_metadata));
1612                            };
1613                            tracing::warn!(
1614                                session_id = %session_id,
1615                                turn_id = %context.turn_id,
1616                                attempt = stream_retry_metadata.attempts + 1,
1617                                max_retries = retry_config.max_retries,
1618                                wait_secs = wait_duration.as_secs_f64(),
1619                                "ReasonAtom: provider stream stall, retrying"
1620                            );
1621                            stream_retry_metadata.record_retry(wait_duration, None);
1622                            tokio::time::sleep(wait_duration).await;
1623                            continue 'stream_attempt;
1624                        }
1625                        compaction_lifecycle
1626                            .fail_if(compaction_started_at.is_some())
1627                            .await;
1628                        return Err(AgentLoopError::llm(stall_error.message));
1629                    },
1630                    _ = keepalive_ticker.tick() => {
1631                        if let Some(ref hb) = self.stream_heartbeater {
1632                            hb.heartbeat(crate::durability::StreamProgress {
1633                                accumulated_len: text.len() + thinking.len(),
1634                                last_delta_at: last_token_at_unix,
1635                            })
1636                            .await;
1637                            last_stream_heartbeat = Instant::now();
1638                        }
1639                        continue;
1640                    },
1641                };
1642                let event = match event {
1643                    Ok(event) => event,
1644                    Err(error) => {
1645                        compaction_lifecycle
1646                            .fail_if(compaction_started_at.is_some())
1647                            .await;
1648                        return Err(error);
1649                    }
1650                };
1651                replay_state.observe(&event);
1652                let advanced_stall_deadline = advances_stall_deadline(&event);
1653                if advanced_stall_deadline {
1654                    stall_sleep
1655                        .as_mut()
1656                        .reset(tokio::time::Instant::now() + stall_timeout);
1657                    last_token_at_unix = unix_now_secs();
1658                }
1659                match event {
1660                    LlmStreamEvent::TextDelta(delta) => {
1661                        if delta.is_empty() {
1662                            continue;
1663                        }
1664                        // Track time-to-first-token on first non-empty delta
1665                        if time_to_first_token_ms.is_none() {
1666                            let ttft = llm_start.elapsed().as_millis() as u64;
1667                            time_to_first_token_ms = Some(ttft);
1668                            tracing::info!(
1669                                session_id = %session_id,
1670                                time_to_first_token_ms = ttft,
1671                                "ReasonAtom: received first token from LLM"
1672                            );
1673                        }
1674                        text.push_str(&delta);
1675                        pending_delta.push_str(&delta);
1676
1677                        // Run output guardrails on the new accumulated text.
1678                        // Cheap by contract — runs in the streaming hot path.
1679                        // On block: suppress the pending delta (the bad text
1680                        // never reaches the client as a delta), record the
1681                        // trip, and break the loop. The replacement message is
1682                        // emitted below after the streaming block.
1683                        if !armed_guardrails.is_empty()
1684                            && let Some(t) =
1685                                evaluate_guardrails(&mut armed_guardrails, &text, &delta)
1686                        {
1687                            tracing::warn!(
1688                                session_id = %session_id,
1689                                turn_id = %context.turn_id,
1690                                guardrail_capability_id = %t.capability_id,
1691                                guardrail_id = %t.guardrail_id,
1692                                reason_code = %t.block.reason_code,
1693                                "ReasonAtom: output guardrail tripped, replacing assistant message"
1694                            );
1695                            pending_delta.clear();
1696                            termination = StreamTermination::GuardrailBlocked(t);
1697                            break;
1698                        }
1699
1700                        // Emit batched delta if interval elapsed
1701                        if !buffer_output_deltas
1702                            && last_delta_emit.elapsed().as_millis() as u64
1703                                >= DELTA_BATCH_INTERVAL_MS
1704                            && !pending_delta.is_empty()
1705                        {
1706                            if let Err(e) = self
1707                                .event_emitter
1708                                .emit(EventRequest::new(
1709                                    session_id,
1710                                    streaming_event_context.clone(),
1711                                    OutputMessageDeltaData {
1712                                        turn_id: context.turn_id,
1713                                        message_id: output_message_id,
1714                                        delta: pending_delta.clone(),
1715                                        accumulated: text.clone(),
1716                                        phase: streamed_phase,
1717                                    },
1718                                ))
1719                                .await
1720                            {
1721                                tracing::warn!(
1722                                    session_id = %session_id,
1723                                    error = %e,
1724                                    "ReasonAtom: failed to emit output.message.delta event"
1725                                );
1726                            }
1727                            pending_delta.clear();
1728                            last_delta_emit = Instant::now();
1729                        }
1730                    }
1731                    LlmStreamEvent::ReasoningDelta { delta, summary: _ } => {
1732                        if delta.is_empty() {
1733                            continue;
1734                        }
1735                        if let Some(t) = append_guarded_thinking_delta(
1736                            &mut armed_guardrails,
1737                            &mut thinking,
1738                            &mut pending_thinking_delta,
1739                            &delta,
1740                        ) {
1741                            tracing::warn!(
1742                                session_id = %session_id,
1743                                guardrail_capability_id = %t.capability_id,
1744                                guardrail_id = %t.guardrail_id,
1745                                "ReasonAtom: output guardrail tripped on thinking stream, replacing assistant message"
1746                            );
1747                            termination = StreamTermination::GuardrailBlocked(t);
1748                            break;
1749                        }
1750                        tracing::debug!(
1751                            session_id = %session_id,
1752                            delta_len = delta.len(),
1753                            total_thinking_len = thinking.len(),
1754                            "ReasonAtom: received ThinkingDelta from LLM"
1755                        );
1756
1757                        // Emit batched thinking delta if interval elapsed
1758                        if last_thinking_delta_emit.elapsed().as_millis() as u64
1759                            >= DELTA_BATCH_INTERVAL_MS
1760                            && !pending_thinking_delta.is_empty()
1761                        {
1762                            if let Err(e) = self
1763                                .event_emitter
1764                                .emit(EventRequest::new(
1765                                    session_id,
1766                                    streaming_event_context.clone(),
1767                                    ReasonThinkingDeltaData {
1768                                        turn_id: context.turn_id,
1769                                        delta: pending_thinking_delta.clone(),
1770                                        accumulated: thinking.clone(),
1771                                    },
1772                                ))
1773                                .await
1774                            {
1775                                tracing::warn!(
1776                                    session_id = %session_id,
1777                                    error = %e,
1778                                    "ReasonAtom: failed to emit reason.thinking.delta event"
1779                                );
1780                            }
1781                            pending_thinking_delta.clear();
1782                            last_thinking_delta_emit = Instant::now();
1783                        }
1784                    }
1785                    LlmStreamEvent::ReasoningItem(item) => {
1786                        if let Some(t) = inspect_guarded_reasoning_item(
1787                            &mut armed_guardrails,
1788                            &mut thinking,
1789                            &item,
1790                        ) {
1791                            tracing::warn!(
1792                                session_id = %session_id,
1793                                guardrail_capability_id = %t.capability_id,
1794                                guardrail_id = %t.guardrail_id,
1795                                "ReasonAtom: output guardrail tripped on completed reasoning item, replacing assistant message"
1796                            );
1797                            termination = StreamTermination::GuardrailBlocked(t);
1798                            break;
1799                        }
1800                        // One durable artifact per provider block, appended in
1801                        // order. Replay walks these; nothing is collapsed into
1802                        // a single per-message slot.
1803                        tracing::debug!(
1804                            session_id = %session_id,
1805                            provider = %item.provider,
1806                            item_id = ?item.item_id,
1807                            has_signature = item.signature.is_some(),
1808                            has_encrypted = item.encrypted.is_some(),
1809                            "ReasonAtom: captured reasoning artifact"
1810                        );
1811                        reasoning.push(item);
1812                    }
1813                    LlmStreamEvent::NativeToolCall(call) => {
1814                        if self.native_async.is_none() {
1815                            return Err(AgentLoopError::config(
1816                                "native async/custom tools require a configured native-call coordinator",
1817                            ));
1818                        }
1819                        let part = crate::message::ToolCallContentPart::from_native(call.clone())?;
1820                        if native_calls.insert(call.id().to_owned(), call).is_none() {
1821                            tool_calls.push(ToolCall {
1822                                id: part.id,
1823                                name: part.name,
1824                                arguments: part.arguments,
1825                            });
1826                        }
1827                    }
1828                    LlmStreamEvent::ToolCalls(calls) => {
1829                        if self.native_async.is_some() {
1830                            for call in &calls {
1831                                native_calls.entry(call.id.clone()).or_insert_with(|| {
1832                                    everruns_provider::native_async::NativeToolCall::Function {
1833                                        call_id: call.id.clone(),
1834                                        name: call.name.clone(),
1835                                        arguments: call.arguments.to_string(),
1836                                        asynchronous: false,
1837                                    }
1838                                });
1839                            }
1840                            for call in calls {
1841                                if !tool_calls.iter().any(|existing| existing.id == call.id) {
1842                                    tool_calls.push(call);
1843                                }
1844                            }
1845                        } else {
1846                            tool_calls = calls;
1847                        }
1848                    }
1849                    LlmStreamEvent::MessagePhase(phase) => {
1850                        // Provider revealed a native phase for the current
1851                        // assistant message mid-stream. Refine the streamed hint
1852                        // monotonically (never flip-flop, never back to None);
1853                        // subsequent output.message.delta events carry it. This is
1854                        // a hint only — it is NOT a completion signal and does not
1855                        // count as stream output. The completed message's phase stays
1856                        // authoritative, and the hint is deliberately not derived
1857                        // from later tool-call presence (EVE-448 anti-pattern).
1858                        streamed_phase = everruns_provider::ExecutionPhase::refine_streamed_hint(
1859                            streamed_phase,
1860                            phase,
1861                        );
1862                    }
1863                    LlmStreamEvent::ProviderCompactionStarted => {
1864                        compaction_lifecycle.start(&mut compaction_started_at).await;
1865                    }
1866                    LlmStreamEvent::Done(metadata) => {
1867                        // Emit any remaining pending delta before completing,
1868                        // unless a post-generation guardrail must first inspect
1869                        // the finalized assistant text.
1870                        if !buffer_output_deltas
1871                            && !pending_delta.is_empty()
1872                            && let Err(e) = self
1873                                .event_emitter
1874                                .emit(EventRequest::new(
1875                                    session_id,
1876                                    streaming_event_context.clone(),
1877                                    OutputMessageDeltaData {
1878                                        turn_id: context.turn_id,
1879                                        message_id: output_message_id,
1880                                        delta: pending_delta.clone(),
1881                                        accumulated: text.clone(),
1882                                        phase: streamed_phase,
1883                                    },
1884                                ))
1885                                .await
1886                        {
1887                            tracing::warn!(
1888                                session_id = %session_id,
1889                                error = %e,
1890                                "ReasonAtom: failed to emit final output.message.delta event"
1891                            );
1892                        }
1893
1894                        // Emit any remaining pending thinking delta before completing
1895                        if !pending_thinking_delta.is_empty()
1896                            && let Err(e) = self
1897                                .event_emitter
1898                                .emit(EventRequest::new(
1899                                    session_id,
1900                                    streaming_event_context.clone(),
1901                                    ReasonThinkingDeltaData {
1902                                        turn_id: context.turn_id,
1903                                        delta: pending_thinking_delta.clone(),
1904                                        accumulated: thinking.clone(),
1905                                    },
1906                                ))
1907                                .await
1908                        {
1909                            tracing::warn!(
1910                                session_id = %session_id,
1911                                error = %e,
1912                                "ReasonAtom: failed to emit final reason.thinking.delta event"
1913                            );
1914                        }
1915
1916                        // Emit reason.thinking.completed if we had any thinking content
1917                        if !thinking.is_empty()
1918                            && let Err(e) = self
1919                                .event_emitter
1920                                .emit(EventRequest::new(
1921                                    session_id,
1922                                    streaming_event_context.clone(),
1923                                    ReasonThinkingCompletedData {
1924                                        turn_id: context.turn_id,
1925                                        thinking: thinking.clone(),
1926                                    },
1927                                ))
1928                                .await
1929                        {
1930                            tracing::warn!(
1931                                session_id = %session_id,
1932                                error = %e,
1933                                "ReasonAtom: failed to emit reason.thinking.completed event"
1934                            );
1935                        }
1936                        termination = StreamTermination::Completed(metadata);
1937                        break;
1938                    }
1939                    LlmStreamEvent::Error(err) => {
1940                        // If we already collected valid tool calls or text before
1941                        // the error arrived, treat it as a partial success. This
1942                        // handles OpenAI Responses API behaviour where a trailing
1943                        // server_error can follow fully-streamed function calls.
1944                        let has_partial_output = compaction_started_at.is_none()
1945                            && (!tool_calls.is_empty() || !text.is_empty());
1946
1947                        if has_partial_output {
1948                            tracing::warn!(
1949                                session_id = %session_id,
1950                                error = %err,
1951                                tool_call_count = tool_calls.len(),
1952                                text_len = text.len(),
1953                                "ReasonAtom: trailing stream error after valid output — treating as partial success"
1954                            );
1955                            // Break out of the stream loop and use the output
1956                            // we already collected. completion_metadata will be
1957                            // None since we never got a Done event.
1958                            termination = StreamTermination::PartialSuccess;
1959                            break;
1960                        }
1961
1962                        if replay_state.should_retry(
1963                            &err,
1964                            stream_retry_metadata.attempts,
1965                            retry_config.max_retries,
1966                        ) {
1967                            let proposed_wait =
1968                                retry_config.calculate_backoff(stream_retry_metadata.attempts);
1969                            let Some(wait_duration) = reserve_retry_wait(
1970                                &retry_config,
1971                                &mut retry_started_at,
1972                                proposed_wait,
1973                            ) else {
1974                                return Err(AgentLoopError::llm_kind(
1975                                    err.kind(),
1976                                    format!(
1977                                        "{err}; automatic recovery time budget exhausted after {} retries; the turn is safe to resume",
1978                                        stream_retry_metadata.attempts
1979                                    ),
1980                                )
1981                                .with_retry_metadata(&stream_retry_metadata));
1982                            };
1983                            tracing::warn!(
1984                                session_id = %session_id,
1985                                turn_id = %context.turn_id,
1986                                attempt = stream_retry_metadata.attempts + 1,
1987                                max_retries = retry_config.max_retries,
1988                                wait_secs = wait_duration.as_secs_f64(),
1989                                error_code = err.code.as_deref().unwrap_or("none"),
1990                                error_status = err.status,
1991                                error = %err,
1992                                "ReasonAtom: transient stream error before output, retrying"
1993                            );
1994                            stream_retry_metadata.record_retry(wait_duration, None);
1995                            tokio::time::sleep(wait_duration).await;
1996                            continue 'stream_attempt;
1997                        }
1998
1999                        // No useful output collected — treat as a real failure.
2000                        let llm_duration_ms = llm_start.elapsed().as_millis() as u64;
2001                        let event_context = EventContext::from_execution_context(context)
2002                            .with_span(
2003                                trace_id.to_string(),
2004                                Uuid::now_v7().to_string(),
2005                                Some(reason_span_id.to_string()),
2006                            );
2007                        let tools_summary: Vec<ToolDefinitionSummary> =
2008                            runtime_agent.tools.iter().map(|t| t.into()).collect();
2009                        let generation_data = LlmGenerationData::failure(
2010                            messages_for_event.clone(),
2011                            tools_summary,
2012                            runtime_agent.model.clone(),
2013                            Some(model_with_provider.provider_type.to_string()),
2014                            err.to_string(),
2015                            Some(llm_duration_ms),
2016                            time_to_first_token_ms,
2017                        );
2018                        let _ = self
2019                            .event_emitter
2020                            .emit(EventRequest::new(
2021                                session_id,
2022                                event_context,
2023                                generation_data,
2024                            ))
2025                            .await;
2026                        if compaction_started_at.is_some() {
2027                            compaction_lifecycle.fail().await;
2028                        }
2029                        return Err(AgentLoopError::llm_kind(err.kind(), err.to_string()));
2030                    }
2031                    // `LlmStreamEvent` is `#[non_exhaustive]`, so a driver may
2032                    // emit a kind this build does not know. Ignoring it keeps
2033                    // the turn streaming rather than aborting; unreachable
2034                    // in-workspace, where every crate shares one provider
2035                    // version.
2036                    _ => {}
2037                }
2038                // Per-event heartbeat after processing the event, so accumulated_len
2039                // reflects the just-received tokens. Throttled to every 5s.
2040                if last_stream_heartbeat.elapsed().as_millis() as u64 >= 5_000
2041                    && let Some(ref hb) = self.stream_heartbeater
2042                {
2043                    hb.heartbeat(crate::durability::StreamProgress {
2044                        accumulated_len: text.len() + thinking.len(),
2045                        last_delta_at: last_token_at_unix,
2046                    })
2047                    .await;
2048                    last_stream_heartbeat = Instant::now();
2049                }
2050            }
2051            compaction_lifecycle
2052                .reject_incomplete(compaction_started_at, &termination)
2053                .await?;
2054            let (mut completion_metadata, tripped) = termination.into_parts();
2055            if let Some(metadata) = completion_metadata.as_mut() {
2056                metadata.retry_metadata =
2057                    merge_retry_metadata(metadata.retry_metadata.take(), &stream_retry_metadata);
2058            }
2059
2060            break 'stream_attempt (
2061                text,
2062                thinking,
2063                reasoning,
2064                tool_calls,
2065                completion_metadata,
2066                time_to_first_token_ms,
2067                pending_delta,
2068                tripped,
2069            );
2070        };
2071        let (mut text, mut thinking, mut reasoning, mut tool_calls) =
2072            (text, thinking, reasoning, tool_calls);
2073        compaction_lifecycle.record_observed(&mut llm_config, compaction_started_at);
2074
2075        // End-of-message citation annotation seam (see knowledge/runtime-resources/citations.md). Runs
2076        // once on the finalized final-answer text to attach claim-level citations
2077        // before guardrails inspect the complete client-visible payload. Skipped
2078        // when a guardrail already tripped,
2079        // when there are no citation providers, when there is no text, or while
2080        // the message still carries tool calls (citations attach to answer
2081        // prose, not intermediate tool-calling turns). The built-in feeds do not
2082        // rewrite text, so streamed deltas stay valid and no buffering is needed;
2083        // a future feed that rewrites (e.g. to strip inline markers) must also
2084        // opt into delta buffering.
2085        let mut citation_annotations: Vec<crate::message::TextAnnotation> = Vec::new();
2086        if tripped.is_none()
2087            && !annotation_providers.is_empty()
2088            && !text.is_empty()
2089            && tool_calls.is_empty()
2090        {
2091            // Apply deterministic capability-owned response filters before
2092            // annotation offsets are computed against the finalized text.
2093            text = filter_response_text(
2094                &self.capability_registry,
2095                &resolved_capability_configs,
2096                text,
2097            );
2098            let collected = collect_annotations(
2099                &annotation_providers,
2100                &runtime_agent.system_prompt,
2101                &text,
2102                &messages,
2103                self.utility_llm_service.as_ref(),
2104            )
2105            .await;
2106            text = collected.text;
2107            citation_annotations = collected.annotations;
2108
2109            // Post-generation guardrails must inspect citation metadata as well
2110            // as prose because annotations are persisted and rendered to clients.
2111            if !citation_annotations.is_empty() && !post_output_providers.is_empty() {
2112                let guarded_output = client_visible_guardrail_text(
2113                    &text,
2114                    &thinking,
2115                    &reasoning,
2116                    &citation_annotations,
2117                );
2118                let ctx = PostGenerationOutputContext {
2119                    system_prompt: &runtime_agent.system_prompt,
2120                    message_text: &guarded_output,
2121                    utility_llm_service: self.utility_llm_service.as_ref(),
2122                    decisions: self.decisions.as_ref(),
2123                };
2124                tripped = evaluate_post_generation_guardrails(&post_output_providers, &ctx).await;
2125            }
2126
2127            // Verification pass: stamp faithfulness verdicts on the collected
2128            // citations (no-op when no citation_verification capability is on).
2129            if tripped.is_none()
2130                && !citation_annotations.is_empty()
2131                && !citation_verifiers.is_empty()
2132            {
2133                citation_annotations = verify_annotations(
2134                    &citation_verifiers,
2135                    &text,
2136                    self.utility_llm_service.as_ref(),
2137                    citation_annotations,
2138                )
2139                .await;
2140            }
2141        }
2142
2143        // Messages without citation annotations still cross the same
2144        // post-generation output seam once.
2145        if tripped.is_none()
2146            && citation_annotations.is_empty()
2147            && !post_output_providers.is_empty()
2148            && (!text.is_empty() || !thinking.is_empty() || !reasoning.is_empty())
2149        {
2150            let guarded_output = client_visible_guardrail_text(&text, &thinking, &reasoning, &[]);
2151            let ctx = PostGenerationOutputContext {
2152                system_prompt: &runtime_agent.system_prompt,
2153                message_text: &guarded_output,
2154                utility_llm_service: self.utility_llm_service.as_ref(),
2155                decisions: self.decisions.as_ref(),
2156            };
2157            tripped = evaluate_post_generation_guardrails(&post_output_providers, &ctx).await;
2158        }
2159
2160        if tripped.is_some() {
2161            citation_annotations.clear();
2162        }
2163
2164        // Completed reasoning metadata is withheld until every output
2165        // guardrail has allowed the message. Otherwise a summary could escape
2166        // through `reason.item` before a later assistant-text block trips.
2167        if tripped.is_none() {
2168            for item in &reasoning {
2169                if let Err(e) = self
2170                    .event_emitter
2171                    .emit(EventRequest::new(
2172                        session_id,
2173                        streaming_event_context.clone(),
2174                        ReasonItemData {
2175                            turn_id: context.turn_id,
2176                            provider: item.provider.clone(),
2177                            model: Some(llm_config.model.clone()),
2178                            item_id: item.item_id.clone().unwrap_or_default(),
2179                            summary: item
2180                                .display_text()
2181                                .filter(|_| !matches!(item.text, Some(ReasoningText::Plain { .. })))
2182                                .into_iter()
2183                                .collect(),
2184                            token_count: item.tokens,
2185                        },
2186                    ))
2187                    .await
2188                {
2189                    tracing::warn!(
2190                        session_id = %session_id,
2191                        error = %e,
2192                        "ReasonAtom: failed to emit reason.item event"
2193                    );
2194                }
2195            }
2196        }
2197
2198        // Release buffered text only after post-generation guardrails allow it.
2199        // If they block, the replacement path below emits only sanitized text.
2200        if buffer_output_deltas
2201            && tripped.is_none()
2202            && !pending_delta.is_empty()
2203            && let Err(e) = self
2204                .event_emitter
2205                .emit(EventRequest::new(
2206                    session_id,
2207                    streaming_event_context.clone(),
2208                    OutputMessageDeltaData {
2209                        turn_id: context.turn_id,
2210                        message_id: output_message_id,
2211                        delta: pending_delta.clone(),
2212                        accumulated: text.clone(),
2213                        phase: streamed_phase,
2214                    },
2215                ))
2216                .await
2217        {
2218            tracing::warn!(
2219                session_id = %session_id,
2220                error = %e,
2221                "ReasonAtom: failed to emit guarded output.message.delta event"
2222            );
2223        }
2224
2225        // If a streaming output guardrail tripped, emit
2226        // output.message.replaced and overwrite the assistant output now so
2227        // every downstream event (llm.generation, output.message.completed)
2228        // carries the replacement instead of the model's withheld tokens.
2229        // The original tokens are never persisted or replayed.
2230        if let Some(ref t) = tripped {
2231            let replaced_event_context = EventContext::from_execution_context(context).with_span(
2232                trace_id.to_string(),
2233                Uuid::now_v7().to_string(),
2234                Some(reason_span_id.to_string()),
2235            );
2236            if let Err(e) = self
2237                .event_emitter
2238                .emit(EventRequest::new(
2239                    session_id,
2240                    replaced_event_context,
2241                    OutputMessageReplacedData {
2242                        turn_id: context.turn_id,
2243                        message_id: output_message_id,
2244                        guardrail_capability_id: t.capability_id.clone(),
2245                        guardrail_id: t.guardrail_id.clone(),
2246                        reason_code: t.block.reason_code.clone(),
2247                        replacement: t.block.replacement.clone(),
2248                    },
2249                ))
2250                .await
2251            {
2252                tracing::warn!(
2253                    session_id = %session_id,
2254                    error = %e,
2255                    "ReasonAtom: failed to emit output.message.replaced event"
2256                );
2257            }
2258            text = t.block.replacement.clone();
2259            tool_calls.clear();
2260            thinking.clear();
2261            reasoning.clear();
2262        }
2263
2264        // Finalized tool-call policy seam. This is the smallest provider-neutral
2265        // point where configured capabilities can normalize a complete call
2266        // batch before the assistant message and downstream events are built.
2267        let rejected_tool_calls = if tool_calls.is_empty() {
2268            Vec::new()
2269        } else {
2270            finalized_calls::apply_finalized_tool_calls_hooks(
2271                &self.capability_registry,
2272                self.event_emitter.as_ref(),
2273                session_id,
2274                context,
2275                &resolved_capability_configs,
2276                &runtime_agent.tools,
2277                &mut tool_calls,
2278                iteration,
2279            )
2280            .await
2281        };
2282        let finalized_tool_calls = tool_calls.clone();
2283        let rejected_tool_call_ids: HashSet<_> = rejected_tool_calls
2284            .iter()
2285            .map(|rejection| rejection.tool_call_id.clone())
2286            .collect();
2287        tool_calls.retain(|call| !rejected_tool_call_ids.contains(&call.id));
2288
2289        let llm_duration_ms = llm_start.elapsed().as_millis() as u64;
2290
2291        let response_id = completion_metadata
2292            .as_ref()
2293            .and_then(|meta| meta.response_id.clone());
2294        let finish_reason = completion_metadata
2295            .as_ref()
2296            .and_then(|meta| meta.finish_reason.clone());
2297
2298        // 15. Convert completion metadata to TokenUsage.
2299        //
2300        // Cost is tracked as two independent values: the provider's authoritative
2301        // inline cost when present (e.g. OpenRouter's usage.cost), and a price-table
2302        // estimate from the model profile computed whenever profile cost data
2303        // exists. Keeping both lets downstream consumers prefer the actual charge
2304        // while still reconciling estimate-vs-actual drift.
2305        let usage = completion_metadata.as_ref().and_then(|meta| {
2306            match (meta.prompt_tokens, meta.completion_tokens) {
2307                (Some(input), Some(output)) => {
2308                    let actual_cost_usd = meta.provider_cost_usd;
2309                    let estimated_cost_usd = crate::model_profiles::estimate_cost_usd(
2310                        &model_with_provider.provider_type,
2311                        &runtime_agent.model,
2312                        input,
2313                        output,
2314                        meta.cache_read_tokens.unwrap_or(0),
2315                        meta.cache_creation_tokens.unwrap_or(0),
2316                    );
2317                    Some(
2318                        TokenUsage::with_cache(
2319                            input,
2320                            output,
2321                            meta.cache_read_tokens,
2322                            meta.cache_creation_tokens,
2323                        )
2324                        .with_cost(actual_cost_usd, estimated_cost_usd),
2325                    )
2326                }
2327                _ => None,
2328            }
2329        });
2330
2331        // 16. Emit llm.generation event (child of reason span)
2332        let event_context = EventContext::from_execution_context(context).with_span(
2333            trace_id.to_string(),
2334            Uuid::now_v7().to_string(),
2335            Some(reason_span_id.to_string()),
2336        );
2337        let tools_summary: Vec<ToolDefinitionSummary> =
2338            runtime_agent.tools.iter().map(|t| t.into()).collect();
2339        let finish_reasons = Some(vec![finish_reason.clone().unwrap_or_else(|| {
2340            if finalized_tool_calls.is_empty() {
2341                "stop".to_string()
2342            } else {
2343                "tool_calls".to_string()
2344            }
2345        })]);
2346        let meta = completion_metadata.as_ref();
2347        let served = meta.and_then(|m| m.response_model.clone());
2348        let retry_info = completion_metadata
2349            .as_ref()
2350            .and_then(|meta| meta.retry_metadata.as_ref())
2351            .filter(|rm| rm.had_retries())
2352            .map(|rm| LlmRetryInfo {
2353                attempts: rm.attempts,
2354                total_wait_ms: rm.total_retry_wait.as_millis() as u64,
2355            });
2356        let mut generation_data = LlmGenerationData::success_with_retry(
2357            messages_for_event.clone(),
2358            tools_summary,
2359            Some(text.clone()).filter(|s| !s.is_empty()),
2360            finalized_tool_calls.clone(),
2361            runtime_agent.model.clone(),
2362            Some(model_with_provider.provider_type.to_string()),
2363            usage.clone(),
2364            Some(llm_duration_ms),
2365            time_to_first_token_ms,
2366            finish_reasons,
2367            response_id.clone(),
2368            retry_info,
2369        )
2370        .with_response_model(served);
2371
2372        // Add compaction info if compaction was performed. Compaction is a
2373        // separate billable model call on the same turn. Preserve whether the
2374        // generation cost was actual or estimated while recording their combined
2375        // best-effort cost for budgets and usage totals. `compaction.cost_usd`
2376        // keeps the split visible (EVE-895).
2377        if let Some(info) = compaction_info {
2378            if let Some(compaction_cost) = info.cost_usd {
2379                match generation_data.metadata.usage.as_mut() {
2380                    Some(usage) => {
2381                        add_compaction_cost(usage, compaction_cost);
2382                    }
2383                    // The generation itself reported no usage — a provider may
2384                    // price compaction without returning usage on the retry.
2385                    // Carry the cost on a usage record of its own rather than
2386                    // dropping it, which is the failure this fixes.
2387                    None => {
2388                        generation_data.metadata.usage = Some(crate::events::TokenUsage {
2389                            input_tokens: 0,
2390                            output_tokens: 0,
2391                            cache_read_tokens: None,
2392                            cache_creation_tokens: None,
2393                            actual_cost_usd: Some(compaction_cost),
2394                            estimated_cost_usd: None,
2395                            effective_cost_usd: None,
2396                        });
2397                    }
2398                }
2399            }
2400            generation_data = generation_data.with_compaction(info);
2401        }
2402
2403        if let Some(request_options) =
2404            build_request_options(&llm_config, &model_with_provider.provider_type.to_string())
2405        {
2406            generation_data = generation_data.with_request_options(request_options);
2407        }
2408
2409        if let Err(e) = self
2410            .event_emitter
2411            .emit(EventRequest::new(
2412                session_id,
2413                event_context,
2414                generation_data,
2415            ))
2416            .await
2417        {
2418            tracing::warn!(
2419                session_id = %session_id,
2420                error = %e,
2421                "ReasonAtom: failed to emit llm.generation event"
2422            );
2423        }
2424
2425        // 17. Build metadata with model and reasoning effort info
2426        let mut metadata = std::collections::HashMap::new();
2427        metadata.insert(
2428            "model".to_string(),
2429            serde_json::Value::String(runtime_agent.model.clone()),
2430        );
2431        if let Some(state) = &llm_config.reasoning_state {
2432            metadata.insert(
2433                reasoning_updates::STATE_KEY.to_string(),
2434                serde_json::json!(state),
2435            );
2436        }
2437        if let Some(effort) = llm_config
2438            .reasoning_state
2439            .as_ref()
2440            .and_then(|state| state.effective)
2441            .or(reasoning_effort)
2442        {
2443            metadata.insert(
2444                "reasoning_effort".to_string(),
2445                serde_json::Value::String(effort.as_str().to_string()),
2446            );
2447        }
2448        // Stamp the provider driver id and provider response id so the chat UI
2449        // can build a deep link to the provider's trace/logs for this message
2450        // (see ProviderTraceConfig). The resolved model carries the driver id,
2451        // not the concrete provider instance id, so the UI keys trace config by
2452        // driver. `response_id` is the provider's generation id (e.g.
2453        // OpenRouter's "gen-..."); absent for providers that do not return one.
2454        metadata.insert(
2455            "provider".to_string(),
2456            serde_json::Value::String(model_with_provider.provider_type.to_string()),
2457        );
2458        if let Some(ref rid) = response_id {
2459            metadata.insert(
2460                "response_id".to_string(),
2461                serde_json::Value::String(rid.clone()),
2462            );
2463        }
2464
2465        // 18. Store and emit output.message.completed event with metadata and usage.
2466        // Apply capability-owned response filters before persisting/returning
2467        // the finalized assistant text.
2468        let text = filter_response_text(
2469            &self.capability_registry,
2470            &resolved_capability_configs,
2471            text,
2472        );
2473        let (provider_opaque_content, provider_checkpoint_candidate) =
2474            provider_managed_compaction::replay_artifacts(
2475                completion_metadata.as_ref(),
2476                tripped.is_none(),
2477                !rejected_tool_calls.is_empty(),
2478            );
2479        let has_tool_calls = !finalized_tool_calls.is_empty();
2480        let mut assistant_message = if has_tool_calls {
2481            RuntimeMessage::assistant_with_tools(&text, finalized_tool_calls.clone())
2482        } else {
2483            RuntimeMessage::assistant(&text)
2484        }
2485        .with_id(output_message_id);
2486        for part in &mut assistant_message.content {
2487            if let crate::message::ContentPart::ToolCall(call) = part {
2488                call.native = native_calls.get(&call.id).cloned();
2489            }
2490        }
2491        // Attach citation annotations produced by the annotation seam above to
2492        // the message's text part (see knowledge/runtime-resources/citations.md).
2493        if !citation_annotations.is_empty() {
2494            for part in assistant_message.content.iter_mut() {
2495                if let crate::message::ContentPart::Text(t) = part {
2496                    t.annotations = std::mem::take(&mut citation_annotations);
2497                    break;
2498                }
2499            }
2500        }
2501        // Use the API-provided phase when available (preserving the provider's value),
2502        // otherwise derive from state: Commentary for intermediate iterations (with tool
2503        // calls), FinalAnswer for the completed response.
2504        let provider_type_for_reasoning = model_with_provider.provider_type.to_string();
2505        // Record where the phase came from. A provider-reported phase is a real
2506        // decision; a derived one is just `has_tool_calls` wearing a
2507        // decision's name, and consumers must be able to tell.
2508        let provider_phase = completion_metadata
2509            .as_ref()
2510            .and_then(|meta| meta.phase.as_deref())
2511            .and_then(everruns_provider::ExecutionPhase::from_provider_str);
2512        let (phase, phase_source) = match provider_phase {
2513            Some(phase) => (phase, everruns_provider::PhaseSource::Provider),
2514            None => (
2515                everruns_provider::ExecutionPhase::from_has_tool_calls(has_tool_calls),
2516                everruns_provider::PhaseSource::Derived,
2517            ),
2518        };
2519        assistant_message.phase = Some(phase);
2520        assistant_message.phase_source = Some(phase_source);
2521        assistant_message.metadata = Some(metadata);
2522        // Reasoning artifacts lead the message content, preserving their order
2523        // among themselves. Every current provider emits reasoning ahead of the
2524        // text and tool calls it produced, and all three require it replayed in
2525        // that position, so leading is the faithful placement.
2526        // Providers that stream reasoning without any replayable artifact
2527        // (Chat Completions `reasoning_content`) would otherwise render live
2528        // and vanish on reload. Persist what was shown, with no replay state,
2529        // so every provider's readable reasoning survives uniformly.
2530        if reasoning.is_empty() && !thinking.is_empty() {
2531            reasoning.push(
2532                ReasoningContentPart::opaque(provider_type_for_reasoning.clone()).with_text(
2533                    ReasoningText::Plain {
2534                        text: thinking.clone(),
2535                    },
2536                ),
2537            );
2538        }
2539        if !reasoning.is_empty() {
2540            let mut content = Vec::with_capacity(reasoning.len() + assistant_message.content.len());
2541            content.extend(reasoning.drain(..).map(ContentPart::Reasoning));
2542            content.append(&mut assistant_message.content);
2543            assistant_message.content = content;
2544        }
2545        if let Some(content) = provider_opaque_content {
2546            assistant_message
2547                .content
2548                .push(ContentPart::ProviderOpaque(content));
2549        }
2550        // Emit output.message.completed event (this stores the message as an event with proper turn context)
2551        // Include token usage for tracking (child of reason span)
2552        let message_event_context = EventContext::from_execution_context(context).with_span(
2553            trace_id.to_string(),
2554            Uuid::now_v7().to_string(),
2555            Some(reason_span_id.to_string()),
2556        );
2557        let mut output_message_data = OutputMessageCompletedData::new(assistant_message);
2558        if let Some(ref u) = usage {
2559            output_message_data = output_message_data.with_usage(u.clone());
2560        }
2561        let result = ReasonResult {
2562            native_counts: None,
2563            success: true,
2564            text,
2565            tool_calls,
2566            has_tool_calls,
2567            tool_definitions: runtime_agent.tools.clone(),
2568            max_iterations: runtime_agent.max_iterations,
2569            error: None,
2570            user_facing_error: None,
2571            error_disclosure: None,
2572            usage,
2573            output_message_id: Some(output_message_id),
2574            time_to_first_token_ms,
2575            response_id,
2576            finish_reason,
2577            locale: resolved_locale,
2578            network_access: runtime_agent.network_access.clone(),
2579            parallel_tool_calls: runtime_agent.parallel_tool_calls,
2580        };
2581        if let Some(coordinator) = &self.native_async {
2582            coordinator
2583                .lock()
2584                .await
2585                .stage_transcript_result(
2586                    serde_json::to_value(&result)
2587                        .map_err(|error| AgentLoopError::store(error.to_string()))?,
2588                )
2589                .await?;
2590        }
2591        let completed_output_event = self
2592            .event_emitter
2593            .emit(EventRequest::new(
2594                session_id,
2595                message_event_context,
2596                output_message_data,
2597            ))
2598            .await?;
2599        compaction_lifecycle
2600            .finish_after_output(
2601                self.compaction_checkpoint_store.as_deref(),
2602                completed_output_event.sequence,
2603                provider_checkpoint_candidate,
2604                compaction_started_at,
2605                restored_checkpoint.is_some(),
2606                completion_metadata.as_ref(),
2607            )
2608            .await;
2609
2610        if let Some(coordinator) = &self.native_async {
2611            coordinator
2612                .lock()
2613                .await
2614                .transcript_committed(&output_message_id.to_string())
2615                .await?;
2616        }
2617        for rejection in rejected_tool_calls {
2618            let Some(call) = finalized_tool_calls
2619                .iter()
2620                .find(|call| call.id == rejection.tool_call_id)
2621            else {
2622                continue;
2623            };
2624            self.event_emitter
2625                .emit(EventRequest::new(
2626                    session_id,
2627                    EventContext::from_execution_context(context),
2628                    ToolCompletedData::failure(
2629                        call.id.clone(),
2630                        call.name.clone(),
2631                        "error".to_string(),
2632                        rejection.error,
2633                        None,
2634                    ),
2635                ))
2636                .await?;
2637        }
2638        tracing::info!(
2639            session_id = %session_id,
2640            turn_id = %context.turn_id,
2641            has_tool_calls = %result.has_tool_calls,
2642            tool_count = %result.tool_calls.len(),
2643            "ReasonAtom: LLM call completed"
2644        );
2645
2646        Ok(result)
2647    }
2648
2649    /// Finalize a partial assistant stream without making a new provider call (EVE-532).
2650    ///
2651    /// Emits `output.message.started`, `output.message.completed` from the persisted
2652    /// `accumulated` text, and `reason.recovered { mode: Finalize }`.
2653    async fn finalize_partial_stream(
2654        &self,
2655        session_id: SessionId,
2656        context: &ExecutionContext,
2657        partial: PartialStreamState,
2658        iteration: u32,
2659        runtime_agent: &crate::RuntimeAgent,
2660        resolved_capability_configs: &[crate::CapabilityRef],
2661    ) -> Result<ReasonResult> {
2662        let event_context = EventContext::from_execution_context(context);
2663        let turn_id = context.turn_id;
2664        let message_id = partial.message_id;
2665
2666        // Signal that output is starting (keeps the streaming protocol intact).
2667        let _ = self
2668            .event_emitter
2669            .emit(EventRequest::new(
2670                session_id,
2671                event_context.clone(),
2672                OutputMessageStartedData {
2673                    reasoning_state: partial.reasoning_state.clone(),
2674                    turn_id,
2675                    message_id,
2676                    model: None,
2677                    iteration: Some(iteration),
2678                    // Recovery/finalize path reconstructs the started signal only;
2679                    // the streamed phase hint is unavailable here (None).
2680                    phase: None,
2681                },
2682            ))
2683            .await;
2684
2685        // Build the assistant message from capability-filtered accumulated text
2686        // and persist it via the canonical event path.
2687        let accumulated = filter_response_text(
2688            &self.capability_registry,
2689            resolved_capability_configs,
2690            partial.accumulated,
2691        );
2692        let mut assistant_message = RuntimeMessage::assistant(&accumulated).with_id(message_id);
2693        if let Some(state) = partial.reasoning_state {
2694            assistant_message.metadata = Some(HashMap::from([
2695                ("model".into(), serde_json::json!("gpt-6-astra")),
2696                ("provider".into(), serde_json::json!("openai")),
2697                (
2698                    reasoning_updates::STATE_KEY.into(),
2699                    serde_json::json!(state),
2700                ),
2701                (
2702                    "reasoning_effort".into(),
2703                    serde_json::json!(state.effective),
2704                ),
2705            ]));
2706        }
2707        let output_message_id = message_id;
2708        self.event_emitter
2709            .emit(EventRequest::new(
2710                session_id,
2711                event_context.clone(),
2712                OutputMessageCompletedData::new(assistant_message),
2713            ))
2714            .await?;
2715
2716        // Emit observability event.
2717        let accumulated_len = accumulated.len();
2718        let _ = self
2719            .event_emitter
2720            .emit(EventRequest::new(
2721                session_id,
2722                event_context.clone(),
2723                ReasonRecoveredData {
2724                    turn_id,
2725                    mode: RecoveryMode::Finalize,
2726                    accumulated_len,
2727                },
2728            ))
2729            .await;
2730
2731        tracing::info!(
2732            session_id = %session_id,
2733            turn_id = %turn_id,
2734            accumulated_len,
2735            "ReasonAtom: finalized partial stream from persisted accumulated text"
2736        );
2737
2738        Ok(ReasonResult {
2739            native_counts: None,
2740            success: true,
2741            text: accumulated,
2742            tool_calls: vec![],
2743            has_tool_calls: false,
2744            tool_definitions: runtime_agent.tools.clone(),
2745            max_iterations: runtime_agent.max_iterations,
2746            error: None,
2747            user_facing_error: None,
2748            error_disclosure: None,
2749            usage: None,
2750            output_message_id: Some(output_message_id),
2751            time_to_first_token_ms: None,
2752            response_id: None,
2753            finish_reason: Some("stop".to_string()),
2754            locale: None,
2755            network_access: None,
2756            // Finalize path has no tool calls, so the preference is irrelevant.
2757            parallel_tool_calls: None,
2758        })
2759    }
2760
2761    /// Resolve image_file references to actual image data
2762    ///
2763    /// This method extracts all image_file IDs from the messages and resolves
2764    /// them to base64-encoded image data using the configured ImageResolver.
2765    ///
2766    /// # Returns
2767    ///
2768    /// A HashMap mapping image IDs to ResolvedImage data. If no ImageResolver
2769    /// is configured, or if resolution fails for some images, those images
2770    /// will simply be missing from the map (and converted to placeholder text).
2771    async fn resolve_images(&self, messages: &[RuntimeMessage]) -> HashMap<Uuid, ResolvedImage> {
2772        let mut resolved = HashMap::new();
2773
2774        // Check if we have an image resolver
2775        let resolver = match &self.image_resolver {
2776            Some(r) => r,
2777            None => return resolved,
2778        };
2779
2780        // Collect all unique image_file IDs from all messages
2781        let image_ids: Vec<Uuid> = messages
2782            .iter()
2783            .flat_map(crate::llm_conversions::extract_image_file_ids)
2784            .collect::<std::collections::HashSet<_>>()
2785            .into_iter()
2786            .collect();
2787
2788        if image_ids.is_empty() {
2789            return resolved;
2790        }
2791
2792        tracing::debug!(
2793            image_count = image_ids.len(),
2794            "ReasonAtom: resolving image_file references"
2795        );
2796
2797        // Resolve each image
2798        for image_id in image_ids {
2799            match resolver.resolve_image(image_id).await {
2800                Ok(Some(image)) => {
2801                    resolved.insert(image_id, image);
2802                }
2803                Ok(None) => {
2804                    tracing::warn!(
2805                        image_id = %image_id,
2806                        "ReasonAtom: image not found during resolution"
2807                    );
2808                }
2809                Err(e) => {
2810                    tracing::warn!(
2811                        image_id = %image_id,
2812                        error = %e,
2813                        "ReasonAtom: failed to resolve image"
2814                    );
2815                }
2816            }
2817        }
2818
2819        tracing::debug!(
2820            resolved_count = resolved.len(),
2821            "ReasonAtom: image resolution complete"
2822        );
2823
2824        resolved
2825    }
2826
2827    async fn resolve_files(&self, messages: &[RuntimeMessage]) -> HashMap<Uuid, ResolvedFile> {
2828        let Some(resolver) = &self.file_resolver else {
2829            return HashMap::new();
2830        };
2831
2832        let file_ids: Vec<Uuid> = messages
2833            .iter()
2834            .flat_map(crate::llm_conversions::extract_file_ids)
2835            .collect::<std::collections::HashSet<_>>()
2836            .into_iter()
2837            .collect();
2838
2839        if file_ids.is_empty() {
2840            return HashMap::new();
2841        }
2842
2843        match resolver.resolve_files(&file_ids).await {
2844            Ok(map) => map,
2845            Err(e) => {
2846                tracing::warn!(
2847                    target: "reason",
2848                    "ReasonAtom: file resolution failed: {e}"
2849                );
2850                HashMap::new()
2851            }
2852        }
2853    }
2854}
2855
2856// ============================================================================
2857// Tests
2858// ============================================================================
2859
2860#[cfg(test)]
2861mod tests;