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