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