1use 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#[derive(Debug, Clone, Serialize, Deserialize)]
108pub struct ReasonInput {
109 pub context: ExecutionContext,
111 pub harness_id: HarnessId,
113 #[serde(skip_serializing_if = "Option::is_none")]
115 pub agent_id: Option<AgentId>,
116 #[serde(default)]
118 pub org_id: i64,
119 #[serde(default)]
123 pub mcp_tool_definitions: Vec<ToolDefinition>,
124 #[serde(skip_serializing_if = "Option::is_none")]
127 pub previous_response_id: Option<String>,
128 #[serde(default = "default_iteration")]
131 pub iteration: u32,
132}
133
134fn default_iteration() -> u32 {
135 1
136}
137
138#[derive(Debug, Clone, Serialize, Deserialize)]
140pub struct NativeExecutionCounts {
141 pub llm_calls: u32,
142 pub tool_calls: u32,
143}
144
145#[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 pub success: bool,
152 pub text: String,
154 #[serde(default)]
156 pub tool_calls: Vec<ToolCall>,
157 pub has_tool_calls: bool,
159 #[serde(default)]
161 pub tool_definitions: Vec<ToolDefinition>,
162 #[serde(default = "default_max_iterations")]
164 pub max_iterations: usize,
165 #[serde(skip_serializing_if = "Option::is_none")]
167 pub error: Option<String>,
168 #[serde(default, skip_serializing_if = "Option::is_none")]
172 pub user_facing_error: Option<UserFacingError>,
173 #[serde(default, skip_serializing_if = "Option::is_none")]
175 pub error_disclosure: Option<ErrorDisclosure>,
176 #[serde(skip_serializing_if = "Option::is_none")]
178 pub usage: Option<TokenUsage>,
179 #[serde(skip_serializing_if = "Option::is_none")]
181 pub output_message_id: Option<MessageId>,
182 #[serde(skip_serializing_if = "Option::is_none")]
184 pub time_to_first_token_ms: Option<u64>,
185 #[serde(skip_serializing_if = "Option::is_none")]
187 pub response_id: Option<String>,
188 #[serde(default, skip_serializing_if = "Option::is_none")]
190 pub finish_reason: Option<String>,
191 #[serde(skip_serializing_if = "Option::is_none")]
193 pub locale: Option<String>,
194 #[serde(default, skip_serializing_if = "Option::is_none")]
196 pub network_access: Option<crate::network_access::NetworkAccessList>,
197 #[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
208pub 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 image_resolver: Option<Arc<dyn ImageResolver>>,
234 file_resolver: Option<Arc<dyn FileResolver>>,
235 stream_heartbeater: Option<Arc<dyn crate::durability::StreamHeartbeater>>,
237 provider_stall_timeout: Option<std::time::Duration>,
239 provider_retry_config: LlmRetryConfig,
241 durable_tool_result_store: Option<Arc<dyn DurableToolResultStore>>,
243 partial_stream_store: Option<Arc<dyn PartialStreamStore>>,
245 reasoning_effort_handle: Option<crate::tool_context::ReasoningEffortHandle>,
249 utility_llm_service: Option<Arc<dyn crate::UtilityLlmService>>,
253 decisions: Option<Arc<dyn crate::DecisionsService>>,
257 schedule_store: Option<Arc<dyn crate::session_services::SessionScheduleStore>>,
263 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 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 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 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 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 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 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 pub fn with_provider_retry_config(mut self, config: LlmRetryConfig) -> Self {
384 self.provider_retry_config = config;
385 self
386 }
387
388 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 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 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 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 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 pub fn name(&self) -> &'static str {
439 "reason"
440 }
441
442 pub async fn execute(&self, input: ReasonInput) -> Result<ReasonResult> {
444 self.execute_inner(input, None).await
445 }
446}
447
448impl ReasonAtom {
449 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 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 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(); 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 let reason_start = Instant::now();
563
564 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, },
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 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 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 let result = match call_result {
643 Ok(result) => {
644 let reason_duration_ms = reason_start.elapsed().as_millis() as u64;
646
647 let completed_context = EventContext::from_execution_context(&context).with_span(
649 trace_id.clone(),
650 reason_span_id.clone(), 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 let reason_duration_ms = reason_start.elapsed().as_millis() as u64;
679
680 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 let is_transient = e.is_transient_llm_error()
699 || (e.llm_error_kind().is_none() && is_transient_error_message(&error_msg));
700
701 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 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 let error_msg_context = EventContext::from_execution_context(&context)
749 .with_span(
750 trace_id.clone(),
751 Uuid::now_v7().to_string(), Some(reason_span_id.clone()), );
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 let completed_context = EventContext::from_execution_context(&context).with_span(
780 trace_id.clone(),
781 reason_span_id.clone(), 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 #[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 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 && (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 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 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 saved.pending = saved.effective;
994 replay.state = saved;
995 }
996 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) => {} Err(e) => {
1019 if reasoning_replay.is_some() {
1020 return Err(e);
1021 }
1022 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 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 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 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 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 let resolved_images = self.resolve_images(&context_messages).await;
1107 let resolved_files = self.resolve_files(&context_messages).await;
1108
1109 let mut llm_messages = Vec::new();
1111
1112 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 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 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 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 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 for (k, v) in &embedder_metadata {
1195 llm_config_builder = llm_config_builder.with_metadata(k, v.clone());
1196 }
1197
1198 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 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 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 let streaming_event_context = EventContext::from_execution_context(context);
1258
1259 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 let buffer_output_deltas = !post_output_providers.is_empty();
1285 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 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 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 let llm_start = Instant::now();
1362
1363 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 const DELTA_BATCH_INTERVAL_MS: u64 = 100;
1398 let retry_config = self.provider_retry_config.clone();
1399 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 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 let mut reasoning: Vec<ReasoningContentPart> = Vec::new();
1574 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 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; let mut last_stream_heartbeat = Instant::now();
1596 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 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 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 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 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 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 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 streamed_phase = everruns_provider::ExecutionPhase::refine_streamed_hint(
1887 streamed_phase,
1888 phase,
1889 );
1890 }
1891 LlmStreamEvent::Done(metadata) => {
1892 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 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 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 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 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 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 _ => {}
2058 }
2059 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 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 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 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 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 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 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 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 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 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 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 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 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 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 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 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 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 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 let provider_type_for_reasoning = model_with_provider.provider_type.to_string();
2525 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 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 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 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 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 phase: None,
2690 },
2691 ))
2692 .await;
2693
2694 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 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 parallel_tool_calls: None,
2767 })
2768 }
2769
2770 async fn resolve_images(&self, messages: &[RuntimeMessage]) -> HashMap<Uuid, ResolvedImage> {
2781 let mut resolved = HashMap::new();
2782
2783 let resolver = match &self.image_resolver {
2785 Some(r) => r,
2786 None => return resolved,
2787 };
2788
2789 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 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#[cfg(test)]
2870mod tests;