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