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, post_generation_guardrail_text,
55};
56use crate::phase_effects::{PhaseEffectEmitter, PhaseEffectSink};
57use crate::runtime_context::{AssembledTurnContext, TurnContextRequest, TurnContextResolver};
58use crate::tool_types::{ToolCall, ToolDefinition};
59use crate::typed_id::{AgentId, HarnessId, MessageId, SessionId};
60use crate::{ErrorDisclosure, UserFacingError, UserFacingErrorContext};
61use crate::{
62 durability::DurableToolResultStore,
63 durability::PartialStreamState,
64 durability::PartialStreamStore,
65 file_services::{FileResolver, ResolvedFile},
66 image_services::ImageResolver,
67 image_services::ResolvedImage,
68};
69use everruns_provider::reasoning::{ReasoningContentPart, ReasoningText};
70
71mod compaction;
72mod error_policy;
73mod finalized_calls;
74mod observability;
75mod output_hooks;
76mod reasoning_updates;
77mod request_controls;
78mod stream_state;
79mod transcript;
80
81use compaction::{
82 ProactiveCompactionContext, ReactiveCompactionContext, apply_proactive_compaction,
83 apply_reactive_compaction,
84};
85use error_policy::{
86 error_disclosure_override, filter_response_text, is_error_placeholder_message,
87 resolve_error_disclosure,
88};
89use observability::{build_request_options, capability_usage_snapshot_records};
90use output_hooks::collect_output_hooks;
91use request_controls::resolve_request_controls;
92use stream_state::{
93 StreamReplayState, StreamTermination, advances_stall_deadline, append_guarded_thinking_delta,
94 inspect_guarded_reasoning_item, merge_retry_metadata,
95};
96use transcript::repair_dangling_tool_calls;
97
98fn client_visible_guardrail_text(
103 text: &str,
104 streamed_reasoning: &str,
105 reasoning: &[ReasoningContentPart],
106 citation_annotations: &[crate::message::TextAnnotation],
107) -> String {
108 let mut guarded = streamed_reasoning.to_string();
109 if guarded.is_empty() {
110 for item_text in reasoning
111 .iter()
112 .filter_map(ReasoningContentPart::display_text)
113 {
114 if !guarded.is_empty() {
115 guarded.push_str("\n\n");
116 }
117 guarded.push_str(&item_text);
118 }
119 }
120
121 let prose = post_generation_guardrail_text(text, citation_annotations);
122 if !guarded.is_empty() && !prose.is_empty() {
123 guarded.push_str("\n\n");
124 }
125 guarded.push_str(&prose);
126 guarded
127}
128
129fn unix_now_secs() -> u64 {
130 std::time::SystemTime::now()
131 .duration_since(std::time::UNIX_EPOCH)
132 .unwrap_or_default()
133 .as_secs()
134}
135
136#[derive(Debug, Clone, Serialize, Deserialize)]
138pub struct ReasonInput {
139 pub context: ExecutionContext,
141 pub harness_id: HarnessId,
143 #[serde(skip_serializing_if = "Option::is_none")]
145 pub agent_id: Option<AgentId>,
146 #[serde(default)]
148 pub org_id: i64,
149 #[serde(default)]
153 pub mcp_tool_definitions: Vec<ToolDefinition>,
154 #[serde(skip_serializing_if = "Option::is_none")]
157 pub previous_response_id: Option<String>,
158 #[serde(default = "default_iteration")]
161 pub iteration: u32,
162}
163
164fn default_iteration() -> u32 {
165 1
166}
167
168#[derive(Debug, Clone, Serialize, Deserialize)]
170pub struct NativeExecutionCounts {
171 pub llm_calls: u32,
172 pub tool_calls: u32,
173}
174
175#[derive(Debug, Clone, Default, Serialize, Deserialize)]
177pub struct ReasonResult {
178 #[serde(default, skip_serializing_if = "Option::is_none")]
179 pub native_counts: Option<NativeExecutionCounts>,
180 pub success: bool,
182 pub text: String,
184 #[serde(default)]
186 pub tool_calls: Vec<ToolCall>,
187 pub has_tool_calls: bool,
189 #[serde(default)]
191 pub tool_definitions: Vec<ToolDefinition>,
192 #[serde(default = "default_max_iterations")]
194 pub max_iterations: usize,
195 #[serde(skip_serializing_if = "Option::is_none")]
197 pub error: Option<String>,
198 #[serde(default, skip_serializing_if = "Option::is_none")]
202 pub user_facing_error: Option<UserFacingError>,
203 #[serde(default, skip_serializing_if = "Option::is_none")]
205 pub error_disclosure: Option<ErrorDisclosure>,
206 #[serde(skip_serializing_if = "Option::is_none")]
208 pub usage: Option<TokenUsage>,
209 #[serde(skip_serializing_if = "Option::is_none")]
211 pub output_message_id: Option<MessageId>,
212 #[serde(skip_serializing_if = "Option::is_none")]
214 pub time_to_first_token_ms: Option<u64>,
215 #[serde(skip_serializing_if = "Option::is_none")]
217 pub response_id: Option<String>,
218 #[serde(default, skip_serializing_if = "Option::is_none")]
220 pub finish_reason: Option<String>,
221 #[serde(skip_serializing_if = "Option::is_none")]
223 pub locale: Option<String>,
224 #[serde(default, skip_serializing_if = "Option::is_none")]
226 pub network_access: Option<crate::network_access::NetworkAccessList>,
227 #[serde(default, skip_serializing_if = "Option::is_none")]
231 pub parallel_tool_calls: Option<bool>,
232}
233
234fn default_max_iterations() -> usize {
235 500
236}
237
238pub struct ReasonAtom {
257 native_async: Option<Arc<tokio::sync::Mutex<crate::native_async::NativeAsyncCoordinator>>>,
258 context_resolver: Arc<dyn TurnContextResolver>,
259 message_retriever: Arc<dyn MessageRetriever>,
260 capability_registry: CapabilityRegistry,
261 event_emitter: PhaseEffectEmitter<dyn PhaseEffectSink>,
262 image_resolver: Option<Arc<dyn ImageResolver>>,
264 file_resolver: Option<Arc<dyn FileResolver>>,
265 stream_heartbeater: Option<Arc<dyn crate::durability::StreamHeartbeater>>,
267 provider_stall_timeout: Option<std::time::Duration>,
269 provider_retry_config: LlmRetryConfig,
271 durable_tool_result_store: Option<Arc<dyn DurableToolResultStore>>,
273 partial_stream_store: Option<Arc<dyn PartialStreamStore>>,
275 reasoning_effort_handle: Option<crate::tool_context::ReasoningEffortHandle>,
279 utility_llm_service: Option<Arc<dyn crate::UtilityLlmService>>,
283 decisions: Option<Arc<dyn crate::DecisionsService>>,
287 schedule_store: Option<Arc<dyn crate::session_services::SessionScheduleStore>>,
293 compaction_checkpoint_store: Option<Arc<dyn crate::CompactionCheckpointStore>>,
295}
296
297impl ReasonAtom {
298 pub fn with_native_async(
299 mut self,
300 coordinator: Arc<tokio::sync::Mutex<crate::native_async::NativeAsyncCoordinator>>,
301 ) -> Self {
302 self.native_async = Some(coordinator);
303 self
304 }
305
306 pub fn new(
308 context_resolver: impl TurnContextResolver + 'static,
309 message_retriever: impl MessageRetriever + 'static,
310 capability_registry: CapabilityRegistry,
311 event_emitter: impl PhaseEffectSink + 'static,
312 ) -> Self {
313 Self {
314 native_async: None,
315 context_resolver: Arc::new(context_resolver),
316 message_retriever: Arc::new(message_retriever),
317 capability_registry,
318 event_emitter: PhaseEffectEmitter::new(Arc::new(event_emitter)),
319 image_resolver: None,
320 file_resolver: None,
321 stream_heartbeater: None,
322 provider_stall_timeout: None,
323 provider_retry_config: LlmRetryConfig::default(),
324 durable_tool_result_store: None,
325 partial_stream_store: None,
326 reasoning_effort_handle: None,
327 utility_llm_service: None,
328 decisions: None,
329 schedule_store: None,
330 compaction_checkpoint_store: None,
331 }
332 }
333
334 pub fn with_schedule_store(
337 mut self,
338 store: Arc<dyn crate::session_services::SessionScheduleStore>,
339 ) -> Self {
340 self.schedule_store = Some(store);
341 self
342 }
343
344 pub fn with_compaction_checkpoint_store(
345 mut self,
346 store: Arc<dyn crate::CompactionCheckpointStore>,
347 ) -> Self {
348 self.compaction_checkpoint_store = Some(store);
349 self
350 }
351
352 fn collect_llm_error_hooks(
358 &self,
359 resolved_capability_configs: &[crate::CapabilityRef],
360 ) -> Vec<(
361 Arc<dyn crate::llm_error_hook::LlmErrorHook>,
362 serde_json::Value,
363 )> {
364 resolved_capability_configs
365 .iter()
366 .filter_map(|cfg| {
367 let cap = self.capability_registry.get(cfg.capability_id())?;
368 let hook = cap.llm_error_hook()?;
369 Some((hook, cfg.config_value().clone()))
370 })
371 .collect()
372 }
373
374 pub fn with_image_resolver(mut self, resolver: Arc<dyn ImageResolver>) -> Self {
387 self.image_resolver = Some(resolver);
388 self
389 }
390
391 pub fn with_file_resolver(mut self, resolver: Arc<dyn FileResolver>) -> Self {
392 self.file_resolver = Some(resolver);
393 self
394 }
395
396 pub fn with_stream_heartbeater(
398 mut self,
399 heartbeater: Arc<dyn crate::durability::StreamHeartbeater>,
400 ) -> Self {
401 self.stream_heartbeater = Some(heartbeater);
402 self
403 }
404
405 pub fn with_provider_stall_timeout(mut self, timeout: std::time::Duration) -> Self {
408 self.provider_stall_timeout = Some(timeout);
409 self
410 }
411
412 pub fn with_provider_retry_config(mut self, config: LlmRetryConfig) -> Self {
414 self.provider_retry_config = config;
415 self
416 }
417
418 pub fn with_durable_tool_result_store(
424 mut self,
425 store: Arc<dyn DurableToolResultStore>,
426 ) -> Self {
427 self.durable_tool_result_store = Some(store);
428 self
429 }
430
431 pub fn with_partial_stream_store(mut self, store: Arc<dyn PartialStreamStore>) -> Self {
433 self.partial_stream_store = Some(store);
434 self
435 }
436
437 pub fn with_reasoning_effort_handle(
444 mut self,
445 handle: crate::tool_context::ReasoningEffortHandle,
446 ) -> Self {
447 self.reasoning_effort_handle = Some(handle);
448 self
449 }
450
451 pub fn with_utility_llm_service(mut self, service: Arc<dyn crate::UtilityLlmService>) -> Self {
454 self.utility_llm_service = Some(service);
455 self
456 }
457
458 pub fn with_decisions(mut self, service: Arc<dyn crate::DecisionsService>) -> Self {
461 self.decisions = Some(service);
462 self
463 }
464}
465
466impl ReasonAtom {
467 pub fn name(&self) -> &'static str {
469 "reason"
470 }
471
472 pub async fn execute(&self, input: ReasonInput) -> Result<ReasonResult> {
474 self.execute_inner(input, None).await
475 }
476}
477
478impl ReasonAtom {
479 pub async fn execute_with_assembled_context(
484 &self,
485 input: ReasonInput,
486 assembled: AssembledTurnContext,
487 ) -> Result<ReasonResult> {
488 self.execute_inner(input, Some(assembled)).await
489 }
490
491 async fn emit_capability_usage_snapshot(
492 &self,
493 session_id: SessionId,
494 context: &ExecutionContext,
495 resolved_capability_configs: &[crate::CapabilityRef],
496 tool_definitions: &[ToolDefinition],
497 ) {
498 let records = capability_usage_snapshot_records(
499 &self.capability_registry,
500 resolved_capability_configs,
501 tool_definitions,
502 );
503 if records.is_empty() {
504 return;
505 }
506
507 if let Err(error) = self
508 .event_emitter
509 .emit(EventRequest::new(
510 session_id,
511 EventContext::from_execution_context(context),
512 CapabilityUsageData { records },
513 ))
514 .await
515 {
516 tracing::warn!(
517 session_id = %session_id,
518 error = %error,
519 "ReasonAtom: failed to emit capability.usage event"
520 );
521 }
522 }
523
524 async fn apply_finalized_tool_call_hooks(
527 &self,
528 session_id: SessionId,
529 context: &ExecutionContext,
530 resolved_capability_configs: &[crate::CapabilityRef],
531 tool_definitions: &[ToolDefinition],
532 tool_calls: &mut [ToolCall],
533 iteration: u32,
534 ) -> Vec<crate::finalized_tool_calls::FinalizedToolCallRejection> {
535 finalized_calls::apply_finalized_tool_calls_hooks(
536 &self.capability_registry,
537 self.event_emitter.as_ref(),
538 session_id,
539 context,
540 resolved_capability_configs,
541 tool_definitions,
542 tool_calls,
543 iteration,
544 )
545 .await
546 }
547
548 async fn execute_inner(
549 &self,
550 input: ReasonInput,
551 assembled: Option<AssembledTurnContext>,
552 ) -> Result<ReasonResult> {
553 let ReasonInput {
554 context,
555 harness_id,
556 agent_id,
557 org_id,
558 mcp_tool_definitions,
559 previous_response_id,
560 iteration,
561 } = input;
562
563 tracing::info!(
564 session_id = %context.session_id,
565 turn_id = %context.turn_id,
566 exec_id = %context.exec_id,
567 harness_id = %harness_id,
568 agent_id = ?agent_id,
569 mcp_tools_count = %mcp_tool_definitions.len(),
570 "ReasonAtom: starting LLM call"
571 );
572
573 let trace_id = context.turn_id.to_string();
581 let reason_span_id = Uuid::now_v7().to_string();
582 let parent_span_id = trace_id.clone(); let event_context = EventContext::from_execution_context(&context).with_span(
586 trace_id.clone(),
587 reason_span_id.clone(),
588 Some(parent_span_id.clone()),
589 );
590
591 let reason_start = Instant::now();
593
594 if let Err(e) = self
596 .event_emitter
597 .emit(EventRequest::new(
598 context.session_id,
599 event_context.clone(),
600 ReasonStartedData {
601 harness_id,
602 agent_id,
603 metadata: None, },
605 ))
606 .await
607 {
608 tracing::warn!(
609 session_id = %context.session_id,
610 error = %e,
611 "ReasonAtom: failed to emit reason.started event"
612 );
613 }
614
615 let assembled = match assembled {
619 Some(assembled) => Ok(assembled),
620 None => {
621 self.context_resolver
622 .resolve_turn_context(TurnContextRequest {
623 session_id: context.session_id,
624 harness_id,
625 agent_id,
626 mcp_tool_definitions: mcp_tool_definitions.clone(),
627 })
628 .await
629 }
630 };
631
632 let (error_disclosure, error_context, error_hooks, call_result) = match assembled {
633 Ok(assembled) => {
634 let error_disclosure = resolve_error_disclosure(
635 &self.capability_registry,
636 &assembled.resolved_capability_configs,
637 error_disclosure_override(&assembled.messages).as_deref(),
638 );
639 let error_hooks =
643 self.collect_llm_error_hooks(&assembled.resolved_capability_configs);
644 let error_context = UserFacingErrorContext::default()
645 .with_provider(assembled.model.provider_type.to_string())
646 .with_model_id(assembled.model.model.clone());
647 let call_result = self
648 .execute_llm_call(
649 context.session_id,
650 harness_id,
651 agent_id,
652 org_id,
653 &context,
654 &trace_id,
655 &reason_span_id,
656 previous_response_id,
657 iteration,
658 assembled,
659 )
660 .await;
661 (error_disclosure, error_context, error_hooks, call_result)
662 }
663 Err(error) => (
664 ErrorDisclosure::default(),
665 UserFacingErrorContext::default(),
666 Vec::new(),
667 Err(error),
668 ),
669 };
670
671 let result = match call_result {
673 Ok(result) => {
674 let reason_duration_ms = reason_start.elapsed().as_millis() as u64;
676
677 let completed_context = EventContext::from_execution_context(&context).with_span(
679 trace_id.clone(),
680 reason_span_id.clone(), Some(parent_span_id.clone()),
682 );
683 if let Err(e) = self
684 .event_emitter
685 .emit(EventRequest::new(
686 context.session_id,
687 completed_context,
688 ReasonCompletedData::success(
689 &result.text,
690 result.has_tool_calls,
691 result.tool_calls.len() as u32,
692 Some(reason_duration_ms),
693 result.usage.clone(),
694 ),
695 ))
696 .await
697 {
698 tracing::warn!(
699 session_id = %context.session_id,
700 error = %e,
701 "ReasonAtom: failed to emit reason.completed event"
702 );
703 }
704 result
705 }
706 Err(e) => {
707 let reason_duration_ms = reason_start.elapsed().as_millis() as u64;
709
710 tracing::warn!(
713 session_id = %context.session_id,
714 turn_id = %context.turn_id,
715 error = %e,
716 "ReasonAtom: LLM call failed"
717 );
718
719 let error_msg = e.to_string();
720 let mut source_error = e.user_facing_error(error_context);
721
722 let is_transient = e.is_transient_llm_error()
729 || (e.llm_error_kind().is_none() && is_transient_error_message(&error_msg));
730
731 if !is_transient && !error_hooks.is_empty() {
737 let services = crate::llm_error_hook::LlmErrorHookServices {
738 schedule_store: self.schedule_store.clone(),
739 };
740 for (hook, config) in &error_hooks {
741 let outcome = {
742 let ctx = crate::llm_error_hook::LlmErrorContext {
743 session_id: context.session_id,
744 error_code: &source_error.code,
745 error_fields: &source_error.fields,
746 config,
747 services: &services,
748 };
749 hook.on_llm_error(&ctx).await
750 };
751 for (key, value) in outcome.extra_error_fields {
752 source_error = source_error.with_field(key, value);
753 }
754 }
755 }
756
757 let user_error = source_error.apply_disclosure(error_disclosure, Some(&error_msg));
758 let user_error_text = user_error.fallback_message();
759
760 let mut output_message_id = None;
761
762 if !is_transient {
763 let mut error_message = RuntimeMessage::assistant(&user_error_text);
765 let mut metadata = std::collections::HashMap::new();
766 user_error.apply_to_message_metadata(&mut metadata);
767 UserFacingError::apply_disclosure_to_message_metadata(
768 &mut metadata,
769 error_disclosure,
770 &source_error.code,
771 );
772 error_message.metadata = Some(metadata);
773
774 output_message_id = Some(error_message.id);
775
776 let error_msg_context = EventContext::from_execution_context(&context)
779 .with_span(
780 trace_id.clone(),
781 Uuid::now_v7().to_string(), Some(reason_span_id.clone()), );
784 if let Err(emit_err) = self
785 .event_emitter
786 .emit(EventRequest::new(
787 context.session_id,
788 error_msg_context,
789 OutputMessageCompletedData::new(error_message)
790 .with_user_facing_error(&user_error)
791 .with_error_disclosure(error_disclosure),
792 ))
793 .await
794 {
795 tracing::warn!(
796 session_id = %context.session_id,
797 error = %emit_err,
798 "ReasonAtom: failed to emit output.message.completed event for error"
799 );
800 }
801 } else {
802 tracing::info!(
803 session_id = %context.session_id,
804 "ReasonAtom: skipping error event for transient LLM error (will be retried)"
805 );
806 }
807
808 let completed_context = EventContext::from_execution_context(&context).with_span(
810 trace_id.clone(),
811 reason_span_id.clone(), Some(parent_span_id.clone()),
813 );
814 if let Err(emit_err) = self
815 .event_emitter
816 .emit(EventRequest::new(
817 context.session_id,
818 completed_context,
819 ReasonCompletedData::failure(error_msg.clone(), Some(reason_duration_ms)),
820 ))
821 .await
822 {
823 tracing::warn!(
824 session_id = %context.session_id,
825 error = %emit_err,
826 "ReasonAtom: failed to emit reason.completed event"
827 );
828 }
829
830 ReasonResult {
831 native_counts: None,
832 success: false,
833 text: user_error_text,
834 tool_calls: vec![],
835 has_tool_calls: false,
836 tool_definitions: vec![],
837 max_iterations: default_max_iterations(),
838 error: Some(error_msg.clone()),
839 user_facing_error: Some(user_error),
840 error_disclosure: Some(error_disclosure),
841 usage: None,
842 output_message_id,
843 time_to_first_token_ms: None,
844 response_id: None,
845 finish_reason: error_msg
846 .to_ascii_lowercase()
847 .contains("model refused")
848 .then(|| "refusal".to_string()),
849 locale: None,
850 network_access: None,
851 parallel_tool_calls: None,
852 }
853 }
854 };
855
856 Ok(result)
857 }
858
859 #[allow(clippy::too_many_arguments)]
861 async fn execute_llm_call(
862 &self,
863 session_id: SessionId,
864 harness_id: HarnessId,
865 agent_id: Option<AgentId>,
866 org_id: i64,
867 context: &ExecutionContext,
868 trace_id: &str,
869 reason_span_id: &str,
870 previous_response_id: Option<String>,
871 iteration: u32,
872 assembled: AssembledTurnContext,
873 ) -> Result<ReasonResult> {
874 let prior_usage = assembled.cumulative_usage();
875 let mut messages = transcript::order_native_results(assembled.messages);
876 let mut message_source_sequence = assembled.message_source_sequence;
877 let model_with_provider = assembled.model;
878 let resolved_model_id = assembled.resolved_model_id;
879 let resolved_locale = assembled.resolved_locale;
880 let compaction_policy = assembled.compaction_policy;
881 let resolved_capability_configs = assembled.resolved_capability_configs;
882 let runtime_agent = assembled.runtime_agent;
883 let embedder_metadata = assembled.embedder_metadata;
884
885 self.emit_capability_usage_snapshot(
886 session_id,
887 context,
888 &resolved_capability_configs,
889 &runtime_agent.tools,
890 )
891 .await;
892
893 let output_hooks =
894 collect_output_hooks(&self.capability_registry, &resolved_capability_configs);
895 let guardrail_providers = output_hooks.streaming;
896 let post_output_providers = output_hooks.post_generation;
897 let annotation_providers = output_hooks.annotations;
898 let citation_verifiers = output_hooks.citation_verifiers;
899
900 let chat_driver = Arc::clone(&model_with_provider.driver);
902 let stateful_response_continuation =
903 previous_response_id.is_some() && chat_driver.supports_stateful_responses();
904 let mut restored_checkpoint: Option<crate::CompactionCheckpoint> = None;
905 let mut checkpoint_suffix_message_count = 0usize;
906 let native_reasoning_compaction = compaction_policy.as_ref().is_none_or(|policy| {
907 matches!(
908 policy.settings().strategy,
909 crate::compaction_policy::CompactionStrategy::Native
910 | crate::compaction_policy::CompactionStrategy::Auto
911 ) && chat_driver.supports_compact()
912 });
913
914 if compaction_policy.is_some()
915 && let Some(store) = self.compaction_checkpoint_store.as_ref()
916 && let Some(checkpoint) = store
917 .get_latest(
918 session_id,
919 model_with_provider.provider_type.as_str(),
920 &model_with_provider.model,
921 )
922 .await?
923 && checkpoint.is_compatible(
924 model_with_provider.provider_type.as_str(),
925 &model_with_provider.model,
926 )
927 && (native_reasoning_compaction || !matches!(
930 &checkpoint.payload,
931 crate::CompactionCheckpointPayload::ProviderOpaque {
932 context: crate::ProviderOpaqueContext::OpenResponsesCompact {
933 reasoning_state: Some(_), ..
934 }
935 }
936 ))
937 {
938 let filters = crate::capabilities::collect_message_filters_only(
939 &resolved_capability_configs,
940 &self.capability_registry,
941 );
942 let mut query =
943 crate::MessageQuery::new(session_id).after_sequence(checkpoint.source_sequence);
944 filters.apply_message_filters(&mut query);
945 let history = self.message_retriever.load_filtered_history(query).await?;
946 messages = history.messages;
947 checkpoint_suffix_message_count = messages.len();
948 filters.apply_post_load_filters(&mut messages);
949 if let crate::CompactionCheckpointPayload::Summary { text } = &checkpoint.payload {
950 messages.insert(
951 0,
952 RuntimeMessage::system(format!(
953 "[CONVERSATION_SUMMARY]\n{text}\n[/CONVERSATION_SUMMARY]"
954 )),
955 );
956 }
957 message_source_sequence = history.source_sequence.or(message_source_sequence);
958 restored_checkpoint = Some(checkpoint);
959 }
960
961 let controls = resolve_request_controls(
962 &messages,
963 self.reasoning_effort_handle.as_ref(),
964 &model_with_provider.provider_type,
965 &model_with_provider.model,
966 );
967 let reasoning_effort = controls.reasoning_effort;
968 let speed = controls.speed;
969 let verbosity = controls.verbosity;
970 let checkpoint_reasoning =
971 restored_checkpoint
972 .as_ref()
973 .and_then(|checkpoint| match &checkpoint.payload {
974 crate::CompactionCheckpointPayload::ProviderOpaque {
975 context:
976 crate::ProviderOpaqueContext::OpenResponsesCompact {
977 reasoning_state, ..
978 },
979 } => reasoning_state.as_ref(),
980 _ => None,
981 });
982 let mut reasoning_replay = reasoning_updates::prepare(
983 &messages,
984 model_with_provider.provider_type.as_str(),
985 &model_with_provider.model,
986 reasoning_effort,
987 self.reasoning_effort_handle
988 .as_ref()
989 .and_then(crate::tool_context::ReasoningEffortHandle::get),
990 checkpoint_reasoning,
991 )
992 .filter(|_| native_reasoning_compaction);
993
994 if let Some(ref store) = self.partial_stream_store {
998 let turn_id_str = context.turn_id.to_string();
999 match store.get_partial_stream(session_id, &turn_id_str).await {
1000 Ok(Some(partial)) if !partial.accumulated.is_empty() => {
1001 return self
1003 .finalize_partial_stream(
1004 session_id,
1005 context,
1006 partial,
1007 iteration,
1008 &runtime_agent,
1009 &resolved_capability_configs,
1010 )
1011 .await;
1012 }
1013 Ok(Some(partial)) => {
1014 if let (Some(replay), Some(mut saved)) =
1015 (reasoning_replay.as_mut(), partial.reasoning_state)
1016 {
1017 saved.pending = saved.effective;
1020 replay.state = saved;
1021 }
1022 let recovery_ctx = EventContext::from_execution_context(context);
1025 let _ = self
1026 .event_emitter
1027 .emit(EventRequest::new(
1028 session_id,
1029 recovery_ctx,
1030 ReasonRecoveredData {
1031 turn_id: context.turn_id,
1032 mode: RecoveryMode::Restart,
1033 accumulated_len: 0,
1034 },
1035 ))
1036 .await;
1037 tracing::info!(
1038 session_id = %session_id,
1039 turn_id = %context.turn_id,
1040 "ReasonAtom: partial stream detected with empty accumulated; restarting clean"
1041 );
1042 }
1043 Ok(None) => {} Err(e) => {
1045 if reasoning_replay.is_some() {
1046 return Err(e);
1047 }
1048 tracing::warn!(
1050 session_id = %session_id,
1051 turn_id = %context.turn_id,
1052 error = %e,
1053 "ReasonAtom: partial-stream store error; proceeding with normal execution"
1054 );
1055 }
1056 }
1057 }
1058
1059 let repair_event_context = EventContext::from_execution_context(context);
1063 let patched_messages = if self.native_async.is_some() {
1064 messages.clone()
1065 } else {
1066 repair_dangling_tool_calls(
1067 &messages,
1068 self.durable_tool_result_store.as_deref(),
1069 self.event_emitter.as_ref(),
1070 session_id,
1071 &repair_event_context,
1072 &context.turn_id.to_string(),
1073 )
1074 .await
1075 };
1076 let raw_tool_result_bytes = compaction_policy
1077 .as_ref()
1078 .map(|policy| policy.total_tool_result_bytes(&patched_messages))
1079 .unwrap_or(0);
1080
1081 let model_view_providers = crate::capabilities::collect_model_view_providers(
1084 &resolved_capability_configs,
1085 &self.capability_registry,
1086 Some(model_with_provider.model.as_str()),
1087 );
1088 let model_view_context = crate::capabilities::ModelViewContext {
1089 session_id,
1090 prior_usage: prior_usage.as_ref(),
1091 };
1092 let mut context_messages =
1093 model_view_providers.apply_model_view(patched_messages, &model_view_context);
1094 context_messages = crate::tool_call_integrity::retain_complete_message_tool_exchanges(
1095 &context_messages,
1096 stateful_response_continuation || restored_checkpoint.is_some(),
1097 );
1098
1099 let mut volatile_suffix_len = 0usize;
1106 {
1107 let facts_ctx = crate::capabilities::FactsContext::new(session_id);
1108 let dynamic_facts = crate::capabilities::collect_dynamic_facts(
1109 &resolved_capability_configs,
1110 &self.capability_registry,
1111 Some(model_with_provider.model.as_str()),
1112 &facts_ctx,
1113 );
1114 if let Some(block) = crate::capabilities::render_facts_block(&dynamic_facts) {
1115 context_messages.push(RuntimeMessage::user(block));
1116 volatile_suffix_len = 1;
1117 }
1118 }
1119
1120 if let Some(context) = runtime_agent.conversation_context.as_ref()
1129 && !context.is_empty()
1130 {
1131 context_messages.insert(0, RuntimeMessage::user(context.clone()));
1132 }
1133
1134 let resolved_images = self.resolve_images(&context_messages).await;
1139 let resolved_files = self.resolve_files(&context_messages).await;
1140
1141 let mut llm_messages = Vec::new();
1143
1144 let has_system_prompt = !runtime_agent.system_prompt.is_empty();
1146 if has_system_prompt {
1147 llm_messages.push(Message {
1148 native_tool_calls: Vec::new(),
1149 role: MessageRole::System,
1150 content: MessageContent::Text(runtime_agent.system_prompt.clone()),
1151 tool_calls: None,
1152 tool_call_id: None,
1153 phase: None,
1154 reasoning: Vec::new(),
1155 configuration_update: None,
1156 });
1157 }
1158
1159 let messages_for_event: Vec<RuntimeMessage> = if has_system_prompt {
1161 std::iter::once(RuntimeMessage::system(&runtime_agent.system_prompt))
1162 .chain(context_messages.iter().cloned())
1163 .collect()
1164 } else {
1165 context_messages.clone()
1166 };
1167
1168 let mut stripped_error_count = 0u32;
1174 for msg in &context_messages {
1175 if is_error_placeholder_message(msg) {
1176 stripped_error_count += 1;
1177 continue;
1178 }
1179 let mut llm_msg = crate::llm_conversions::llm_message_from_message_with_attachments(
1180 msg,
1181 &resolved_images,
1182 &resolved_files,
1183 );
1184 llm_msg.configuration_update = reasoning_replay
1185 .as_ref()
1186 .and_then(|replay| replay.transitions.get(&msg.id).copied());
1187 if msg.role == RuntimeMessageRole::User
1188 && let Some(ref actor) = msg.external_actor
1189 {
1190 llm_msg.prepend_text_prefix(&format!("[{}] ", actor.display_label()));
1191 }
1192 llm_messages.push(llm_msg);
1193 }
1194 if stripped_error_count > 0 {
1195 tracing::info!(
1196 session_id = %session_id,
1197 stripped_error_count,
1198 "ReasonAtom: stripped error placeholder messages from LLM input"
1199 );
1200 }
1201
1202 llm_messages = crate::tool_call_integrity::retain_complete_llm_tool_exchanges_for_request(
1207 llm_messages,
1208 stateful_response_continuation || restored_checkpoint.is_some(),
1209 );
1210
1211 let mut llm_config_builder =
1213 crate::llm_conversions::llm_call_config_builder_from_agent(&runtime_agent);
1214 if let Some(effort) = reasoning_effort {
1215 llm_config_builder = llm_config_builder.reasoning_effort(effort);
1216 }
1217 if let Some(speed) = speed {
1218 llm_config_builder = llm_config_builder.speed(speed);
1219 }
1220 if let Some(verbosity) = verbosity {
1221 llm_config_builder = llm_config_builder.verbosity(verbosity);
1222 }
1223
1224 for (k, v) in &embedder_metadata {
1226 llm_config_builder = llm_config_builder.with_metadata(k, v.clone());
1227 }
1228
1229 llm_config_builder = llm_config_builder
1233 .with_metadata("session_id", session_id.to_string())
1234 .with_metadata("harness_id", harness_id.to_string())
1235 .with_metadata("turn_id", context.turn_id.to_string())
1236 .with_metadata("exec_id", context.exec_id.to_string())
1237 .with_metadata("org_id", format!("org_{:032x}", org_id));
1238 if let Some(agent_id) = agent_id {
1239 llm_config_builder = llm_config_builder.with_metadata("agent_id", agent_id.to_string());
1240 }
1241
1242 if let Some(model_id) = &resolved_model_id {
1244 llm_config_builder = llm_config_builder.with_metadata("model_id", model_id.to_string());
1245 }
1246
1247 let mut llm_config = llm_config_builder
1248 .previous_response_id(previous_response_id.clone())
1249 .volatile_suffix_len(volatile_suffix_len)
1250 .build();
1251 if let Some(replay) = &reasoning_replay {
1252 llm_config.reasoning_effort = replay.state.baseline;
1253 llm_config.reasoning_state = Some(replay.state.clone());
1254 if replay.reset_continuation {
1255 llm_config.previous_response_id = None;
1256 }
1257 } else if messages
1258 .iter()
1259 .rev()
1260 .find(|message| {
1261 message.role == RuntimeMessageRole::Agent && !is_error_placeholder_message(message)
1262 })
1263 .and_then(|message| message.metadata.as_ref())
1264 .is_some_and(|metadata| metadata.contains_key(reasoning_updates::STATE_KEY))
1265 {
1266 llm_config.previous_response_id = None;
1269 }
1270 if let Some(checkpoint) = restored_checkpoint.as_ref()
1271 && let crate::CompactionCheckpointPayload::ProviderOpaque { context } =
1272 &checkpoint.payload
1273 {
1274 llm_config.previous_response_id = None;
1275 llm_config.provider_opaque_context = Some(context.clone());
1276 }
1277
1278 tracing::debug!(
1279 session_id = %session_id,
1280 turn_id = %context.turn_id,
1281 model = %runtime_agent.model,
1282 message_count = %llm_messages.len(),
1283 "ReasonAtom: calling LLM"
1284 );
1285
1286 let streaming_event_context = EventContext::from_execution_context(context);
1289
1290 let mut armed_guardrails: Vec<ArmedGuardrail> = Vec::new();
1297 for (cap_id, cfg, provider) in &guardrail_providers {
1298 let ctx = OutputGuardrailContext {
1299 system_prompt: &runtime_agent.system_prompt,
1300 config: cfg,
1301 };
1302 let guardrail_id = provider.id().to_string();
1303 if let Some(run) = provider.arm(&ctx) {
1304 armed_guardrails.push(ArmedGuardrail {
1305 capability_id: cap_id.clone(),
1306 guardrail_id,
1307 run,
1308 });
1309 }
1310 }
1311 let buffer_output_deltas = !post_output_providers.is_empty();
1316 let output_message_id = MessageId::new();
1320 tracing::info!(
1321 session_id = %session_id,
1322 turn_id = %context.turn_id,
1323 "ReasonAtom: emitting output.message.started event"
1324 );
1325 if let Err(e) = self
1326 .event_emitter
1327 .emit(EventRequest::new(
1328 session_id,
1329 streaming_event_context.clone(),
1330 OutputMessageStartedData {
1331 reasoning_state: llm_config.reasoning_state.clone(),
1332 turn_id: context.turn_id,
1333 message_id: output_message_id,
1334 model: Some(runtime_agent.model.clone()),
1335 iteration: Some(iteration),
1336 phase: None,
1339 },
1340 ))
1341 .await
1342 {
1343 if llm_config.reasoning_state.is_some() {
1344 return Err(e);
1345 }
1346 tracing::warn!(
1347 session_id = %session_id,
1348 error = %e,
1349 "ReasonAtom: failed to emit output.message.started event"
1350 );
1351 } else {
1352 tracing::info!(
1353 session_id = %session_id,
1354 "ReasonAtom: output.message.started event emitted successfully"
1355 );
1356 }
1357
1358 let thinking_enabled = reasoning_effort.is_some();
1360 if thinking_enabled {
1361 tracing::info!(
1362 session_id = %session_id,
1363 turn_id = %context.turn_id,
1364 "ReasonAtom: emitting reason.thinking.started event"
1365 );
1366 if let Err(e) = self
1367 .event_emitter
1368 .emit(EventRequest::new(
1369 session_id,
1370 streaming_event_context.clone(),
1371 ReasonThinkingStartedData {
1372 turn_id: context.turn_id,
1373 model: Some(runtime_agent.model.clone()),
1374 },
1375 ))
1376 .await
1377 {
1378 tracing::warn!(
1379 session_id = %session_id,
1380 error = %e,
1381 "ReasonAtom: failed to emit reason.thinking.started event"
1382 );
1383 } else {
1384 tracing::info!(
1385 session_id = %session_id,
1386 "ReasonAtom: reason.thinking.started event emitted successfully"
1387 );
1388 }
1389 }
1390
1391 let llm_start = Instant::now();
1393
1394 let mut compaction_info: Option<LlmCompactionInfo> = None;
1398 let mut llm_messages_for_call = llm_messages.clone();
1399
1400 if let Some(policy) = compaction_policy.as_deref() {
1401 compaction_info = apply_proactive_compaction(
1402 ProactiveCompactionContext {
1403 chat_driver: chat_driver.as_ref(),
1404 policy,
1405 checkpoint_store: self.compaction_checkpoint_store.as_ref(),
1406 event_emitter: self.event_emitter.as_ref(),
1407 event_context: &streaming_event_context,
1408 session_id,
1409 message_source_sequence,
1410 provider_type: model_with_provider.provider_type.as_str(),
1411 model: &model_with_provider.model,
1412 system_prompt: has_system_prompt
1413 .then_some(runtime_agent.system_prompt.as_str()),
1414 stateful_response_continuation,
1415 checkpoint_restored: restored_checkpoint.is_some(),
1416 checkpoint_suffix_message_count,
1417 raw_tool_result_bytes,
1418 prior_usage: prior_usage.as_ref(),
1419 },
1420 &mut llm_messages_for_call,
1421 &mut llm_config,
1422 )
1423 .await?;
1424 }
1425
1426 const DELTA_BATCH_INTERVAL_MS: u64 = 100;
1429 let retry_config = self.provider_retry_config.clone();
1430 let has_provider_executed_tools = llm_config
1438 .driver_options
1439 .get("openrouter/routing")
1440 .and_then(|raw| raw.get("server_tools"))
1441 .and_then(|tools| tools.as_array())
1442 .is_some_and(|tools| !tools.is_empty());
1443 let mut stream_retry_metadata = RetryMetadata::default();
1444 let mut retry_started_at = None;
1445 let mut streamed_phase: Option<everruns_provider::ExecutionPhase> = None;
1450 let mut native_calls = std::collections::BTreeMap::new();
1451 let (
1452 text,
1453 thinking,
1454 reasoning,
1455 tool_calls,
1456 completion_metadata,
1457 time_to_first_token_ms,
1458 pending_delta,
1459 mut tripped,
1460 ) = 'stream_attempt: loop {
1461 let stream_result = if let Some(remaining) =
1462 remaining_retry_time(&retry_config, retry_started_at)
1463 {
1464 match tokio::time::timeout(
1465 remaining,
1466 chat_driver.chat_completion_stream(
1467 &crate::ProviderEndpoint::default(),
1468 llm_messages_for_call.clone(),
1469 &llm_config,
1470 ),
1471 )
1472 .await
1473 {
1474 Ok(result) => result,
1475 Err(_) => {
1476 return Err(AgentLoopError::llm_kind(
1477 crate::error::LlmErrorKind::Unavailable,
1478 format!(
1479 "provider retry time budget exhausted after {} retries over {:.1}s; the turn is safe to resume",
1480 stream_retry_metadata.attempts,
1481 retry_config.max_retry_elapsed.as_secs_f64()
1482 ),
1483 )
1484 .with_retry_metadata(&stream_retry_metadata));
1485 }
1486 }
1487 } else {
1488 chat_driver
1489 .chat_completion_stream(
1490 &crate::ProviderEndpoint::default(),
1491 llm_messages_for_call.clone(),
1492 &llm_config,
1493 )
1494 .await
1495 };
1496 let mut stream = match stream_result {
1497 Ok(stream) => stream,
1498 Err(e) if e.is_request_too_large() => {
1499 let Some(policy) = compaction_policy.as_deref() else {
1500 tracing::warn!(
1501 session_id = %session_id,
1502 turn_id = %context.turn_id,
1503 "ReasonAtom: context too large and compaction capability is not enabled"
1504 );
1505 return Err(e);
1506 };
1507 let outcome = apply_reactive_compaction(
1508 ReactiveCompactionContext {
1509 chat_driver: chat_driver.as_ref(),
1510 policy,
1511 checkpoint_store: self.compaction_checkpoint_store.as_ref(),
1512 event_emitter: self.event_emitter.as_ref(),
1513 event_context: &streaming_event_context,
1514 session_id,
1515 message_source_sequence,
1516 provider_type: model_with_provider.provider_type.as_str(),
1517 model: &model_with_provider.model,
1518 summarization_model_fallback: &runtime_agent.model,
1519 system_prompt: has_system_prompt
1520 .then_some(runtime_agent.system_prompt.as_str()),
1521 stateful_response_continuation,
1522 },
1523 &mut llm_messages_for_call,
1524 &mut llm_config,
1525 )
1526 .await?;
1527 let Some(outcome) = outcome else {
1528 return Err(e);
1529 };
1530 if outcome.generation_info.is_some() {
1531 compaction_info = outcome.generation_info;
1532 }
1533
1534 chat_driver
1535 .chat_completion_stream(
1536 &crate::ProviderEndpoint::default(),
1537 llm_messages_for_call.clone(),
1538 &llm_config,
1539 )
1540 .await?
1541 }
1542 Err(e)
1543 if e.is_transient_llm_error()
1544 && !e.llm_retry_handled()
1545 && !has_provider_executed_tools
1546 && stream_retry_metadata.attempts < retry_config.max_retries =>
1547 {
1548 let proposed_wait =
1549 retry_config.calculate_backoff(stream_retry_metadata.attempts);
1550 let Some(wait_duration) =
1551 reserve_retry_wait(&retry_config, &mut retry_started_at, proposed_wait)
1552 else {
1553 return Err(AgentLoopError::llm_kind(
1554 e.llm_error_kind()
1555 .unwrap_or(crate::error::LlmErrorKind::Unavailable),
1556 format!(
1557 "{e}; automatic recovery time budget exhausted after {} retries; the turn is safe to resume",
1558 stream_retry_metadata.attempts
1559 ),
1560 )
1561 .with_retry_metadata(&stream_retry_metadata));
1562 };
1563 tracing::warn!(
1564 session_id = %session_id,
1565 turn_id = %context.turn_id,
1566 attempt = stream_retry_metadata.attempts + 1,
1567 max_retries = retry_config.max_retries,
1568 wait_secs = wait_duration.as_secs_f64(),
1569 error = %e,
1570 "ReasonAtom: transient provider failure before stream, retrying"
1571 );
1572 stream_retry_metadata.record_retry(wait_duration, None);
1573 tokio::time::sleep(wait_duration).await;
1574 continue 'stream_attempt;
1575 }
1576 Err(e) => return Err(e),
1577 };
1578
1579 if let Some(coordinator) = &self.native_async {
1580 coordinator
1581 .lock()
1582 .await
1583 .begin_transcript_response(output_message_id.to_string())
1584 .await?;
1585 let coordinator = coordinator.clone();
1586 stream = Box::pin(futures::stream::unfold(
1587 Some((coordinator, stream)),
1588 |state| async move {
1589 let (coordinator, mut source) = state?;
1590 let event = coordinator
1591 .lock()
1592 .await
1593 .next_response_event(&mut source)
1594 .await;
1595 let finished = matches!(&event, Ok(LlmStreamEvent::Done(_)) | Err(_));
1596 Some((event, (!finished).then_some((coordinator, source))))
1597 },
1598 ));
1599 }
1600 let mut text = String::new();
1601 let mut reasoning: Vec<ReasoningContentPart> = Vec::new();
1605 let mut thinking = String::new();
1607 let mut tool_calls = Vec::new();
1608 let mut termination = StreamTermination::Exhausted;
1609 let mut replay_state = StreamReplayState::for_request(has_provider_executed_tools);
1610 let mut pending_delta = String::new();
1611 let mut pending_thinking_delta = String::new();
1612 let mut last_delta_emit = Instant::now();
1613 let mut last_thinking_delta_emit = Instant::now();
1614 let mut time_to_first_token_ms: Option<u64> = None;
1615
1616 let stall_timeout = self
1618 .provider_stall_timeout
1619 .unwrap_or(std::time::Duration::from_secs(120));
1620 let initial_stall_timeout = remaining_retry_time(&retry_config, retry_started_at)
1621 .map_or(stall_timeout, |remaining| remaining.min(stall_timeout));
1622 let mut stall_sleep = Box::pin(tokio::time::sleep(initial_stall_timeout));
1623 let mut keepalive_ticker = tokio::time::interval(std::time::Duration::from_secs(12));
1624 keepalive_ticker.set_missed_tick_behavior(tokio::time::MissedTickBehavior::Skip);
1625 keepalive_ticker.tick().await; let mut last_stream_heartbeat = Instant::now();
1627 let mut last_token_at_unix: u64 = unix_now_secs();
1632
1633 loop {
1634 let event = tokio::select! {
1635 biased;
1636 next = stream.next() => match next {
1637 Some(e) => e,
1638 None => break,
1639 },
1640 _ = &mut stall_sleep => {
1641 let stall_error =
1651 crate::driver_registry::LlmStreamError::new(format!(
1652 "provider stream stall: no tokens for {}s",
1653 stall_timeout.as_secs()
1654 ));
1655 tracing::warn!(
1656 session_id = %session_id,
1657 turn_id = %context.turn_id,
1658 stall_secs = stall_timeout.as_secs(),
1659 "ReasonAtom: provider stream stall timeout"
1660 );
1661 if replay_state.should_retry(
1662 &stall_error,
1663 stream_retry_metadata.attempts,
1664 retry_config.max_retries,
1665 ) {
1666 let proposed_wait = retry_config
1667 .calculate_backoff(stream_retry_metadata.attempts);
1668 let Some(wait_duration) = reserve_retry_wait(
1669 &retry_config,
1670 &mut retry_started_at,
1671 proposed_wait,
1672 ) else {
1673 return Err(AgentLoopError::llm_kind(
1674 crate::error::LlmErrorKind::Unavailable,
1675 format!(
1676 "{}; automatic recovery time budget exhausted after {} retries; the turn is safe to resume",
1677 stall_error.message,
1678 stream_retry_metadata.attempts
1679 ),
1680 )
1681 .with_retry_metadata(&stream_retry_metadata));
1682 };
1683 tracing::warn!(
1684 session_id = %session_id,
1685 turn_id = %context.turn_id,
1686 attempt = stream_retry_metadata.attempts + 1,
1687 max_retries = retry_config.max_retries,
1688 wait_secs = wait_duration.as_secs_f64(),
1689 "ReasonAtom: provider stream stall, retrying"
1690 );
1691 stream_retry_metadata.record_retry(wait_duration, None);
1692 tokio::time::sleep(wait_duration).await;
1693 continue 'stream_attempt;
1694 }
1695 return Err(AgentLoopError::llm(stall_error.message));
1696 },
1697 _ = keepalive_ticker.tick() => {
1698 if let Some(ref hb) = self.stream_heartbeater {
1699 hb.heartbeat(crate::durability::StreamProgress {
1700 accumulated_len: text.len() + thinking.len(),
1701 last_delta_at: last_token_at_unix,
1702 })
1703 .await;
1704 last_stream_heartbeat = Instant::now();
1705 }
1706 continue;
1707 },
1708 };
1709 let event = event?;
1710 replay_state.observe(&event);
1711 let advanced_stall_deadline = advances_stall_deadline(&event);
1712 if advanced_stall_deadline {
1713 stall_sleep
1714 .as_mut()
1715 .reset(tokio::time::Instant::now() + stall_timeout);
1716 last_token_at_unix = unix_now_secs();
1717 }
1718 match event {
1719 LlmStreamEvent::TextDelta(delta) => {
1720 if delta.is_empty() {
1721 continue;
1722 }
1723 if time_to_first_token_ms.is_none() {
1725 let ttft = llm_start.elapsed().as_millis() as u64;
1726 time_to_first_token_ms = Some(ttft);
1727 tracing::info!(
1728 session_id = %session_id,
1729 time_to_first_token_ms = ttft,
1730 "ReasonAtom: received first token from LLM"
1731 );
1732 }
1733 text.push_str(&delta);
1734 pending_delta.push_str(&delta);
1735
1736 if !armed_guardrails.is_empty()
1743 && let Some(t) =
1744 evaluate_guardrails(&mut armed_guardrails, &text, &delta)
1745 {
1746 tracing::warn!(
1747 session_id = %session_id,
1748 turn_id = %context.turn_id,
1749 guardrail_capability_id = %t.capability_id,
1750 guardrail_id = %t.guardrail_id,
1751 reason_code = %t.block.reason_code,
1752 "ReasonAtom: output guardrail tripped, replacing assistant message"
1753 );
1754 pending_delta.clear();
1755 termination = StreamTermination::GuardrailBlocked(t);
1756 break;
1757 }
1758
1759 if !buffer_output_deltas
1761 && last_delta_emit.elapsed().as_millis() as u64
1762 >= DELTA_BATCH_INTERVAL_MS
1763 && !pending_delta.is_empty()
1764 {
1765 if let Err(e) = self
1766 .event_emitter
1767 .emit(EventRequest::new(
1768 session_id,
1769 streaming_event_context.clone(),
1770 OutputMessageDeltaData {
1771 turn_id: context.turn_id,
1772 message_id: output_message_id,
1773 delta: pending_delta.clone(),
1774 accumulated: text.clone(),
1775 phase: streamed_phase,
1776 },
1777 ))
1778 .await
1779 {
1780 tracing::warn!(
1781 session_id = %session_id,
1782 error = %e,
1783 "ReasonAtom: failed to emit output.message.delta event"
1784 );
1785 }
1786 pending_delta.clear();
1787 last_delta_emit = Instant::now();
1788 }
1789 }
1790 LlmStreamEvent::ReasoningDelta { delta, summary: _ } => {
1791 if delta.is_empty() {
1792 continue;
1793 }
1794 if let Some(t) = append_guarded_thinking_delta(
1795 &mut armed_guardrails,
1796 &mut thinking,
1797 &mut pending_thinking_delta,
1798 &delta,
1799 ) {
1800 tracing::warn!(
1801 session_id = %session_id,
1802 guardrail_capability_id = %t.capability_id,
1803 guardrail_id = %t.guardrail_id,
1804 "ReasonAtom: output guardrail tripped on thinking stream, replacing assistant message"
1805 );
1806 termination = StreamTermination::GuardrailBlocked(t);
1807 break;
1808 }
1809 tracing::debug!(
1810 session_id = %session_id,
1811 delta_len = delta.len(),
1812 total_thinking_len = thinking.len(),
1813 "ReasonAtom: received ThinkingDelta from LLM"
1814 );
1815
1816 if last_thinking_delta_emit.elapsed().as_millis() as u64
1818 >= DELTA_BATCH_INTERVAL_MS
1819 && !pending_thinking_delta.is_empty()
1820 {
1821 if let Err(e) = self
1822 .event_emitter
1823 .emit(EventRequest::new(
1824 session_id,
1825 streaming_event_context.clone(),
1826 ReasonThinkingDeltaData {
1827 turn_id: context.turn_id,
1828 delta: pending_thinking_delta.clone(),
1829 accumulated: thinking.clone(),
1830 },
1831 ))
1832 .await
1833 {
1834 tracing::warn!(
1835 session_id = %session_id,
1836 error = %e,
1837 "ReasonAtom: failed to emit reason.thinking.delta event"
1838 );
1839 }
1840 pending_thinking_delta.clear();
1841 last_thinking_delta_emit = Instant::now();
1842 }
1843 }
1844 LlmStreamEvent::ReasoningItem(item) => {
1845 if let Some(t) = inspect_guarded_reasoning_item(
1846 &mut armed_guardrails,
1847 &mut thinking,
1848 &item,
1849 ) {
1850 tracing::warn!(
1851 session_id = %session_id,
1852 guardrail_capability_id = %t.capability_id,
1853 guardrail_id = %t.guardrail_id,
1854 "ReasonAtom: output guardrail tripped on completed reasoning item, replacing assistant message"
1855 );
1856 termination = StreamTermination::GuardrailBlocked(t);
1857 break;
1858 }
1859 tracing::debug!(
1863 session_id = %session_id,
1864 provider = %item.provider,
1865 item_id = ?item.item_id,
1866 has_signature = item.signature.is_some(),
1867 has_encrypted = item.encrypted.is_some(),
1868 "ReasonAtom: captured reasoning artifact"
1869 );
1870 reasoning.push(item);
1871 }
1872 LlmStreamEvent::NativeToolCall(call) => {
1873 if self.native_async.is_none() {
1874 return Err(AgentLoopError::config(
1875 "native async/custom tools require a configured native-call coordinator",
1876 ));
1877 }
1878 let part = crate::message::ToolCallContentPart::from_native(call.clone())?;
1879 if native_calls.insert(call.id().to_owned(), call).is_none() {
1880 tool_calls.push(ToolCall {
1881 id: part.id,
1882 name: part.name,
1883 arguments: part.arguments,
1884 });
1885 }
1886 }
1887 LlmStreamEvent::ToolCalls(calls) => {
1888 if self.native_async.is_some() {
1889 for call in &calls {
1890 native_calls.entry(call.id.clone()).or_insert_with(|| {
1891 everruns_provider::native_async::NativeToolCall::Function {
1892 call_id: call.id.clone(),
1893 name: call.name.clone(),
1894 arguments: call.arguments.to_string(),
1895 asynchronous: false,
1896 }
1897 });
1898 }
1899 for call in calls {
1900 if !tool_calls.iter().any(|existing| existing.id == call.id) {
1901 tool_calls.push(call);
1902 }
1903 }
1904 } else {
1905 tool_calls = calls;
1906 }
1907 }
1908 LlmStreamEvent::MessagePhase(phase) => {
1909 streamed_phase = everruns_provider::ExecutionPhase::refine_streamed_hint(
1918 streamed_phase,
1919 phase,
1920 );
1921 }
1922 LlmStreamEvent::Done(metadata) => {
1923 if !buffer_output_deltas
1927 && !pending_delta.is_empty()
1928 && let Err(e) = self
1929 .event_emitter
1930 .emit(EventRequest::new(
1931 session_id,
1932 streaming_event_context.clone(),
1933 OutputMessageDeltaData {
1934 turn_id: context.turn_id,
1935 message_id: output_message_id,
1936 delta: pending_delta.clone(),
1937 accumulated: text.clone(),
1938 phase: streamed_phase,
1939 },
1940 ))
1941 .await
1942 {
1943 tracing::warn!(
1944 session_id = %session_id,
1945 error = %e,
1946 "ReasonAtom: failed to emit final output.message.delta event"
1947 );
1948 }
1949
1950 if !pending_thinking_delta.is_empty()
1952 && let Err(e) = self
1953 .event_emitter
1954 .emit(EventRequest::new(
1955 session_id,
1956 streaming_event_context.clone(),
1957 ReasonThinkingDeltaData {
1958 turn_id: context.turn_id,
1959 delta: pending_thinking_delta.clone(),
1960 accumulated: thinking.clone(),
1961 },
1962 ))
1963 .await
1964 {
1965 tracing::warn!(
1966 session_id = %session_id,
1967 error = %e,
1968 "ReasonAtom: failed to emit final reason.thinking.delta event"
1969 );
1970 }
1971
1972 if !thinking.is_empty()
1974 && let Err(e) = self
1975 .event_emitter
1976 .emit(EventRequest::new(
1977 session_id,
1978 streaming_event_context.clone(),
1979 ReasonThinkingCompletedData {
1980 turn_id: context.turn_id,
1981 thinking: thinking.clone(),
1982 },
1983 ))
1984 .await
1985 {
1986 tracing::warn!(
1987 session_id = %session_id,
1988 error = %e,
1989 "ReasonAtom: failed to emit reason.thinking.completed event"
1990 );
1991 }
1992 termination = StreamTermination::Completed(metadata);
1993 break;
1994 }
1995 LlmStreamEvent::Error(err) => {
1996 let has_partial_output = !tool_calls.is_empty() || !text.is_empty();
2001
2002 if has_partial_output {
2003 tracing::warn!(
2004 session_id = %session_id,
2005 error = %err,
2006 tool_call_count = tool_calls.len(),
2007 text_len = text.len(),
2008 "ReasonAtom: trailing stream error after valid output — treating as partial success"
2009 );
2010 termination = StreamTermination::PartialSuccess;
2014 break;
2015 }
2016
2017 if replay_state.should_retry(
2018 &err,
2019 stream_retry_metadata.attempts,
2020 retry_config.max_retries,
2021 ) {
2022 let proposed_wait =
2023 retry_config.calculate_backoff(stream_retry_metadata.attempts);
2024 let Some(wait_duration) = reserve_retry_wait(
2025 &retry_config,
2026 &mut retry_started_at,
2027 proposed_wait,
2028 ) else {
2029 return Err(AgentLoopError::llm_kind(
2030 err.kind(),
2031 format!(
2032 "{err}; automatic recovery time budget exhausted after {} retries; the turn is safe to resume",
2033 stream_retry_metadata.attempts
2034 ),
2035 )
2036 .with_retry_metadata(&stream_retry_metadata));
2037 };
2038 tracing::warn!(
2039 session_id = %session_id,
2040 turn_id = %context.turn_id,
2041 attempt = stream_retry_metadata.attempts + 1,
2042 max_retries = retry_config.max_retries,
2043 wait_secs = wait_duration.as_secs_f64(),
2044 error_code = err.code.as_deref().unwrap_or("none"),
2045 error_status = err.status,
2046 error = %err,
2047 "ReasonAtom: transient stream error before output, retrying"
2048 );
2049 stream_retry_metadata.record_retry(wait_duration, None);
2050 tokio::time::sleep(wait_duration).await;
2051 continue 'stream_attempt;
2052 }
2053
2054 let llm_duration_ms = llm_start.elapsed().as_millis() as u64;
2056 let event_context = EventContext::from_execution_context(context)
2057 .with_span(
2058 trace_id.to_string(),
2059 Uuid::now_v7().to_string(),
2060 Some(reason_span_id.to_string()),
2061 );
2062 let tools_summary: Vec<ToolDefinitionSummary> =
2063 runtime_agent.tools.iter().map(|t| t.into()).collect();
2064 let generation_data = LlmGenerationData::failure(
2065 messages_for_event.clone(),
2066 tools_summary,
2067 runtime_agent.model.clone(),
2068 Some(model_with_provider.provider_type.to_string()),
2069 err.to_string(),
2070 Some(llm_duration_ms),
2071 time_to_first_token_ms,
2072 );
2073 let _ = self
2074 .event_emitter
2075 .emit(EventRequest::new(
2076 session_id,
2077 event_context,
2078 generation_data,
2079 ))
2080 .await;
2081 return Err(AgentLoopError::llm_kind(err.kind(), err.to_string()));
2082 }
2083 _ => {}
2089 }
2090 if last_stream_heartbeat.elapsed().as_millis() as u64 >= 5_000
2093 && let Some(ref hb) = self.stream_heartbeater
2094 {
2095 hb.heartbeat(crate::durability::StreamProgress {
2096 accumulated_len: text.len() + thinking.len(),
2097 last_delta_at: last_token_at_unix,
2098 })
2099 .await;
2100 last_stream_heartbeat = Instant::now();
2101 }
2102 }
2103 let (mut completion_metadata, tripped) = termination.into_parts();
2104 if let Some(metadata) = completion_metadata.as_mut() {
2105 metadata.retry_metadata =
2106 merge_retry_metadata(metadata.retry_metadata.take(), &stream_retry_metadata);
2107 }
2108
2109 break 'stream_attempt (
2110 text,
2111 thinking,
2112 reasoning,
2113 tool_calls,
2114 completion_metadata,
2115 time_to_first_token_ms,
2116 pending_delta,
2117 tripped,
2118 );
2119 };
2120 let (mut text, mut thinking, mut reasoning, mut tool_calls) =
2121 (text, thinking, reasoning, tool_calls);
2122
2123 let mut citation_annotations: Vec<crate::message::TextAnnotation> = Vec::new();
2134 if tripped.is_none()
2135 && !annotation_providers.is_empty()
2136 && !text.is_empty()
2137 && tool_calls.is_empty()
2138 {
2139 text = filter_response_text(
2142 &self.capability_registry,
2143 &resolved_capability_configs,
2144 text,
2145 );
2146 let collected = collect_annotations(
2147 &annotation_providers,
2148 &runtime_agent.system_prompt,
2149 &text,
2150 &messages,
2151 self.utility_llm_service.as_ref(),
2152 )
2153 .await;
2154 text = collected.text;
2155 citation_annotations = collected.annotations;
2156
2157 if !citation_annotations.is_empty() && !post_output_providers.is_empty() {
2160 let guarded_output = client_visible_guardrail_text(
2161 &text,
2162 &thinking,
2163 &reasoning,
2164 &citation_annotations,
2165 );
2166 let ctx = PostGenerationOutputContext {
2167 system_prompt: &runtime_agent.system_prompt,
2168 message_text: &guarded_output,
2169 utility_llm_service: self.utility_llm_service.as_ref(),
2170 decisions: self.decisions.as_ref(),
2171 };
2172 tripped = evaluate_post_generation_guardrails(&post_output_providers, &ctx).await;
2173 }
2174
2175 if tripped.is_none()
2178 && !citation_annotations.is_empty()
2179 && !citation_verifiers.is_empty()
2180 {
2181 citation_annotations = verify_annotations(
2182 &citation_verifiers,
2183 &text,
2184 self.utility_llm_service.as_ref(),
2185 citation_annotations,
2186 )
2187 .await;
2188 }
2189 }
2190
2191 if tripped.is_none()
2194 && citation_annotations.is_empty()
2195 && !post_output_providers.is_empty()
2196 && (!text.is_empty() || !thinking.is_empty() || !reasoning.is_empty())
2197 {
2198 let guarded_output = client_visible_guardrail_text(&text, &thinking, &reasoning, &[]);
2199 let ctx = PostGenerationOutputContext {
2200 system_prompt: &runtime_agent.system_prompt,
2201 message_text: &guarded_output,
2202 utility_llm_service: self.utility_llm_service.as_ref(),
2203 decisions: self.decisions.as_ref(),
2204 };
2205 tripped = evaluate_post_generation_guardrails(&post_output_providers, &ctx).await;
2206 }
2207
2208 if tripped.is_some() {
2209 citation_annotations.clear();
2210 }
2211
2212 if tripped.is_none() {
2216 for item in &reasoning {
2217 if let Err(e) = self
2218 .event_emitter
2219 .emit(EventRequest::new(
2220 session_id,
2221 streaming_event_context.clone(),
2222 ReasonItemData {
2223 turn_id: context.turn_id,
2224 provider: item.provider.clone(),
2225 model: Some(llm_config.model.clone()),
2226 item_id: item.item_id.clone().unwrap_or_default(),
2227 summary: item
2228 .display_text()
2229 .filter(|_| !matches!(item.text, Some(ReasoningText::Plain { .. })))
2230 .into_iter()
2231 .collect(),
2232 token_count: item.tokens,
2233 },
2234 ))
2235 .await
2236 {
2237 tracing::warn!(
2238 session_id = %session_id,
2239 error = %e,
2240 "ReasonAtom: failed to emit reason.item event"
2241 );
2242 }
2243 }
2244 }
2245
2246 if buffer_output_deltas
2249 && tripped.is_none()
2250 && !pending_delta.is_empty()
2251 && let Err(e) = self
2252 .event_emitter
2253 .emit(EventRequest::new(
2254 session_id,
2255 streaming_event_context.clone(),
2256 OutputMessageDeltaData {
2257 turn_id: context.turn_id,
2258 message_id: output_message_id,
2259 delta: pending_delta.clone(),
2260 accumulated: text.clone(),
2261 phase: streamed_phase,
2262 },
2263 ))
2264 .await
2265 {
2266 tracing::warn!(
2267 session_id = %session_id,
2268 error = %e,
2269 "ReasonAtom: failed to emit guarded output.message.delta event"
2270 );
2271 }
2272
2273 if let Some(ref t) = tripped {
2279 let replaced_event_context = EventContext::from_execution_context(context).with_span(
2280 trace_id.to_string(),
2281 Uuid::now_v7().to_string(),
2282 Some(reason_span_id.to_string()),
2283 );
2284 if let Err(e) = self
2285 .event_emitter
2286 .emit(EventRequest::new(
2287 session_id,
2288 replaced_event_context,
2289 OutputMessageReplacedData {
2290 turn_id: context.turn_id,
2291 message_id: output_message_id,
2292 guardrail_capability_id: t.capability_id.clone(),
2293 guardrail_id: t.guardrail_id.clone(),
2294 reason_code: t.block.reason_code.clone(),
2295 replacement: t.block.replacement.clone(),
2296 },
2297 ))
2298 .await
2299 {
2300 tracing::warn!(
2301 session_id = %session_id,
2302 error = %e,
2303 "ReasonAtom: failed to emit output.message.replaced event"
2304 );
2305 }
2306 text = t.block.replacement.clone();
2307 tool_calls.clear();
2308 thinking.clear();
2309 reasoning.clear();
2310 }
2311
2312 let rejected_tool_calls = if tool_calls.is_empty() {
2316 Vec::new()
2317 } else {
2318 self.apply_finalized_tool_call_hooks(
2319 session_id,
2320 context,
2321 &resolved_capability_configs,
2322 &runtime_agent.tools,
2323 &mut tool_calls,
2324 iteration,
2325 )
2326 .await
2327 };
2328 let finalized_tool_calls = tool_calls.clone();
2329 let rejected_tool_call_ids: HashSet<_> = rejected_tool_calls
2330 .iter()
2331 .map(|rejection| rejection.tool_call_id.clone())
2332 .collect();
2333 tool_calls.retain(|call| !rejected_tool_call_ids.contains(&call.id));
2334
2335 let llm_duration_ms = llm_start.elapsed().as_millis() as u64;
2336
2337 let response_id = completion_metadata
2338 .as_ref()
2339 .and_then(|meta| meta.response_id.clone());
2340 let finish_reason = completion_metadata
2341 .as_ref()
2342 .and_then(|meta| meta.finish_reason.clone());
2343
2344 let usage = completion_metadata.as_ref().and_then(|meta| {
2352 match (meta.prompt_tokens, meta.completion_tokens) {
2353 (Some(input), Some(output)) => {
2354 let actual_cost_usd = meta.provider_cost_usd;
2355 let estimated_cost_usd = crate::model_profiles::estimate_cost_usd(
2356 &model_with_provider.provider_type,
2357 &runtime_agent.model,
2358 input,
2359 output,
2360 meta.cache_read_tokens.unwrap_or(0),
2361 meta.cache_creation_tokens.unwrap_or(0),
2362 );
2363 Some(
2364 TokenUsage::with_cache(
2365 input,
2366 output,
2367 meta.cache_read_tokens,
2368 meta.cache_creation_tokens,
2369 )
2370 .with_cost(actual_cost_usd, estimated_cost_usd),
2371 )
2372 }
2373 _ => None,
2374 }
2375 });
2376
2377 let event_context = EventContext::from_execution_context(context).with_span(
2379 trace_id.to_string(),
2380 Uuid::now_v7().to_string(),
2381 Some(reason_span_id.to_string()),
2382 );
2383 let tools_summary: Vec<ToolDefinitionSummary> =
2384 runtime_agent.tools.iter().map(|t| t.into()).collect();
2385 let finish_reasons = Some(vec![finish_reason.clone().unwrap_or_else(|| {
2386 if finalized_tool_calls.is_empty() {
2387 "stop".to_string()
2388 } else {
2389 "tool_calls".to_string()
2390 }
2391 })]);
2392 let meta = completion_metadata.as_ref();
2393 let served = meta.and_then(|m| m.response_model.clone());
2394 let retry_info = completion_metadata
2395 .as_ref()
2396 .and_then(|meta| meta.retry_metadata.as_ref())
2397 .filter(|rm| rm.had_retries())
2398 .map(|rm| LlmRetryInfo {
2399 attempts: rm.attempts,
2400 total_wait_ms: rm.total_retry_wait.as_millis() as u64,
2401 });
2402 let mut generation_data = LlmGenerationData::success_with_retry(
2403 messages_for_event.clone(),
2404 tools_summary,
2405 Some(text.clone()).filter(|s| !s.is_empty()),
2406 finalized_tool_calls.clone(),
2407 runtime_agent.model.clone(),
2408 Some(model_with_provider.provider_type.to_string()),
2409 usage.clone(),
2410 Some(llm_duration_ms),
2411 time_to_first_token_ms,
2412 finish_reasons,
2413 response_id.clone(),
2414 retry_info,
2415 )
2416 .with_response_model(served);
2417
2418 if let Some(info) = compaction_info {
2424 if let Some(compaction_cost) = info.cost_usd {
2425 match generation_data.metadata.usage.as_mut() {
2426 Some(usage) => {
2427 add_compaction_cost(usage, compaction_cost);
2428 }
2429 None => {
2434 generation_data.metadata.usage = Some(crate::events::TokenUsage {
2435 input_tokens: 0,
2436 output_tokens: 0,
2437 cache_read_tokens: None,
2438 cache_creation_tokens: None,
2439 actual_cost_usd: Some(compaction_cost),
2440 estimated_cost_usd: None,
2441 effective_cost_usd: None,
2442 });
2443 }
2444 }
2445 }
2446 generation_data = generation_data.with_compaction(info);
2447 }
2448
2449 if let Some(request_options) =
2450 build_request_options(&llm_config, &model_with_provider.provider_type.to_string())
2451 {
2452 generation_data = generation_data.with_request_options(request_options);
2453 }
2454
2455 if let Err(e) = self
2456 .event_emitter
2457 .emit(EventRequest::new(
2458 session_id,
2459 event_context,
2460 generation_data,
2461 ))
2462 .await
2463 {
2464 tracing::warn!(
2465 session_id = %session_id,
2466 error = %e,
2467 "ReasonAtom: failed to emit llm.generation event"
2468 );
2469 }
2470
2471 let mut metadata = std::collections::HashMap::new();
2473 metadata.insert(
2474 "model".to_string(),
2475 serde_json::Value::String(runtime_agent.model.clone()),
2476 );
2477 if let Some(state) = &llm_config.reasoning_state {
2478 metadata.insert(
2479 reasoning_updates::STATE_KEY.to_string(),
2480 serde_json::json!(state),
2481 );
2482 }
2483 if let Some(effort) = llm_config
2484 .reasoning_state
2485 .as_ref()
2486 .and_then(|state| state.effective)
2487 .or(reasoning_effort)
2488 {
2489 metadata.insert(
2490 "reasoning_effort".to_string(),
2491 serde_json::Value::String(effort.as_str().to_string()),
2492 );
2493 }
2494 metadata.insert(
2501 "provider".to_string(),
2502 serde_json::Value::String(model_with_provider.provider_type.to_string()),
2503 );
2504 if let Some(ref rid) = response_id {
2505 metadata.insert(
2506 "response_id".to_string(),
2507 serde_json::Value::String(rid.clone()),
2508 );
2509 }
2510
2511 let text = filter_response_text(
2515 &self.capability_registry,
2516 &resolved_capability_configs,
2517 text,
2518 );
2519 let has_tool_calls = !finalized_tool_calls.is_empty();
2520 let mut assistant_message = if has_tool_calls {
2521 RuntimeMessage::assistant_with_tools(&text, finalized_tool_calls.clone())
2522 } else {
2523 RuntimeMessage::assistant(&text)
2524 }
2525 .with_id(output_message_id);
2526 for part in &mut assistant_message.content {
2527 if let crate::message::ContentPart::ToolCall(call) = part {
2528 call.native = native_calls.get(&call.id).cloned();
2529 }
2530 }
2531 if !citation_annotations.is_empty() {
2534 for part in assistant_message.content.iter_mut() {
2535 if let crate::message::ContentPart::Text(t) = part {
2536 t.annotations = std::mem::take(&mut citation_annotations);
2537 break;
2538 }
2539 }
2540 }
2541 let provider_type_for_reasoning = model_with_provider.provider_type.to_string();
2545 let provider_phase = completion_metadata
2549 .as_ref()
2550 .and_then(|meta| meta.phase.as_deref())
2551 .and_then(everruns_provider::ExecutionPhase::from_provider_str);
2552 let (phase, phase_source) = match provider_phase {
2553 Some(phase) => (phase, everruns_provider::PhaseSource::Provider),
2554 None => (
2555 everruns_provider::ExecutionPhase::from_has_tool_calls(has_tool_calls),
2556 everruns_provider::PhaseSource::Derived,
2557 ),
2558 };
2559 assistant_message.phase = Some(phase);
2560 assistant_message.phase_source = Some(phase_source);
2561 assistant_message.metadata = Some(metadata);
2562 if reasoning.is_empty() && !thinking.is_empty() {
2571 reasoning.push(
2572 ReasoningContentPart::opaque(provider_type_for_reasoning.clone()).with_text(
2573 ReasoningText::Plain {
2574 text: thinking.clone(),
2575 },
2576 ),
2577 );
2578 }
2579 if !reasoning.is_empty() {
2580 let mut content = Vec::with_capacity(reasoning.len() + assistant_message.content.len());
2581 content.extend(reasoning.drain(..).map(ContentPart::Reasoning));
2582 content.append(&mut assistant_message.content);
2583 assistant_message.content = content;
2584 }
2585 let message_event_context = EventContext::from_execution_context(context).with_span(
2588 trace_id.to_string(),
2589 Uuid::now_v7().to_string(),
2590 Some(reason_span_id.to_string()),
2591 );
2592 let mut output_message_data = OutputMessageCompletedData::new(assistant_message);
2593 if let Some(ref u) = usage {
2594 output_message_data = output_message_data.with_usage(u.clone());
2595 }
2596 let result = ReasonResult {
2597 native_counts: None,
2598 success: true,
2599 text,
2600 tool_calls,
2601 has_tool_calls,
2602 tool_definitions: runtime_agent.tools.clone(),
2603 max_iterations: runtime_agent.max_iterations,
2604 error: None,
2605 user_facing_error: None,
2606 error_disclosure: None,
2607 usage,
2608 output_message_id: Some(output_message_id),
2609 time_to_first_token_ms,
2610 response_id,
2611 finish_reason,
2612 locale: resolved_locale,
2613 network_access: runtime_agent.network_access.clone(),
2614 parallel_tool_calls: runtime_agent.parallel_tool_calls,
2615 };
2616 if let Some(coordinator) = &self.native_async {
2617 coordinator
2618 .lock()
2619 .await
2620 .stage_transcript_result(
2621 serde_json::to_value(&result)
2622 .map_err(|error| AgentLoopError::store(error.to_string()))?,
2623 )
2624 .await?;
2625 }
2626 self.event_emitter
2627 .emit(EventRequest::new(
2628 session_id,
2629 message_event_context,
2630 output_message_data,
2631 ))
2632 .await?;
2633
2634 if let Some(coordinator) = &self.native_async {
2635 coordinator
2636 .lock()
2637 .await
2638 .transcript_committed(&output_message_id.to_string())
2639 .await?;
2640 }
2641 for rejection in rejected_tool_calls {
2642 let Some(call) = finalized_tool_calls
2643 .iter()
2644 .find(|call| call.id == rejection.tool_call_id)
2645 else {
2646 continue;
2647 };
2648 self.event_emitter
2649 .emit(EventRequest::new(
2650 session_id,
2651 EventContext::from_execution_context(context),
2652 ToolCompletedData::failure(
2653 call.id.clone(),
2654 call.name.clone(),
2655 "error".to_string(),
2656 rejection.error,
2657 None,
2658 ),
2659 ))
2660 .await?;
2661 }
2662 tracing::info!(
2663 session_id = %session_id,
2664 turn_id = %context.turn_id,
2665 has_tool_calls = %result.has_tool_calls,
2666 tool_count = %result.tool_calls.len(),
2667 "ReasonAtom: LLM call completed"
2668 );
2669
2670 Ok(result)
2671 }
2672
2673 async fn finalize_partial_stream(
2678 &self,
2679 session_id: SessionId,
2680 context: &ExecutionContext,
2681 partial: PartialStreamState,
2682 iteration: u32,
2683 runtime_agent: &crate::RuntimeAgent,
2684 resolved_capability_configs: &[crate::CapabilityRef],
2685 ) -> Result<ReasonResult> {
2686 let event_context = EventContext::from_execution_context(context);
2687 let turn_id = context.turn_id;
2688 let message_id = partial.message_id;
2689
2690 let _ = self
2692 .event_emitter
2693 .emit(EventRequest::new(
2694 session_id,
2695 event_context.clone(),
2696 OutputMessageStartedData {
2697 reasoning_state: partial.reasoning_state.clone(),
2698 turn_id,
2699 message_id,
2700 model: None,
2701 iteration: Some(iteration),
2702 phase: None,
2705 },
2706 ))
2707 .await;
2708
2709 let accumulated = filter_response_text(
2712 &self.capability_registry,
2713 resolved_capability_configs,
2714 partial.accumulated,
2715 );
2716 let mut assistant_message = RuntimeMessage::assistant(&accumulated).with_id(message_id);
2717 if let Some(state) = partial.reasoning_state {
2718 assistant_message.metadata = Some(HashMap::from([
2719 ("model".into(), serde_json::json!("gpt-6-astra")),
2720 ("provider".into(), serde_json::json!("openai")),
2721 (
2722 reasoning_updates::STATE_KEY.into(),
2723 serde_json::json!(state),
2724 ),
2725 (
2726 "reasoning_effort".into(),
2727 serde_json::json!(state.effective),
2728 ),
2729 ]));
2730 }
2731 let output_message_id = message_id;
2732 self.event_emitter
2733 .emit(EventRequest::new(
2734 session_id,
2735 event_context.clone(),
2736 OutputMessageCompletedData::new(assistant_message),
2737 ))
2738 .await?;
2739
2740 let accumulated_len = accumulated.len();
2742 let _ = self
2743 .event_emitter
2744 .emit(EventRequest::new(
2745 session_id,
2746 event_context.clone(),
2747 ReasonRecoveredData {
2748 turn_id,
2749 mode: RecoveryMode::Finalize,
2750 accumulated_len,
2751 },
2752 ))
2753 .await;
2754
2755 tracing::info!(
2756 session_id = %session_id,
2757 turn_id = %turn_id,
2758 accumulated_len,
2759 "ReasonAtom: finalized partial stream from persisted accumulated text"
2760 );
2761
2762 Ok(ReasonResult {
2763 native_counts: None,
2764 success: true,
2765 text: accumulated,
2766 tool_calls: vec![],
2767 has_tool_calls: false,
2768 tool_definitions: runtime_agent.tools.clone(),
2769 max_iterations: runtime_agent.max_iterations,
2770 error: None,
2771 user_facing_error: None,
2772 error_disclosure: None,
2773 usage: None,
2774 output_message_id: Some(output_message_id),
2775 time_to_first_token_ms: None,
2776 response_id: None,
2777 finish_reason: Some("stop".to_string()),
2778 locale: None,
2779 network_access: None,
2780 parallel_tool_calls: None,
2782 })
2783 }
2784
2785 async fn resolve_images(&self, messages: &[RuntimeMessage]) -> HashMap<Uuid, ResolvedImage> {
2796 let mut resolved = HashMap::new();
2797
2798 let resolver = match &self.image_resolver {
2800 Some(r) => r,
2801 None => return resolved,
2802 };
2803
2804 let image_ids: Vec<Uuid> = messages
2806 .iter()
2807 .flat_map(crate::llm_conversions::extract_image_file_ids)
2808 .collect::<std::collections::HashSet<_>>()
2809 .into_iter()
2810 .collect();
2811
2812 if image_ids.is_empty() {
2813 return resolved;
2814 }
2815
2816 tracing::debug!(
2817 image_count = image_ids.len(),
2818 "ReasonAtom: resolving image_file references"
2819 );
2820
2821 for image_id in image_ids {
2823 match resolver.resolve_image(image_id).await {
2824 Ok(Some(image)) => {
2825 resolved.insert(image_id, image);
2826 }
2827 Ok(None) => {
2828 tracing::warn!(
2829 image_id = %image_id,
2830 "ReasonAtom: image not found during resolution"
2831 );
2832 }
2833 Err(e) => {
2834 tracing::warn!(
2835 image_id = %image_id,
2836 error = %e,
2837 "ReasonAtom: failed to resolve image"
2838 );
2839 }
2840 }
2841 }
2842
2843 tracing::debug!(
2844 resolved_count = resolved.len(),
2845 "ReasonAtom: image resolution complete"
2846 );
2847
2848 resolved
2849 }
2850
2851 async fn resolve_files(&self, messages: &[RuntimeMessage]) -> HashMap<Uuid, ResolvedFile> {
2852 let Some(resolver) = &self.file_resolver else {
2853 return HashMap::new();
2854 };
2855
2856 let file_ids: Vec<Uuid> = messages
2857 .iter()
2858 .flat_map(crate::llm_conversions::extract_file_ids)
2859 .collect::<std::collections::HashSet<_>>()
2860 .into_iter()
2861 .collect();
2862
2863 if file_ids.is_empty() {
2864 return HashMap::new();
2865 }
2866
2867 match resolver.resolve_files(&file_ids).await {
2868 Ok(map) => map,
2869 Err(e) => {
2870 tracing::warn!(
2871 target: "reason",
2872 "ReasonAtom: file resolution failed: {e}"
2873 );
2874 HashMap::new()
2875 }
2876 }
2877 }
2878}
2879
2880#[cfg(test)]
2885mod tests;