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