Skip to main content

everruns_engine/execution/
reason.rs

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