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 EventContext, EventRequest, LlmCompactionInfo, LlmGenerationData, LlmRetryInfo,
41 OutputMessageCompletedData, OutputMessageDeltaData, OutputMessageReplacedData,
42 OutputMessageStartedData, ReasonCompletedData, ReasonItemData, ReasonRecoveredData,
43 ReasonStartedData, ReasonThinkingCompletedData, ReasonThinkingDeltaData,
44 ReasonThinkingStartedData, RecoveryMode, TokenUsage, ToolCompletedData, ToolDefinitionSummary,
45};
46use crate::llm_retry::{
47 LlmRetryConfig, RetryMetadata, is_transient_error_message, remaining_retry_time,
48 reserve_retry_wait,
49};
50use crate::message::{ContentPart, RuntimeMessage, RuntimeMessageRole};
51use crate::message_retriever::MessageRetriever;
52use crate::output_guardrail::{
53 ArmedGuardrail, OutputGuardrailContext, PostGenerationOutputContext, evaluate_guardrails,
54 evaluate_post_generation_guardrails,
55};
56use crate::phase_effects::{PhaseEffectEmitter, PhaseEffectSink};
57use crate::runtime_context::{AssembledTurnContext, TurnContextRequest, TurnContextResolver};
58use crate::tool_types::{ToolCall, ToolDefinition};
59use crate::typed_id::{AgentId, HarnessId, MessageId, SessionId};
60use crate::{ErrorDisclosure, UserFacingError, UserFacingErrorContext};
61use crate::{
62 durability::DurableToolResultStore,
63 durability::PartialStreamState,
64 durability::PartialStreamStore,
65 file_services::{FileResolver, ResolvedFile},
66 image_services::ImageResolver,
67 image_services::ResolvedImage,
68};
69use everruns_provider::reasoning::{ReasoningContentPart, ReasoningText};
70
71mod background_call;
72mod compaction;
73mod error_policy;
74mod facts;
75mod finalized_calls;
76mod hosted_tools;
77mod observability;
78mod output_hooks;
79mod provider_managed_compaction;
80mod reasoning_updates;
81mod request_controls;
82mod stream_state;
83mod transcript;
84
85use compaction::{
86 ProactiveCompactionContext, ReactiveCompactionContext, apply_proactive_compaction,
87 apply_reactive_compaction,
88};
89use error_policy::{filter_response_text, is_error_placeholder_message};
90#[cfg(test)]
91use observability::capability_usage_snapshot_records;
92pub use observability::capability_usage_snapshot_records as capability_usage_records;
93use observability::{build_request_options, emit_capability_usage_snapshot};
94use output_hooks::{client_visible_guardrail_text, collect_output_hooks};
95use request_controls::resolve_request_controls;
96use stream_state::{
97 StreamReplayState, StreamTermination, advances_stall_deadline, append_guarded_thinking_delta,
98 inspect_guarded_reasoning_item, merge_retry_metadata,
99};
100use transcript::repair_dangling_tool_calls;
101
102fn unix_now_secs() -> u64 {
103 std::time::SystemTime::now()
104 .duration_since(std::time::UNIX_EPOCH)
105 .unwrap_or_default()
106 .as_secs()
107}
108
109#[derive(Debug, Clone, Serialize, Deserialize)]
111pub struct ReasonInput {
112 pub context: ExecutionContext,
114 pub harness_id: HarnessId,
116 #[serde(skip_serializing_if = "Option::is_none")]
118 pub agent_id: Option<AgentId>,
119 #[serde(default)]
121 pub org_id: i64,
122 #[serde(default)]
126 pub mcp_tool_definitions: Vec<ToolDefinition>,
127 #[serde(skip_serializing_if = "Option::is_none")]
130 pub previous_response_id: Option<String>,
131 #[serde(default = "default_iteration")]
134 pub iteration: u32,
135}
136
137fn default_iteration() -> u32 {
138 1
139}
140
141#[derive(Debug, Clone, Serialize, Deserialize)]
143pub struct NativeExecutionCounts {
144 pub llm_calls: u32,
145 pub tool_calls: u32,
146}
147
148#[derive(Debug, Clone, Default, Serialize, Deserialize)]
150pub struct ReasonResult {
151 #[serde(default, skip_serializing_if = "Option::is_none")]
152 pub native_counts: Option<NativeExecutionCounts>,
153 pub success: bool,
155 pub text: String,
157 #[serde(default)]
159 pub tool_calls: Vec<ToolCall>,
160 pub has_tool_calls: bool,
162 #[serde(default)]
164 pub tool_definitions: Vec<ToolDefinition>,
165 #[serde(default = "default_max_iterations")]
167 pub max_iterations: usize,
168 #[serde(skip_serializing_if = "Option::is_none")]
170 pub error: Option<String>,
171 #[serde(default, skip_serializing_if = "Option::is_none")]
175 pub user_facing_error: Option<UserFacingError>,
176 #[serde(default, skip_serializing_if = "Option::is_none")]
178 pub error_disclosure: Option<ErrorDisclosure>,
179 #[serde(skip_serializing_if = "Option::is_none")]
181 pub usage: Option<TokenUsage>,
182 #[serde(skip_serializing_if = "Option::is_none")]
184 pub output_message_id: Option<MessageId>,
185 #[serde(skip_serializing_if = "Option::is_none")]
187 pub time_to_first_token_ms: Option<u64>,
188 #[serde(skip_serializing_if = "Option::is_none")]
190 pub response_id: Option<String>,
191 #[serde(default, skip_serializing_if = "Option::is_none")]
193 pub finish_reason: Option<String>,
194 #[serde(skip_serializing_if = "Option::is_none")]
196 pub locale: Option<String>,
197 #[serde(default, skip_serializing_if = "Option::is_none")]
199 pub network_access: Option<crate::network_access::NetworkAccessList>,
200 #[serde(default, skip_serializing_if = "Option::is_none")]
204 pub parallel_tool_calls: Option<bool>,
205 #[serde(default, skip_serializing_if = "std::ops::Not::not")]
207 pub waiting_for_tool_results: bool,
208}
209
210fn default_max_iterations() -> usize {
211 500
212}
213
214pub struct ReasonAtom {
233 native_async: Option<Arc<tokio::sync::Mutex<crate::native_async::NativeAsyncCoordinator>>>,
234 context_resolver: Arc<dyn TurnContextResolver>,
235 message_retriever: Arc<dyn MessageRetriever>,
236 capability_registry: CapabilityRegistry,
237 event_emitter: PhaseEffectEmitter<dyn PhaseEffectSink>,
238 image_resolver: Option<Arc<dyn ImageResolver>>,
240 file_resolver: Option<Arc<dyn FileResolver>>,
241 stream_heartbeater: Option<Arc<dyn crate::durability::StreamHeartbeater>>,
243 provider_stall_timeout: Option<std::time::Duration>,
245 provider_retry_config: LlmRetryConfig,
247 durable_tool_result_store: Option<Arc<dyn DurableToolResultStore>>,
249 partial_stream_store: Option<Arc<dyn PartialStreamStore>>,
251 reasoning_effort_handle: Option<crate::tool_context::ReasoningEffortHandle>,
255 utility_llm_service: Option<Arc<dyn crate::UtilityLlmService>>,
259 decisions: Option<Arc<dyn crate::DecisionsService>>,
263 schedule_store: Option<Arc<dyn crate::session_services::SessionScheduleStore>>,
269 compaction_checkpoint_store: Option<Arc<dyn crate::CompactionCheckpointStore>>,
271 background_call: everruns_provider::background_call::BackgroundCallContext,
273}
274
275impl ReasonAtom {
276 pub fn with_native_async(
277 mut self,
278 coordinator: Arc<tokio::sync::Mutex<crate::native_async::NativeAsyncCoordinator>>,
279 ) -> Self {
280 self.native_async = Some(coordinator);
281 self
282 }
283
284 pub fn new(
286 context_resolver: impl TurnContextResolver + 'static,
287 message_retriever: impl MessageRetriever + 'static,
288 capability_registry: CapabilityRegistry,
289 event_emitter: impl PhaseEffectSink + 'static,
290 ) -> Self {
291 Self {
292 native_async: None,
293 context_resolver: Arc::new(context_resolver),
294 message_retriever: Arc::new(message_retriever),
295 capability_registry,
296 event_emitter: PhaseEffectEmitter::new(Arc::new(event_emitter)),
297 image_resolver: None,
298 file_resolver: None,
299 stream_heartbeater: None,
300 provider_stall_timeout: None,
301 provider_retry_config: LlmRetryConfig::default(),
302 durable_tool_result_store: None,
303 partial_stream_store: None,
304 reasoning_effort_handle: None,
305 utility_llm_service: None,
306 decisions: None,
307 schedule_store: None,
308 compaction_checkpoint_store: None,
309 background_call: Default::default(),
310 }
311 }
312
313 pub fn with_schedule_store(
316 mut self,
317 store: Arc<dyn crate::session_services::SessionScheduleStore>,
318 ) -> Self {
319 self.schedule_store = Some(store);
320 self
321 }
322
323 pub fn with_compaction_checkpoint_store(
324 mut self,
325 store: Arc<dyn crate::CompactionCheckpointStore>,
326 ) -> Self {
327 self.compaction_checkpoint_store = Some(store);
328 self
329 }
330
331 pub fn with_image_resolver(mut self, resolver: Arc<dyn ImageResolver>) -> Self {
344 self.image_resolver = Some(resolver);
345 self
346 }
347
348 pub fn with_file_resolver(mut self, resolver: Arc<dyn FileResolver>) -> Self {
349 self.file_resolver = Some(resolver);
350 self
351 }
352
353 pub fn with_stream_heartbeater(
355 mut self,
356 heartbeater: Arc<dyn crate::durability::StreamHeartbeater>,
357 ) -> Self {
358 self.stream_heartbeater = Some(heartbeater);
359 self
360 }
361
362 pub fn with_provider_stall_timeout(mut self, timeout: std::time::Duration) -> Self {
365 self.provider_stall_timeout = Some(timeout);
366 self
367 }
368
369 pub fn with_provider_retry_config(mut self, config: LlmRetryConfig) -> Self {
371 self.provider_retry_config = config;
372 self
373 }
374
375 pub fn with_durable_tool_result_store(
381 mut self,
382 store: Arc<dyn DurableToolResultStore>,
383 ) -> Self {
384 self.durable_tool_result_store = Some(store);
385 self
386 }
387
388 pub fn with_partial_stream_store(mut self, store: Arc<dyn PartialStreamStore>) -> Self {
390 self.partial_stream_store = Some(store);
391 self
392 }
393
394 pub fn with_reasoning_effort_handle(
401 mut self,
402 handle: crate::tool_context::ReasoningEffortHandle,
403 ) -> Self {
404 self.reasoning_effort_handle = Some(handle);
405 self
406 }
407
408 pub fn with_utility_llm_service(mut self, service: Arc<dyn crate::UtilityLlmService>) -> Self {
411 self.utility_llm_service = Some(service);
412 self
413 }
414
415 pub fn with_decisions(mut self, service: Arc<dyn crate::DecisionsService>) -> Self {
418 self.decisions = Some(service);
419 self
420 }
421}
422
423impl ReasonAtom {
424 pub fn name(&self) -> &'static str {
426 "reason"
427 }
428
429 pub async fn execute(&self, input: ReasonInput) -> Result<ReasonResult> {
431 self.execute_inner(input, None).await
432 }
433}
434
435impl ReasonAtom {
436 pub async fn execute_with_assembled_context(
441 &self,
442 input: ReasonInput,
443 assembled: AssembledTurnContext,
444 ) -> Result<ReasonResult> {
445 self.execute_inner(input, Some(assembled)).await
446 }
447
448 async fn execute_inner(
449 &self,
450 input: ReasonInput,
451 assembled: Option<AssembledTurnContext>,
452 ) -> Result<ReasonResult> {
453 let ReasonInput {
454 context,
455 harness_id,
456 agent_id,
457 org_id,
458 mcp_tool_definitions,
459 previous_response_id,
460 iteration,
461 } = input;
462
463 tracing::info!(
464 session_id = %context.session_id,
465 turn_id = %context.turn_id,
466 exec_id = %context.exec_id,
467 harness_id = %harness_id,
468 agent_id = ?agent_id,
469 mcp_tools_count = %mcp_tool_definitions.len(),
470 "ReasonAtom: starting LLM call"
471 );
472
473 let trace_id = context.turn_id.to_string();
481 let reason_span_id = Uuid::now_v7().to_string();
482 let parent_span_id = trace_id.clone(); let event_context = EventContext::from_execution_context(&context).with_span(
486 trace_id.clone(),
487 reason_span_id.clone(),
488 Some(parent_span_id.clone()),
489 );
490
491 let reason_start = Instant::now();
493
494 if let Err(e) = self
496 .event_emitter
497 .emit(EventRequest::new(
498 context.session_id,
499 event_context.clone(),
500 ReasonStartedData {
501 harness_id,
502 agent_id,
503 metadata: None, },
505 ))
506 .await
507 {
508 tracing::warn!(
509 session_id = %context.session_id,
510 error = %e,
511 "ReasonAtom: failed to emit reason.started event"
512 );
513 }
514
515 let assembled = match assembled {
519 Some(assembled) => Ok(assembled),
520 None => {
521 self.context_resolver
522 .resolve_turn_context(TurnContextRequest {
523 session_id: context.session_id,
524 harness_id,
525 agent_id,
526 mcp_tool_definitions: mcp_tool_definitions.clone(),
527 allow_provider_managed_reduction: true,
528 })
529 .await
530 }
531 };
532
533 let (error_disclosure, error_context, error_hooks, call_result) = match assembled {
534 Ok(assembled) => {
535 let outcome = provider_managed_compaction::execute_with_fallback(
536 provider_managed_compaction::Call {
537 atom: self,
538 session_id: context.session_id,
539 harness_id,
540 agent_id,
541 org_id,
542 context: &context,
543 trace_id: &trace_id,
544 reason_span_id: &reason_span_id,
545 previous_response_id,
546 iteration,
547 mcp_tool_definitions: &mcp_tool_definitions,
548 assembled,
549 },
550 )
551 .await;
552 (
553 outcome.disclosure,
554 outcome.error_context,
555 outcome.error_hooks,
556 outcome.result,
557 )
558 }
559 Err(error) => (
560 ErrorDisclosure::default(),
561 UserFacingErrorContext::default(),
562 Vec::new(),
563 Err(error),
564 ),
565 };
566
567 let result = match call_result {
569 Ok(result) => {
570 let reason_duration_ms = reason_start.elapsed().as_millis() as u64;
572
573 let completed_context = EventContext::from_execution_context(&context).with_span(
575 trace_id.clone(),
576 reason_span_id.clone(), Some(parent_span_id.clone()),
578 );
579 if let Err(e) = self
580 .event_emitter
581 .emit(EventRequest::new(
582 context.session_id,
583 completed_context,
584 ReasonCompletedData::success(
585 &result.text,
586 result.has_tool_calls,
587 result.tool_calls.len() as u32,
588 Some(reason_duration_ms),
589 result.usage.clone(),
590 ),
591 ))
592 .await
593 {
594 tracing::warn!(
595 session_id = %context.session_id,
596 error = %e,
597 "ReasonAtom: failed to emit reason.completed event"
598 );
599 }
600 result
601 }
602 Err(e) => {
603 let reason_duration_ms = reason_start.elapsed().as_millis() as u64;
605
606 tracing::warn!(
609 session_id = %context.session_id,
610 turn_id = %context.turn_id,
611 error = %e,
612 "ReasonAtom: LLM call failed"
613 );
614
615 let error_msg = e.to_string();
616 let mut source_error = e.user_facing_error(error_context);
617
618 let is_transient = e.is_transient_llm_error()
625 || (e.llm_error_kind().is_none() && is_transient_error_message(&error_msg));
626
627 if !is_transient && !error_hooks.is_empty() {
633 let services = crate::llm_error_hook::LlmErrorHookServices {
634 schedule_store: self.schedule_store.clone(),
635 };
636 for (hook, config) in &error_hooks {
637 let outcome = {
638 let ctx = crate::llm_error_hook::LlmErrorContext {
639 session_id: context.session_id,
640 error_code: &source_error.code,
641 error_fields: &source_error.fields,
642 config,
643 services: &services,
644 };
645 hook.on_llm_error(&ctx).await
646 };
647 for (key, value) in outcome.extra_error_fields {
648 source_error = source_error.with_field(key, value);
649 }
650 }
651 }
652
653 let user_error = source_error.apply_disclosure(error_disclosure, Some(&error_msg));
654 let user_error_text = user_error.fallback_message();
655
656 let mut output_message_id = None;
657
658 if !is_transient {
659 let mut error_message = RuntimeMessage::assistant(&user_error_text);
661 let mut metadata = std::collections::HashMap::new();
662 user_error.apply_to_message_metadata(&mut metadata);
663 UserFacingError::apply_disclosure_to_message_metadata(
664 &mut metadata,
665 error_disclosure,
666 &source_error.code,
667 );
668 error_message.metadata = Some(metadata);
669
670 output_message_id = Some(error_message.id);
671
672 let error_msg_context = EventContext::from_execution_context(&context)
675 .with_span(
676 trace_id.clone(),
677 Uuid::now_v7().to_string(), Some(reason_span_id.clone()), );
680 if let Err(emit_err) = self
681 .event_emitter
682 .emit(EventRequest::new(
683 context.session_id,
684 error_msg_context,
685 OutputMessageCompletedData::new(error_message)
686 .with_user_facing_error(&user_error)
687 .with_error_disclosure(error_disclosure),
688 ))
689 .await
690 {
691 tracing::warn!(
692 session_id = %context.session_id,
693 error = %emit_err,
694 "ReasonAtom: failed to emit output.message.completed event for error"
695 );
696 }
697 } else {
698 tracing::info!(
699 session_id = %context.session_id,
700 "ReasonAtom: skipping error event for transient LLM error (will be retried)"
701 );
702 }
703
704 let completed_context = EventContext::from_execution_context(&context).with_span(
706 trace_id.clone(),
707 reason_span_id.clone(), Some(parent_span_id.clone()),
709 );
710 if let Err(emit_err) = self
711 .event_emitter
712 .emit(EventRequest::new(
713 context.session_id,
714 completed_context,
715 ReasonCompletedData::failure(error_msg.clone(), Some(reason_duration_ms)),
716 ))
717 .await
718 {
719 tracing::warn!(
720 session_id = %context.session_id,
721 error = %emit_err,
722 "ReasonAtom: failed to emit reason.completed event"
723 );
724 }
725
726 ReasonResult {
727 native_counts: None,
728 success: false,
729 text: user_error_text,
730 tool_calls: vec![],
731 has_tool_calls: false,
732 tool_definitions: vec![],
733 max_iterations: default_max_iterations(),
734 error: Some(error_msg.clone()),
735 user_facing_error: Some(user_error),
736 error_disclosure: Some(error_disclosure),
737 usage: None,
738 output_message_id,
739 time_to_first_token_ms: None,
740 response_id: None,
741 finish_reason: error_msg
742 .to_ascii_lowercase()
743 .contains("model refused")
744 .then(|| "refusal".to_string()),
745 ..ReasonResult::default()
746 }
747 }
748 };
749
750 Ok(result)
751 }
752
753 #[allow(clippy::too_many_arguments)]
755 async fn execute_llm_call(
756 &self,
757 session_id: SessionId,
758 harness_id: HarnessId,
759 agent_id: Option<AgentId>,
760 org_id: i64,
761 context: &ExecutionContext,
762 trace_id: &str,
763 reason_span_id: &str,
764 previous_response_id: Option<String>,
765 iteration: u32,
766 assembled: AssembledTurnContext,
767 ) -> Result<ReasonResult> {
768 let prior_usage = assembled.cumulative_usage();
769 let mut messages = transcript::order_native_results(assembled.messages);
770 let mut message_source_sequence = assembled.message_source_sequence;
771 let model_with_provider = assembled.model;
772 let provider_managed = model_with_provider
773 .provider_managed_reduction_option
774 .is_some();
775 let supports_clear_at = facts::supports_clear_at(
776 &model_with_provider.provider_type,
777 &model_with_provider.model,
778 );
779 let resolved_model_id = assembled.resolved_model_id;
780 let resolved_locale = assembled.resolved_locale;
781 let compaction_policy = assembled.compaction_policy;
782 let resolved_capability_configs = assembled.resolved_capability_configs;
783 let runtime_agent = assembled.runtime_agent;
784 let embedder_metadata = assembled.embedder_metadata;
785
786 emit_capability_usage_snapshot(
787 self.event_emitter.as_ref(),
788 &self.capability_registry,
789 session_id,
790 context,
791 &resolved_capability_configs,
792 &runtime_agent.tools,
793 )
794 .await;
795
796 let output_hooks =
797 collect_output_hooks(&self.capability_registry, &resolved_capability_configs);
798 let guardrail_providers = output_hooks.streaming;
799 let post_output_providers = output_hooks.post_generation;
800 let annotation_providers = output_hooks.annotations;
801 let citation_verifiers = output_hooks.citation_verifiers;
802
803 let chat_driver = Arc::clone(&model_with_provider.driver);
805 let stateful_response_continuation =
806 previous_response_id.is_some() && chat_driver.supports_stateful_responses();
807 let mut restored_checkpoint: Option<crate::CompactionCheckpoint> = None;
808 let mut checkpoint_suffix_message_count = 0usize;
809 let native_reasoning_compaction = compaction_policy.as_ref().is_none_or(|policy| {
810 matches!(
811 policy.settings().strategy,
812 crate::compaction_policy::CompactionStrategy::Native
813 | crate::compaction_policy::CompactionStrategy::Auto
814 ) && chat_driver.supports_compact()
815 });
816
817 let checkpoint_format = super::provider_checkpoint::format_version(provider_managed);
818 if (compaction_policy.is_some() || provider_managed)
819 && let Some(store) = self.compaction_checkpoint_store.as_ref()
820 && let Some(checkpoint) = store
821 .get_latest_format(
822 session_id,
823 model_with_provider.provider_type.as_str(),
824 &model_with_provider.model,
825 checkpoint_format,
826 )
827 .await?
828 && super::provider_checkpoint::is_restorable(
829 &checkpoint,
830 chat_driver.as_ref(),
831 model_with_provider.provider_type.as_str(),
832 &model_with_provider.model,
833 provider_managed,
834 native_reasoning_compaction,
835 )
836 {
837 let filters = crate::capabilities::collect_message_filters_only_with_context(
838 &resolved_capability_configs,
839 &self.capability_registry,
840 provider_managed,
841 );
842 let mut query =
843 crate::MessageQuery::new(session_id).after_sequence(checkpoint.source_sequence);
844 filters.apply_message_filters(&mut query);
845 let history = self.message_retriever.load_filtered_history(query).await?;
846 messages = history.messages;
847 checkpoint_suffix_message_count = messages.len();
848 filters.apply_post_load_filters(&mut messages);
849 if let crate::CompactionCheckpointPayload::Summary { text } = &checkpoint.payload {
850 messages.insert(
851 0,
852 RuntimeMessage::system(format!(
853 "[CONVERSATION_SUMMARY]\n{text}\n[/CONVERSATION_SUMMARY]"
854 )),
855 );
856 }
857 message_source_sequence = history.source_sequence.or(message_source_sequence);
858 restored_checkpoint = Some(checkpoint);
859 }
860
861 let controls = resolve_request_controls(
862 &messages,
863 self.reasoning_effort_handle.as_ref(),
864 &model_with_provider.provider_type,
865 &model_with_provider.model,
866 );
867 let reasoning_effort = controls.reasoning_effort;
868 let speed = controls.speed;
869 let verbosity = controls.verbosity;
870 let checkpoint_reasoning =
871 super::provider_checkpoint::reasoning_state(restored_checkpoint.as_ref());
872 let mut reasoning_replay = reasoning_updates::prepare(
873 &messages,
874 model_with_provider.provider_type.as_str(),
875 &model_with_provider.model,
876 reasoning_effort,
877 self.reasoning_effort_handle
878 .as_ref()
879 .and_then(crate::tool_context::ReasoningEffortHandle::get),
880 checkpoint_reasoning,
881 )
882 .filter(|_| native_reasoning_compaction);
883
884 if let Some(ref store) = self.partial_stream_store {
888 let turn_id_str = context.turn_id.to_string();
889 match store.get_partial_stream(session_id, &turn_id_str).await {
890 Ok(Some(partial)) if !partial.accumulated.is_empty() => {
891 return self
893 .finalize_partial_stream(
894 session_id,
895 context,
896 partial,
897 iteration,
898 &runtime_agent,
899 &resolved_capability_configs,
900 )
901 .await;
902 }
903 Ok(Some(partial)) => {
904 if let (Some(replay), Some(mut saved)) =
905 (reasoning_replay.as_mut(), partial.reasoning_state)
906 {
907 saved.pending = saved.effective;
910 replay.state = saved;
911 }
912 let recovery_ctx = EventContext::from_execution_context(context);
915 let _ = self
916 .event_emitter
917 .emit(EventRequest::new(
918 session_id,
919 recovery_ctx,
920 ReasonRecoveredData {
921 turn_id: context.turn_id,
922 mode: RecoveryMode::Restart,
923 accumulated_len: 0,
924 },
925 ))
926 .await;
927 tracing::info!(
928 session_id = %session_id,
929 turn_id = %context.turn_id,
930 "ReasonAtom: partial stream detected with empty accumulated; restarting clean"
931 );
932 }
933 Ok(None) => {} Err(e) => {
935 if reasoning_replay.is_some() {
936 return Err(e);
937 }
938 tracing::warn!(
940 session_id = %session_id,
941 turn_id = %context.turn_id,
942 error = %e,
943 "ReasonAtom: partial-stream store error; proceeding with normal execution"
944 );
945 }
946 }
947 }
948
949 let repair_event_context = EventContext::from_execution_context(context);
953 let patched_messages = if self.native_async.is_some() {
954 messages.clone()
955 } else {
956 repair_dangling_tool_calls(
957 &messages,
958 self.durable_tool_result_store.as_deref(),
959 self.event_emitter.as_ref(),
960 session_id,
961 &repair_event_context,
962 &context.turn_id.to_string(),
963 )
964 .await
965 };
966 let raw_tool_result_bytes = compaction_policy
967 .as_ref()
968 .map(|policy| policy.total_tool_result_bytes(&patched_messages))
969 .unwrap_or(0);
970
971 let model_view_providers = crate::capabilities::collect_model_view_providers(
974 &resolved_capability_configs,
975 &self.capability_registry,
976 Some(model_with_provider.model.as_str()),
977 );
978 let model_view_context = crate::capabilities::ModelViewContext {
979 session_id,
980 prior_usage: prior_usage.as_ref(),
981 provider_managed_reduction: provider_managed,
982 };
983 let mut context_messages =
984 model_view_providers.apply_model_view(patched_messages, &model_view_context);
985 context_messages = crate::tool_call_integrity::retain_complete_message_tool_exchanges(
986 &context_messages,
987 stateful_response_continuation || restored_checkpoint.is_some(),
988 );
989
990 let render_facts = |at| {
994 crate::capabilities::render_facts_block(&crate::capabilities::collect_dynamic_facts(
995 &resolved_capability_configs,
996 &self.capability_registry,
997 Some(model_with_provider.model.as_str()),
998 &crate::capabilities::FactsContext::new(session_id).at(at),
999 ))
1000 };
1001 let (context_messages, volatile_suffix_len) =
1002 facts::interleave_facts(context_messages, render_facts, supports_clear_at);
1003 let mut context_messages = context_messages;
1004
1005 if let Some(context) = runtime_agent.conversation_context.as_ref()
1014 && !context.is_empty()
1015 {
1016 context_messages.insert(0, RuntimeMessage::user(context.clone()));
1017 }
1018
1019 let resolved_images = self.resolve_images(&context_messages).await;
1024 let resolved_files = self.resolve_files(&context_messages).await;
1025
1026 let mut llm_messages = Vec::new();
1028
1029 let has_system_prompt = !runtime_agent.system_prompt.is_empty();
1031 if has_system_prompt {
1032 llm_messages.push(Message {
1033 native_tool_calls: Vec::new(),
1034 role: MessageRole::System,
1035 content: MessageContent::Text(runtime_agent.system_prompt.clone()),
1036 tool_calls: None,
1037 tool_call_id: None,
1038 phase: None,
1039 reasoning: Vec::new(),
1040 configuration_update: None,
1041 });
1042 }
1043
1044 let messages_for_event: Vec<RuntimeMessage> = if has_system_prompt {
1046 std::iter::once(RuntimeMessage::system(&runtime_agent.system_prompt))
1047 .chain(context_messages.iter().cloned())
1048 .collect()
1049 } else {
1050 context_messages.clone()
1051 };
1052
1053 let mut stripped_error_count = 0u32;
1059 for msg in &context_messages {
1060 if is_error_placeholder_message(msg) {
1061 stripped_error_count += 1;
1062 continue;
1063 }
1064 let mut llm_msg = crate::llm_conversions::llm_message_from_message_with_attachments(
1065 msg,
1066 &resolved_images,
1067 &resolved_files,
1068 );
1069 llm_msg.configuration_update = reasoning_replay
1070 .as_ref()
1071 .and_then(|replay| replay.transitions.get(&msg.id).copied());
1072 if msg.role == RuntimeMessageRole::User
1073 && let Some(ref actor) = msg.external_actor
1074 {
1075 llm_msg.prepend_text_prefix(&format!("[{}] ", actor.display_label()));
1076 }
1077 facts::mark_turn_scoped(&mut llm_msg, msg, supports_clear_at);
1078 llm_messages.push(llm_msg);
1079 }
1080 if stripped_error_count > 0 {
1081 tracing::info!(
1082 session_id = %session_id,
1083 stripped_error_count,
1084 "ReasonAtom: stripped error placeholder messages from LLM input"
1085 );
1086 }
1087
1088 llm_messages = crate::tool_call_integrity::retain_complete_llm_tool_exchanges_for_request(
1093 llm_messages,
1094 stateful_response_continuation || restored_checkpoint.is_some(),
1095 );
1096
1097 let mut llm_config_builder =
1099 crate::llm_conversions::llm_call_config_builder_from_agent(&runtime_agent);
1100 if let Some(effort) = reasoning_effort {
1101 llm_config_builder = llm_config_builder.reasoning_effort(effort);
1102 }
1103 if let Some(speed) = speed {
1104 llm_config_builder = llm_config_builder.speed(speed);
1105 }
1106 if let Some(verbosity) = verbosity {
1107 llm_config_builder = llm_config_builder.verbosity(verbosity);
1108 }
1109
1110 for (k, v) in &embedder_metadata {
1112 llm_config_builder = llm_config_builder.with_metadata(k, v.clone());
1113 }
1114
1115 llm_config_builder = llm_config_builder
1119 .with_metadata("session_id", session_id.to_string())
1120 .with_metadata("harness_id", harness_id.to_string())
1121 .with_metadata("turn_id", context.turn_id.to_string())
1122 .with_metadata("exec_id", context.exec_id.to_string())
1123 .with_metadata("org_id", format!("org_{:032x}", org_id));
1124 if let Some(agent_id) = agent_id {
1125 llm_config_builder = llm_config_builder.with_metadata("agent_id", agent_id.to_string());
1126 }
1127
1128 if let Some(model_id) = &resolved_model_id {
1130 llm_config_builder = llm_config_builder.with_metadata("model_id", model_id.to_string());
1131 }
1132
1133 let mut llm_config = llm_config_builder
1134 .previous_response_id(previous_response_id.clone())
1135 .volatile_suffix_len(volatile_suffix_len)
1136 .build();
1137 llm_config.background_call = self.background_call.clone();
1138 if let Some(replay) = &reasoning_replay {
1139 llm_config.reasoning_effort = replay.state.baseline;
1140 llm_config.reasoning_state = Some(replay.state.clone());
1141 if replay.reset_continuation {
1142 llm_config.previous_response_id = None;
1143 }
1144 } else if messages
1145 .iter()
1146 .rev()
1147 .find(|message| {
1148 message.role == RuntimeMessageRole::Agent && !is_error_placeholder_message(message)
1149 })
1150 .and_then(|message| message.metadata.as_ref())
1151 .is_some_and(|metadata| metadata.contains_key(reasoning_updates::STATE_KEY))
1152 {
1153 llm_config.previous_response_id = None;
1156 }
1157 if let Some(checkpoint) = restored_checkpoint.as_ref()
1158 && let crate::CompactionCheckpointPayload::ProviderOpaque { context } =
1159 &checkpoint.payload
1160 {
1161 llm_config.previous_response_id = None;
1162 llm_config.provider_opaque_context = Some(context.clone());
1163 }
1164 provider_managed_compaction::decorate_request(
1165 &mut llm_config,
1166 provider_managed,
1167 restored_checkpoint.is_some(),
1168 );
1169 hosted_tools::ensure_hosted_tools_supported(
1170 &llm_config.driver_options,
1171 model_with_provider.provider_type.as_str(),
1172 )?;
1173
1174 tracing::debug!(
1175 session_id = %session_id,
1176 turn_id = %context.turn_id,
1177 model = %runtime_agent.model,
1178 message_count = %llm_messages.len(),
1179 "ReasonAtom: calling LLM"
1180 );
1181
1182 let streaming_event_context = EventContext::from_execution_context(context);
1185
1186 let mut armed_guardrails: Vec<ArmedGuardrail> = Vec::new();
1193 for (cap_id, cfg, provider) in &guardrail_providers {
1194 let ctx = OutputGuardrailContext {
1195 system_prompt: &runtime_agent.system_prompt,
1196 config: cfg,
1197 };
1198 let guardrail_id = provider.id().to_string();
1199 if let Some(run) = provider.arm(&ctx) {
1200 armed_guardrails.push(ArmedGuardrail {
1201 capability_id: cap_id.clone(),
1202 guardrail_id,
1203 run,
1204 });
1205 }
1206 }
1207 let buffer_output_deltas = !post_output_providers.is_empty();
1212 let output_message_id = MessageId::new();
1216 tracing::info!(
1217 session_id = %session_id,
1218 turn_id = %context.turn_id,
1219 "ReasonAtom: emitting output.message.started event"
1220 );
1221 if let Err(e) = self
1222 .event_emitter
1223 .emit(EventRequest::new(
1224 session_id,
1225 streaming_event_context.clone(),
1226 OutputMessageStartedData {
1227 reasoning_state: llm_config.reasoning_state.clone(),
1228 turn_id: context.turn_id,
1229 message_id: output_message_id,
1230 model: Some(runtime_agent.model.clone()),
1231 iteration: Some(iteration),
1232 phase: None,
1235 },
1236 ))
1237 .await
1238 {
1239 if llm_config.reasoning_state.is_some() {
1240 return Err(e);
1241 }
1242 tracing::warn!(
1243 session_id = %session_id,
1244 error = %e,
1245 "ReasonAtom: failed to emit output.message.started event"
1246 );
1247 } else {
1248 tracing::info!(
1249 session_id = %session_id,
1250 "ReasonAtom: output.message.started event emitted successfully"
1251 );
1252 }
1253
1254 let thinking_enabled = reasoning_effort.is_some();
1256 if thinking_enabled {
1257 tracing::info!(
1258 session_id = %session_id,
1259 turn_id = %context.turn_id,
1260 "ReasonAtom: emitting reason.thinking.started event"
1261 );
1262 if let Err(e) = self
1263 .event_emitter
1264 .emit(EventRequest::new(
1265 session_id,
1266 streaming_event_context.clone(),
1267 ReasonThinkingStartedData {
1268 turn_id: context.turn_id,
1269 model: Some(runtime_agent.model.clone()),
1270 },
1271 ))
1272 .await
1273 {
1274 tracing::warn!(
1275 session_id = %session_id,
1276 error = %e,
1277 "ReasonAtom: failed to emit reason.thinking.started event"
1278 );
1279 } else {
1280 tracing::info!(
1281 session_id = %session_id,
1282 "ReasonAtom: reason.thinking.started event emitted successfully"
1283 );
1284 }
1285 }
1286
1287 let llm_start = Instant::now();
1289
1290 let mut compaction_info: Option<LlmCompactionInfo> = None;
1294 let mut llm_messages_for_call = llm_messages.clone();
1295 let compaction_lifecycle = provider_managed_compaction::Lifecycle::new(
1296 self.event_emitter.as_ref(),
1297 session_id,
1298 &streaming_event_context,
1299 &model_with_provider.model,
1300 model_with_provider.provider_type.as_str(),
1301 llm_messages_for_call.len(),
1302 message_source_sequence,
1303 &llm_config,
1304 );
1305
1306 if !provider_managed && let Some(policy) = compaction_policy.as_deref() {
1307 compaction_info = apply_proactive_compaction(
1308 ProactiveCompactionContext {
1309 chat_driver: chat_driver.as_ref(),
1310 policy,
1311 checkpoint_store: self.compaction_checkpoint_store.as_ref(),
1312 event_emitter: self.event_emitter.as_ref(),
1313 event_context: &streaming_event_context,
1314 session_id,
1315 message_source_sequence,
1316 provider_type: model_with_provider.provider_type.as_str(),
1317 model: &model_with_provider.model,
1318 system_prompt: has_system_prompt
1319 .then_some(runtime_agent.system_prompt.as_str()),
1320 stateful_response_continuation,
1321 checkpoint_restored: restored_checkpoint.is_some(),
1322 checkpoint_suffix_message_count,
1323 raw_tool_result_bytes,
1324 prior_usage: prior_usage.as_ref(),
1325 },
1326 &mut llm_messages_for_call,
1327 &mut llm_config,
1328 )
1329 .await?;
1330 }
1331
1332 const DELTA_BATCH_INTERVAL_MS: u64 = 100;
1335 let retry_config = self.provider_retry_config.clone();
1336 let has_provider_executed_tools =
1337 hosted_tools::has_provider_executed_tools(&llm_config.driver_options);
1338 let mut stream_retry_metadata = RetryMetadata::default();
1339 let mut retry_started_at = None;
1340 let mut streamed_phase: Option<everruns_provider::ExecutionPhase> = None;
1345 let mut native_calls = std::collections::BTreeMap::new();
1346 let mut compaction_started_at: Option<Instant> = None;
1347 let (
1348 text,
1349 thinking,
1350 reasoning,
1351 tool_calls,
1352 completion_metadata,
1353 time_to_first_token_ms,
1354 pending_delta,
1355 mut tripped,
1356 ) = 'stream_attempt: loop {
1357 let stream_result = if let Some(remaining) =
1358 remaining_retry_time(&retry_config, retry_started_at)
1359 {
1360 match tokio::time::timeout(
1361 remaining,
1362 chat_driver.chat_completion_stream(
1363 &crate::ProviderEndpoint::default(),
1364 llm_messages_for_call.clone(),
1365 &llm_config,
1366 ),
1367 )
1368 .await
1369 {
1370 Ok(result) => result,
1371 Err(_) => {
1372 compaction_lifecycle.fail_if(provider_managed).await;
1373 return Err(AgentLoopError::llm_kind(
1374 crate::error::LlmErrorKind::Unavailable,
1375 format!(
1376 "provider retry time budget exhausted after {} retries over {:.1}s; the turn is safe to resume",
1377 stream_retry_metadata.attempts,
1378 retry_config.max_retry_elapsed.as_secs_f64()
1379 ),
1380 )
1381 .with_retry_metadata(&stream_retry_metadata));
1382 }
1383 }
1384 } else {
1385 chat_driver
1386 .chat_completion_stream(
1387 &crate::ProviderEndpoint::default(),
1388 llm_messages_for_call.clone(),
1389 &llm_config,
1390 )
1391 .await
1392 };
1393 let mut stream = match stream_result {
1394 Ok(stream) => stream,
1395 Err(e) if e.is_request_too_large() => {
1396 compaction_lifecycle.fail_if(provider_managed).await;
1397 if provider_managed {
1398 return Err(e);
1399 }
1400 let Some(policy) = compaction_policy.as_deref() else {
1401 tracing::warn!(
1402 session_id = %session_id,
1403 turn_id = %context.turn_id,
1404 "ReasonAtom: context too large and compaction capability is not enabled"
1405 );
1406 return Err(e);
1407 };
1408 let outcome = apply_reactive_compaction(
1409 ReactiveCompactionContext {
1410 chat_driver: chat_driver.as_ref(),
1411 policy,
1412 checkpoint_store: self.compaction_checkpoint_store.as_ref(),
1413 event_emitter: self.event_emitter.as_ref(),
1414 event_context: &streaming_event_context,
1415 session_id,
1416 message_source_sequence,
1417 provider_type: model_with_provider.provider_type.as_str(),
1418 model: &model_with_provider.model,
1419 summarization_model_fallback: &runtime_agent.model,
1420 system_prompt: has_system_prompt
1421 .then_some(runtime_agent.system_prompt.as_str()),
1422 stateful_response_continuation,
1423 },
1424 &mut llm_messages_for_call,
1425 &mut llm_config,
1426 )
1427 .await?;
1428 let Some(outcome) = outcome else {
1429 return Err(e);
1430 };
1431 if outcome.generation_info.is_some() {
1432 compaction_info = outcome.generation_info;
1433 }
1434
1435 chat_driver
1436 .chat_completion_stream(
1437 &crate::ProviderEndpoint::default(),
1438 llm_messages_for_call.clone(),
1439 &llm_config,
1440 )
1441 .await?
1442 }
1443 Err(e)
1444 if e.is_transient_llm_error()
1445 && !e.llm_retry_handled()
1446 && !has_provider_executed_tools
1447 && stream_retry_metadata.attempts < retry_config.max_retries =>
1448 {
1449 let proposed_wait =
1450 retry_config.calculate_backoff(stream_retry_metadata.attempts);
1451 let Some(wait_duration) =
1452 reserve_retry_wait(&retry_config, &mut retry_started_at, proposed_wait)
1453 else {
1454 compaction_lifecycle.fail_if(provider_managed).await;
1455 return Err(AgentLoopError::llm_kind(
1456 e.llm_error_kind()
1457 .unwrap_or(crate::error::LlmErrorKind::Unavailable),
1458 format!(
1459 "{e}; automatic recovery time budget exhausted after {} retries; the turn is safe to resume",
1460 stream_retry_metadata.attempts
1461 ),
1462 )
1463 .with_retry_metadata(&stream_retry_metadata));
1464 };
1465 tracing::warn!(
1466 session_id = %session_id,
1467 turn_id = %context.turn_id,
1468 attempt = stream_retry_metadata.attempts + 1,
1469 max_retries = retry_config.max_retries,
1470 wait_secs = wait_duration.as_secs_f64(),
1471 error = %e,
1472 "ReasonAtom: transient provider failure before stream, retrying"
1473 );
1474 stream_retry_metadata.record_retry(wait_duration, None);
1475 tokio::time::sleep(wait_duration).await;
1476 continue 'stream_attempt;
1477 }
1478 Err(error) => {
1479 return compaction_lifecycle
1480 .fail_start(error, provider_managed)
1481 .await;
1482 }
1483 };
1484
1485 if let Some(coordinator) = &self.native_async {
1486 coordinator
1487 .lock()
1488 .await
1489 .begin_transcript_response(output_message_id.to_string())
1490 .await?;
1491 let coordinator = coordinator.clone();
1492 stream = Box::pin(futures::stream::unfold(
1493 Some((coordinator, stream)),
1494 |state| async move {
1495 let (coordinator, mut source) = state?;
1496 let event = coordinator
1497 .lock()
1498 .await
1499 .next_response_event(&mut source)
1500 .await;
1501 let finished = matches!(&event, Ok(LlmStreamEvent::Done(_)) | Err(_));
1502 Some((event, (!finished).then_some((coordinator, source))))
1503 },
1504 ));
1505 }
1506 let mut text = String::new();
1507 let mut reasoning: Vec<ReasoningContentPart> = Vec::new();
1511 let mut thinking = String::new();
1513 let mut tool_calls = Vec::new();
1514 let mut termination = StreamTermination::Exhausted;
1515 let mut replay_state = StreamReplayState::for_request(has_provider_executed_tools);
1516 let mut pending_delta = String::new();
1517 let mut pending_thinking_delta = String::new();
1518 let mut last_delta_emit = Instant::now();
1519 let mut last_thinking_delta_emit = Instant::now();
1520 let mut time_to_first_token_ms: Option<u64> = None;
1521
1522 let stall_timeout = self
1524 .provider_stall_timeout
1525 .unwrap_or(std::time::Duration::from_secs(120));
1526 let initial_stall_timeout = remaining_retry_time(&retry_config, retry_started_at)
1527 .map_or(stall_timeout, |remaining| remaining.min(stall_timeout));
1528 let mut stall_sleep = Box::pin(tokio::time::sleep(initial_stall_timeout));
1529 let mut keepalive_ticker = tokio::time::interval(std::time::Duration::from_secs(12));
1530 keepalive_ticker.set_missed_tick_behavior(tokio::time::MissedTickBehavior::Skip);
1531 keepalive_ticker.tick().await; let mut last_stream_heartbeat = Instant::now();
1533 let mut last_token_at_unix: u64 = unix_now_secs();
1538
1539 loop {
1540 let event = tokio::select! {
1541 biased;
1542 next = stream.next() => match next {
1543 Some(e) => e,
1544 None => break,
1545 },
1546 _ = &mut stall_sleep => {
1547 let stall_error =
1557 crate::driver_registry::LlmStreamError::new(format!(
1558 "provider stream stall: no tokens for {}s",
1559 stall_timeout.as_secs()
1560 ));
1561 tracing::warn!(
1562 session_id = %session_id,
1563 turn_id = %context.turn_id,
1564 stall_secs = stall_timeout.as_secs(),
1565 "ReasonAtom: provider stream stall timeout"
1566 );
1567 if replay_state.should_retry(
1568 &stall_error,
1569 stream_retry_metadata.attempts,
1570 retry_config.max_retries,
1571 ) {
1572 let proposed_wait = retry_config
1573 .calculate_backoff(stream_retry_metadata.attempts);
1574 let Some(wait_duration) = reserve_retry_wait(
1575 &retry_config,
1576 &mut retry_started_at,
1577 proposed_wait,
1578 ) else {
1579 compaction_lifecycle
1580 .fail_if(provider_managed)
1581 .await;
1582 return Err(AgentLoopError::llm_kind(
1583 crate::error::LlmErrorKind::Unavailable,
1584 format!(
1585 "{}; automatic recovery time budget exhausted after {} retries; the turn is safe to resume",
1586 stall_error.message,
1587 stream_retry_metadata.attempts
1588 ),
1589 )
1590 .with_retry_metadata(&stream_retry_metadata));
1591 };
1592 tracing::warn!(
1593 session_id = %session_id,
1594 turn_id = %context.turn_id,
1595 attempt = stream_retry_metadata.attempts + 1,
1596 max_retries = retry_config.max_retries,
1597 wait_secs = wait_duration.as_secs_f64(),
1598 "ReasonAtom: provider stream stall, retrying"
1599 );
1600 stream_retry_metadata.record_retry(wait_duration, None);
1601 tokio::time::sleep(wait_duration).await;
1602 continue 'stream_attempt;
1603 }
1604 compaction_lifecycle
1605 .fail_if(compaction_started_at.is_some())
1606 .await;
1607 return Err(AgentLoopError::llm(stall_error.message));
1608 },
1609 _ = keepalive_ticker.tick() => {
1610 if let Some(ref hb) = self.stream_heartbeater {
1611 hb.heartbeat(crate::durability::StreamProgress {
1612 accumulated_len: text.len() + thinking.len(),
1613 last_delta_at: last_token_at_unix,
1614 })
1615 .await;
1616 last_stream_heartbeat = Instant::now();
1617 }
1618 continue;
1619 },
1620 };
1621 let event = match event {
1622 Ok(event) => event,
1623 Err(error) => {
1624 compaction_lifecycle
1625 .fail_if(compaction_started_at.is_some())
1626 .await;
1627 return Err(error);
1628 }
1629 };
1630 replay_state.observe(&event);
1631 let advanced_stall_deadline = advances_stall_deadline(&event);
1632 if advanced_stall_deadline {
1633 stall_sleep
1634 .as_mut()
1635 .reset(tokio::time::Instant::now() + stall_timeout);
1636 last_token_at_unix = unix_now_secs();
1637 }
1638 match event {
1639 LlmStreamEvent::TextDelta(delta) => {
1640 if delta.is_empty() {
1641 continue;
1642 }
1643 if time_to_first_token_ms.is_none() {
1645 let ttft = llm_start.elapsed().as_millis() as u64;
1646 time_to_first_token_ms = Some(ttft);
1647 tracing::info!(
1648 session_id = %session_id,
1649 time_to_first_token_ms = ttft,
1650 "ReasonAtom: received first token from LLM"
1651 );
1652 }
1653 text.push_str(&delta);
1654 pending_delta.push_str(&delta);
1655
1656 if !armed_guardrails.is_empty()
1663 && let Some(t) =
1664 evaluate_guardrails(&mut armed_guardrails, &text, &delta)
1665 {
1666 tracing::warn!(
1667 session_id = %session_id,
1668 turn_id = %context.turn_id,
1669 guardrail_capability_id = %t.capability_id,
1670 guardrail_id = %t.guardrail_id,
1671 reason_code = %t.block.reason_code,
1672 "ReasonAtom: output guardrail tripped, replacing assistant message"
1673 );
1674 pending_delta.clear();
1675 termination = StreamTermination::GuardrailBlocked(t);
1676 break;
1677 }
1678
1679 if !buffer_output_deltas
1681 && last_delta_emit.elapsed().as_millis() as u64
1682 >= DELTA_BATCH_INTERVAL_MS
1683 && !pending_delta.is_empty()
1684 {
1685 if let Err(e) = self
1686 .event_emitter
1687 .emit(EventRequest::new(
1688 session_id,
1689 streaming_event_context.clone(),
1690 OutputMessageDeltaData {
1691 turn_id: context.turn_id,
1692 message_id: output_message_id,
1693 delta: pending_delta.clone(),
1694 accumulated: text.clone(),
1695 phase: streamed_phase,
1696 },
1697 ))
1698 .await
1699 {
1700 tracing::warn!(
1701 session_id = %session_id,
1702 error = %e,
1703 "ReasonAtom: failed to emit output.message.delta event"
1704 );
1705 }
1706 pending_delta.clear();
1707 last_delta_emit = Instant::now();
1708 }
1709 }
1710 LlmStreamEvent::ReasoningDelta { delta, summary: _ } => {
1711 if delta.is_empty() {
1712 continue;
1713 }
1714 if let Some(t) = append_guarded_thinking_delta(
1715 &mut armed_guardrails,
1716 &mut thinking,
1717 &mut pending_thinking_delta,
1718 &delta,
1719 ) {
1720 tracing::warn!(
1721 session_id = %session_id,
1722 guardrail_capability_id = %t.capability_id,
1723 guardrail_id = %t.guardrail_id,
1724 "ReasonAtom: output guardrail tripped on thinking stream, replacing assistant message"
1725 );
1726 termination = StreamTermination::GuardrailBlocked(t);
1727 break;
1728 }
1729 tracing::debug!(
1730 session_id = %session_id,
1731 delta_len = delta.len(),
1732 total_thinking_len = thinking.len(),
1733 "ReasonAtom: received ThinkingDelta from LLM"
1734 );
1735
1736 if last_thinking_delta_emit.elapsed().as_millis() as u64
1738 >= DELTA_BATCH_INTERVAL_MS
1739 && !pending_thinking_delta.is_empty()
1740 {
1741 if let Err(e) = self
1742 .event_emitter
1743 .emit(EventRequest::new(
1744 session_id,
1745 streaming_event_context.clone(),
1746 ReasonThinkingDeltaData {
1747 turn_id: context.turn_id,
1748 delta: pending_thinking_delta.clone(),
1749 accumulated: thinking.clone(),
1750 },
1751 ))
1752 .await
1753 {
1754 tracing::warn!(
1755 session_id = %session_id,
1756 error = %e,
1757 "ReasonAtom: failed to emit reason.thinking.delta event"
1758 );
1759 }
1760 pending_thinking_delta.clear();
1761 last_thinking_delta_emit = Instant::now();
1762 }
1763 }
1764 LlmStreamEvent::ReasoningItem(item) => {
1765 if let Some(t) = inspect_guarded_reasoning_item(
1766 &mut armed_guardrails,
1767 &mut thinking,
1768 &item,
1769 ) {
1770 tracing::warn!(
1771 session_id = %session_id,
1772 guardrail_capability_id = %t.capability_id,
1773 guardrail_id = %t.guardrail_id,
1774 "ReasonAtom: output guardrail tripped on completed reasoning item, replacing assistant message"
1775 );
1776 termination = StreamTermination::GuardrailBlocked(t);
1777 break;
1778 }
1779 tracing::debug!(
1783 session_id = %session_id,
1784 provider = %item.provider,
1785 item_id = ?item.item_id,
1786 has_signature = item.signature.is_some(),
1787 has_encrypted = item.encrypted.is_some(),
1788 "ReasonAtom: captured reasoning artifact"
1789 );
1790 reasoning.push(item);
1791 }
1792 LlmStreamEvent::NativeToolCall(call) => {
1793 if self.native_async.is_none() {
1794 return Err(AgentLoopError::config(
1795 "native async/custom tools require a configured native-call coordinator",
1796 ));
1797 }
1798 let part = crate::message::ToolCallContentPart::from_native(call.clone())?;
1799 if native_calls.insert(call.id().to_owned(), call).is_none() {
1800 tool_calls.push(ToolCall {
1801 id: part.id,
1802 name: part.name,
1803 arguments: part.arguments,
1804 });
1805 }
1806 }
1807 LlmStreamEvent::ToolCalls(calls) => {
1808 if self.native_async.is_some() {
1809 for call in &calls {
1810 native_calls.entry(call.id.clone()).or_insert_with(|| {
1811 everruns_provider::native_async::NativeToolCall::Function {
1812 call_id: call.id.clone(),
1813 name: call.name.clone(),
1814 arguments: call.arguments.to_string(),
1815 asynchronous: false,
1816 }
1817 });
1818 }
1819 for call in calls {
1820 if !tool_calls.iter().any(|existing| existing.id == call.id) {
1821 tool_calls.push(call);
1822 }
1823 }
1824 } else {
1825 tool_calls = calls;
1826 }
1827 }
1828 LlmStreamEvent::MessagePhase(phase) => {
1829 streamed_phase = everruns_provider::ExecutionPhase::refine_streamed_hint(
1838 streamed_phase,
1839 phase,
1840 );
1841 }
1842 LlmStreamEvent::ProviderCompactionStarted => {
1843 compaction_lifecycle.start(&mut compaction_started_at).await;
1844 }
1845 LlmStreamEvent::HostedToolCall(call) => {
1846 let request = hosted_tools::hosted_call_event(
1847 session_id,
1848 &streaming_event_context,
1849 context.turn_id,
1850 call,
1851 );
1852 if let Err(error) = self.event_emitter.emit(request).await {
1853 tracing::warn!(%session_id, %error, "ReasonAtom: failed to emit tool.hosted_call");
1854 }
1855 }
1856 LlmStreamEvent::Done(metadata) => {
1857 if !buffer_output_deltas
1861 && !pending_delta.is_empty()
1862 && let Err(e) = self
1863 .event_emitter
1864 .emit(EventRequest::new(
1865 session_id,
1866 streaming_event_context.clone(),
1867 OutputMessageDeltaData {
1868 turn_id: context.turn_id,
1869 message_id: output_message_id,
1870 delta: pending_delta.clone(),
1871 accumulated: text.clone(),
1872 phase: streamed_phase,
1873 },
1874 ))
1875 .await
1876 {
1877 tracing::warn!(
1878 session_id = %session_id,
1879 error = %e,
1880 "ReasonAtom: failed to emit final output.message.delta event"
1881 );
1882 }
1883
1884 if !pending_thinking_delta.is_empty()
1886 && let Err(e) = self
1887 .event_emitter
1888 .emit(EventRequest::new(
1889 session_id,
1890 streaming_event_context.clone(),
1891 ReasonThinkingDeltaData {
1892 turn_id: context.turn_id,
1893 delta: pending_thinking_delta.clone(),
1894 accumulated: thinking.clone(),
1895 },
1896 ))
1897 .await
1898 {
1899 tracing::warn!(
1900 session_id = %session_id,
1901 error = %e,
1902 "ReasonAtom: failed to emit final reason.thinking.delta event"
1903 );
1904 }
1905
1906 if !thinking.is_empty()
1908 && let Err(e) = self
1909 .event_emitter
1910 .emit(EventRequest::new(
1911 session_id,
1912 streaming_event_context.clone(),
1913 ReasonThinkingCompletedData {
1914 turn_id: context.turn_id,
1915 thinking: thinking.clone(),
1916 },
1917 ))
1918 .await
1919 {
1920 tracing::warn!(
1921 session_id = %session_id,
1922 error = %e,
1923 "ReasonAtom: failed to emit reason.thinking.completed event"
1924 );
1925 }
1926 termination = StreamTermination::Completed(metadata);
1927 break;
1928 }
1929 LlmStreamEvent::Error(err) => {
1930 let has_partial_output = compaction_started_at.is_none()
1935 && (!tool_calls.is_empty() || !text.is_empty());
1936
1937 if has_partial_output {
1938 tracing::warn!(
1939 session_id = %session_id,
1940 error = %err,
1941 tool_call_count = tool_calls.len(),
1942 text_len = text.len(),
1943 "ReasonAtom: trailing stream error after valid output — treating as partial success"
1944 );
1945 termination = StreamTermination::PartialSuccess;
1949 break;
1950 }
1951
1952 if replay_state.should_retry(
1953 &err,
1954 stream_retry_metadata.attempts,
1955 retry_config.max_retries,
1956 ) {
1957 let proposed_wait =
1958 retry_config.calculate_backoff(stream_retry_metadata.attempts);
1959 let Some(wait_duration) = reserve_retry_wait(
1960 &retry_config,
1961 &mut retry_started_at,
1962 proposed_wait,
1963 ) else {
1964 return Err(AgentLoopError::llm_kind(
1965 err.kind(),
1966 format!(
1967 "{err}; automatic recovery time budget exhausted after {} retries; the turn is safe to resume",
1968 stream_retry_metadata.attempts
1969 ),
1970 )
1971 .with_retry_metadata(&stream_retry_metadata));
1972 };
1973 tracing::warn!(
1974 session_id = %session_id,
1975 turn_id = %context.turn_id,
1976 attempt = stream_retry_metadata.attempts + 1,
1977 max_retries = retry_config.max_retries,
1978 wait_secs = wait_duration.as_secs_f64(),
1979 error_code = err.code.as_deref().unwrap_or("none"),
1980 error_status = err.status,
1981 error = %err,
1982 "ReasonAtom: transient stream error before output, retrying"
1983 );
1984 stream_retry_metadata.record_retry(wait_duration, None);
1985 tokio::time::sleep(wait_duration).await;
1986 continue 'stream_attempt;
1987 }
1988
1989 let llm_duration_ms = llm_start.elapsed().as_millis() as u64;
1991 let event_context = EventContext::from_execution_context(context)
1992 .with_span(
1993 trace_id.to_string(),
1994 Uuid::now_v7().to_string(),
1995 Some(reason_span_id.to_string()),
1996 );
1997 let tools_summary: Vec<ToolDefinitionSummary> =
1998 runtime_agent.tools.iter().map(|t| t.into()).collect();
1999 let generation_data = LlmGenerationData::failure(
2000 messages_for_event.clone(),
2001 tools_summary,
2002 runtime_agent.model.clone(),
2003 Some(model_with_provider.provider_type.to_string()),
2004 err.to_string(),
2005 Some(llm_duration_ms),
2006 time_to_first_token_ms,
2007 );
2008 let _ = self
2009 .event_emitter
2010 .emit(EventRequest::new(
2011 session_id,
2012 event_context,
2013 generation_data,
2014 ))
2015 .await;
2016 if compaction_started_at.is_some() {
2017 compaction_lifecycle.fail().await;
2018 }
2019 return Err(AgentLoopError::llm_kind(err.kind(), err.to_string()));
2020 }
2021 _ => {}
2027 }
2028 if last_stream_heartbeat.elapsed().as_millis() as u64 >= 5_000
2031 && let Some(ref hb) = self.stream_heartbeater
2032 {
2033 hb.heartbeat(crate::durability::StreamProgress {
2034 accumulated_len: text.len() + thinking.len(),
2035 last_delta_at: last_token_at_unix,
2036 })
2037 .await;
2038 last_stream_heartbeat = Instant::now();
2039 }
2040 }
2041 compaction_lifecycle
2042 .reject_incomplete(compaction_started_at, &termination)
2043 .await?;
2044 let (mut completion_metadata, tripped) = termination.into_parts();
2045 if let Some(metadata) = completion_metadata.as_mut() {
2046 metadata.retry_metadata =
2047 merge_retry_metadata(metadata.retry_metadata.take(), &stream_retry_metadata);
2048 }
2049
2050 break 'stream_attempt (
2051 text,
2052 thinking,
2053 reasoning,
2054 tool_calls,
2055 completion_metadata,
2056 time_to_first_token_ms,
2057 pending_delta,
2058 tripped,
2059 );
2060 };
2061 let (mut text, mut thinking, mut reasoning, mut tool_calls) =
2062 (text, thinking, reasoning, tool_calls);
2063 compaction_lifecycle.record_observed(&mut llm_config, compaction_started_at);
2064
2065 let mut citation_annotations: Vec<crate::message::TextAnnotation> = Vec::new();
2076 if tripped.is_none()
2077 && !annotation_providers.is_empty()
2078 && !text.is_empty()
2079 && tool_calls.is_empty()
2080 {
2081 text = filter_response_text(
2084 &self.capability_registry,
2085 &resolved_capability_configs,
2086 text,
2087 );
2088 let collected = collect_annotations(
2089 &annotation_providers,
2090 &runtime_agent.system_prompt,
2091 &text,
2092 &messages,
2093 self.utility_llm_service.as_ref(),
2094 )
2095 .await;
2096 text = collected.text;
2097 citation_annotations = collected.annotations;
2098
2099 if !citation_annotations.is_empty() && !post_output_providers.is_empty() {
2102 let guarded_output = client_visible_guardrail_text(
2103 &text,
2104 &thinking,
2105 &reasoning,
2106 &citation_annotations,
2107 );
2108 let ctx = PostGenerationOutputContext {
2109 system_prompt: &runtime_agent.system_prompt,
2110 message_text: &guarded_output,
2111 utility_llm_service: self.utility_llm_service.as_ref(),
2112 decisions: self.decisions.as_ref(),
2113 };
2114 tripped = evaluate_post_generation_guardrails(&post_output_providers, &ctx).await;
2115 }
2116
2117 if tripped.is_none()
2120 && !citation_annotations.is_empty()
2121 && !citation_verifiers.is_empty()
2122 {
2123 citation_annotations = verify_annotations(
2124 &citation_verifiers,
2125 &text,
2126 self.utility_llm_service.as_ref(),
2127 citation_annotations,
2128 )
2129 .await;
2130 }
2131 }
2132
2133 if tripped.is_none()
2136 && citation_annotations.is_empty()
2137 && !post_output_providers.is_empty()
2138 && (!text.is_empty() || !thinking.is_empty() || !reasoning.is_empty())
2139 {
2140 let guarded_output = client_visible_guardrail_text(&text, &thinking, &reasoning, &[]);
2141 let ctx = PostGenerationOutputContext {
2142 system_prompt: &runtime_agent.system_prompt,
2143 message_text: &guarded_output,
2144 utility_llm_service: self.utility_llm_service.as_ref(),
2145 decisions: self.decisions.as_ref(),
2146 };
2147 tripped = evaluate_post_generation_guardrails(&post_output_providers, &ctx).await;
2148 }
2149
2150 if tripped.is_some() {
2151 citation_annotations.clear();
2152 }
2153
2154 if tripped.is_none() {
2158 for item in &reasoning {
2159 if let Err(e) = self
2160 .event_emitter
2161 .emit(EventRequest::new(
2162 session_id,
2163 streaming_event_context.clone(),
2164 ReasonItemData {
2165 turn_id: context.turn_id,
2166 provider: item.provider.clone(),
2167 model: Some(llm_config.model.clone()),
2168 item_id: item.item_id.clone().unwrap_or_default(),
2169 summary: item
2170 .display_text()
2171 .filter(|_| !matches!(item.text, Some(ReasoningText::Plain { .. })))
2172 .into_iter()
2173 .collect(),
2174 token_count: item.tokens,
2175 },
2176 ))
2177 .await
2178 {
2179 tracing::warn!(
2180 session_id = %session_id,
2181 error = %e,
2182 "ReasonAtom: failed to emit reason.item event"
2183 );
2184 }
2185 }
2186 }
2187
2188 if buffer_output_deltas
2191 && tripped.is_none()
2192 && !pending_delta.is_empty()
2193 && let Err(e) = self
2194 .event_emitter
2195 .emit(EventRequest::new(
2196 session_id,
2197 streaming_event_context.clone(),
2198 OutputMessageDeltaData {
2199 turn_id: context.turn_id,
2200 message_id: output_message_id,
2201 delta: pending_delta.clone(),
2202 accumulated: text.clone(),
2203 phase: streamed_phase,
2204 },
2205 ))
2206 .await
2207 {
2208 tracing::warn!(
2209 session_id = %session_id,
2210 error = %e,
2211 "ReasonAtom: failed to emit guarded output.message.delta event"
2212 );
2213 }
2214
2215 if let Some(ref t) = tripped {
2221 let replaced_event_context = EventContext::from_execution_context(context).with_span(
2222 trace_id.to_string(),
2223 Uuid::now_v7().to_string(),
2224 Some(reason_span_id.to_string()),
2225 );
2226 if let Err(e) = self
2227 .event_emitter
2228 .emit(EventRequest::new(
2229 session_id,
2230 replaced_event_context,
2231 OutputMessageReplacedData {
2232 turn_id: context.turn_id,
2233 message_id: output_message_id,
2234 guardrail_capability_id: t.capability_id.clone(),
2235 guardrail_id: t.guardrail_id.clone(),
2236 reason_code: t.block.reason_code.clone(),
2237 replacement: t.block.replacement.clone(),
2238 },
2239 ))
2240 .await
2241 {
2242 tracing::warn!(
2243 session_id = %session_id,
2244 error = %e,
2245 "ReasonAtom: failed to emit output.message.replaced event"
2246 );
2247 }
2248 text = t.block.replacement.clone();
2249 tool_calls.clear();
2250 thinking.clear();
2251 reasoning.clear();
2252 }
2253
2254 let rejected_tool_calls = if tool_calls.is_empty() {
2258 Vec::new()
2259 } else {
2260 finalized_calls::apply_finalized_tool_calls_hooks(
2261 &self.capability_registry,
2262 self.event_emitter.as_ref(),
2263 session_id,
2264 context,
2265 &resolved_capability_configs,
2266 &runtime_agent.tools,
2267 &mut tool_calls,
2268 iteration,
2269 )
2270 .await
2271 };
2272 let finalized_tool_calls = tool_calls.clone();
2273 let rejected_tool_call_ids: HashSet<_> = rejected_tool_calls
2274 .iter()
2275 .map(|rejection| rejection.tool_call_id.clone())
2276 .collect();
2277 tool_calls.retain(|call| !rejected_tool_call_ids.contains(&call.id));
2278
2279 let llm_duration_ms = llm_start.elapsed().as_millis() as u64;
2280
2281 let response_id = completion_metadata
2282 .as_ref()
2283 .and_then(|meta| meta.response_id.clone());
2284 let finish_reason = completion_metadata
2285 .as_ref()
2286 .and_then(|meta| meta.finish_reason.clone());
2287
2288 let usage = completion_metadata.as_ref().and_then(|meta| {
2290 hosted_tools::completion_usage(
2291 meta,
2292 &model_with_provider.provider_type,
2293 &runtime_agent.model,
2294 )
2295 });
2296
2297 let event_context = EventContext::from_execution_context(context).with_span(
2299 trace_id.to_string(),
2300 Uuid::now_v7().to_string(),
2301 Some(reason_span_id.to_string()),
2302 );
2303 let tools_summary: Vec<ToolDefinitionSummary> =
2304 runtime_agent.tools.iter().map(|t| t.into()).collect();
2305 let finish_reasons = Some(vec![finish_reason.clone().unwrap_or_else(|| {
2306 if finalized_tool_calls.is_empty() {
2307 "stop".to_string()
2308 } else {
2309 "tool_calls".to_string()
2310 }
2311 })]);
2312 let meta = completion_metadata.as_ref();
2313 let served = meta.and_then(|m| m.response_model.clone());
2314 let retry_info = completion_metadata
2315 .as_ref()
2316 .and_then(|meta| meta.retry_metadata.as_ref())
2317 .filter(|rm| rm.had_retries())
2318 .map(|rm| LlmRetryInfo {
2319 attempts: rm.attempts,
2320 total_wait_ms: rm.total_retry_wait.as_millis() as u64,
2321 });
2322 let mut generation_data = LlmGenerationData::success_with_retry(
2323 messages_for_event.clone(),
2324 tools_summary,
2325 Some(text.clone()).filter(|s| !s.is_empty()),
2326 finalized_tool_calls.clone(),
2327 runtime_agent.model.clone(),
2328 Some(model_with_provider.provider_type.to_string()),
2329 usage.clone(),
2330 Some(llm_duration_ms),
2331 time_to_first_token_ms,
2332 finish_reasons,
2333 response_id.clone(),
2334 retry_info,
2335 )
2336 .with_response_model(served);
2337
2338 if let Some(info) = compaction_info {
2344 if let Some(compaction_cost) = info.cost_usd {
2345 match generation_data.metadata.usage.as_mut() {
2346 Some(usage) => {
2347 add_compaction_cost(usage, compaction_cost);
2348 }
2349 None => {
2354 generation_data.metadata.usage = Some(crate::events::TokenUsage {
2355 input_tokens: 0,
2356 output_tokens: 0,
2357 cache_read_tokens: None,
2358 cache_creation_tokens: None,
2359 actual_cost_usd: Some(compaction_cost),
2360 estimated_cost_usd: None,
2361 effective_cost_usd: None,
2362 });
2363 }
2364 }
2365 }
2366 generation_data = generation_data.with_compaction(info);
2367 }
2368
2369 if let Some(request_options) =
2370 build_request_options(&llm_config, &model_with_provider.provider_type.to_string())
2371 {
2372 generation_data = generation_data.with_request_options(request_options);
2373 }
2374
2375 if let Err(e) = self
2376 .event_emitter
2377 .emit(EventRequest::new(
2378 session_id,
2379 event_context,
2380 generation_data,
2381 ))
2382 .await
2383 {
2384 tracing::warn!(
2385 session_id = %session_id,
2386 error = %e,
2387 "ReasonAtom: failed to emit llm.generation event"
2388 );
2389 }
2390
2391 let mut metadata = std::collections::HashMap::new();
2393 metadata.insert(
2394 "model".to_string(),
2395 serde_json::Value::String(runtime_agent.model.clone()),
2396 );
2397 if let Some(state) = &llm_config.reasoning_state {
2398 metadata.insert(
2399 reasoning_updates::STATE_KEY.to_string(),
2400 serde_json::json!(state),
2401 );
2402 }
2403 if let Some(effort) = llm_config
2404 .reasoning_state
2405 .as_ref()
2406 .and_then(|state| state.effective)
2407 .or(reasoning_effort)
2408 {
2409 metadata.insert(
2410 "reasoning_effort".to_string(),
2411 serde_json::Value::String(effort.as_str().to_string()),
2412 );
2413 }
2414 metadata.insert(
2421 "provider".to_string(),
2422 serde_json::Value::String(model_with_provider.provider_type.to_string()),
2423 );
2424 if let Some(ref rid) = response_id {
2425 metadata.insert(
2426 "response_id".to_string(),
2427 serde_json::Value::String(rid.clone()),
2428 );
2429 }
2430
2431 let text = filter_response_text(
2435 &self.capability_registry,
2436 &resolved_capability_configs,
2437 text,
2438 );
2439 let (provider_opaque_content, provider_checkpoint_candidate) =
2440 provider_managed_compaction::replay_artifacts(
2441 completion_metadata.as_ref(),
2442 tripped.is_none(),
2443 !rejected_tool_calls.is_empty(),
2444 );
2445 let has_tool_calls = !finalized_tool_calls.is_empty();
2446 let mut assistant_message = if has_tool_calls {
2447 RuntimeMessage::assistant_with_tools(&text, finalized_tool_calls.clone())
2448 } else {
2449 RuntimeMessage::assistant(&text)
2450 }
2451 .with_id(output_message_id);
2452 for part in &mut assistant_message.content {
2453 if let crate::message::ContentPart::ToolCall(call) = part {
2454 call.native = native_calls.get(&call.id).cloned();
2455 }
2456 }
2457 if !citation_annotations.is_empty() {
2460 for part in assistant_message.content.iter_mut() {
2461 if let crate::message::ContentPart::Text(t) = part {
2462 t.annotations = std::mem::take(&mut citation_annotations);
2463 break;
2464 }
2465 }
2466 }
2467 let provider_type_for_reasoning = model_with_provider.provider_type.to_string();
2471 let provider_phase = completion_metadata
2475 .as_ref()
2476 .and_then(|meta| meta.phase.as_deref())
2477 .and_then(everruns_provider::ExecutionPhase::from_provider_str);
2478 let (phase, phase_source) = match provider_phase {
2479 Some(phase) => (phase, everruns_provider::PhaseSource::Provider),
2480 None => (
2481 everruns_provider::ExecutionPhase::from_has_tool_calls(has_tool_calls),
2482 everruns_provider::PhaseSource::Derived,
2483 ),
2484 };
2485 assistant_message.phase = Some(phase);
2486 assistant_message.phase_source = Some(phase_source);
2487 assistant_message.metadata = Some(metadata);
2488 if reasoning.is_empty() && !thinking.is_empty() {
2497 reasoning.push(
2498 ReasoningContentPart::opaque(provider_type_for_reasoning.clone()).with_text(
2499 ReasoningText::Plain {
2500 text: thinking.clone(),
2501 },
2502 ),
2503 );
2504 }
2505 if !reasoning.is_empty() {
2506 let mut content = Vec::with_capacity(reasoning.len() + assistant_message.content.len());
2507 content.extend(reasoning.drain(..).map(ContentPart::Reasoning));
2508 content.append(&mut assistant_message.content);
2509 assistant_message.content = content;
2510 }
2511 if let Some(content) = provider_opaque_content {
2512 assistant_message
2513 .content
2514 .push(ContentPart::ProviderOpaque(content));
2515 }
2516 let message_event_context = EventContext::from_execution_context(context).with_span(
2519 trace_id.to_string(),
2520 Uuid::now_v7().to_string(),
2521 Some(reason_span_id.to_string()),
2522 );
2523 let mut output_message_data = OutputMessageCompletedData::new(assistant_message);
2524 if let Some(ref u) = usage {
2525 output_message_data = output_message_data.with_usage(u.clone());
2526 }
2527 let result = ReasonResult {
2528 native_counts: None,
2529 success: true,
2530 text,
2531 tool_calls,
2532 has_tool_calls,
2533 tool_definitions: runtime_agent.tools.clone(),
2534 max_iterations: runtime_agent.max_iterations,
2535 usage,
2536 output_message_id: Some(output_message_id),
2537 time_to_first_token_ms,
2538 response_id,
2539 finish_reason,
2540 locale: resolved_locale,
2541 network_access: runtime_agent.network_access.clone(),
2542 parallel_tool_calls: runtime_agent.parallel_tool_calls,
2543 ..ReasonResult::default()
2544 };
2545 if let Some(coordinator) = &self.native_async {
2546 coordinator
2547 .lock()
2548 .await
2549 .stage_transcript_result(
2550 serde_json::to_value(&result)
2551 .map_err(|error| AgentLoopError::store(error.to_string()))?,
2552 )
2553 .await?;
2554 }
2555 let completed_output_event = self
2556 .event_emitter
2557 .emit(EventRequest::new(
2558 session_id,
2559 message_event_context,
2560 output_message_data,
2561 ))
2562 .await?;
2563 compaction_lifecycle
2564 .finish_after_output(
2565 self.compaction_checkpoint_store.as_deref(),
2566 completed_output_event.sequence,
2567 provider_checkpoint_candidate,
2568 compaction_started_at,
2569 restored_checkpoint.is_some(),
2570 completion_metadata.as_ref(),
2571 )
2572 .await;
2573
2574 if let Some(coordinator) = &self.native_async {
2575 coordinator
2576 .lock()
2577 .await
2578 .transcript_committed(&output_message_id.to_string())
2579 .await?;
2580 }
2581 for rejection in rejected_tool_calls {
2582 let Some(call) = finalized_tool_calls
2583 .iter()
2584 .find(|call| call.id == rejection.tool_call_id)
2585 else {
2586 continue;
2587 };
2588 self.event_emitter
2589 .emit(EventRequest::new(
2590 session_id,
2591 EventContext::from_execution_context(context),
2592 ToolCompletedData::failure(
2593 call.id.clone(),
2594 call.name.clone(),
2595 "error".to_string(),
2596 rejection.error,
2597 None,
2598 ),
2599 ))
2600 .await?;
2601 }
2602 tracing::info!(
2603 session_id = %session_id,
2604 turn_id = %context.turn_id,
2605 has_tool_calls = %result.has_tool_calls,
2606 tool_count = %result.tool_calls.len(),
2607 "ReasonAtom: LLM call completed"
2608 );
2609
2610 Ok(result)
2611 }
2612
2613 async fn finalize_partial_stream(
2618 &self,
2619 session_id: SessionId,
2620 context: &ExecutionContext,
2621 partial: PartialStreamState,
2622 iteration: u32,
2623 runtime_agent: &crate::RuntimeAgent,
2624 resolved_capability_configs: &[crate::CapabilityRef],
2625 ) -> Result<ReasonResult> {
2626 let event_context = EventContext::from_execution_context(context);
2627 let turn_id = context.turn_id;
2628 let message_id = partial.message_id;
2629
2630 let _ = self
2632 .event_emitter
2633 .emit(EventRequest::new(
2634 session_id,
2635 event_context.clone(),
2636 OutputMessageStartedData {
2637 reasoning_state: partial.reasoning_state.clone(),
2638 turn_id,
2639 message_id,
2640 model: None,
2641 iteration: Some(iteration),
2642 phase: None,
2645 },
2646 ))
2647 .await;
2648
2649 let accumulated = filter_response_text(
2652 &self.capability_registry,
2653 resolved_capability_configs,
2654 partial.accumulated,
2655 );
2656 let mut assistant_message = RuntimeMessage::assistant(&accumulated).with_id(message_id);
2657 if let Some(state) = partial.reasoning_state {
2658 assistant_message.metadata = Some(HashMap::from([
2659 ("model".into(), serde_json::json!("gpt-6-astra")),
2660 ("provider".into(), serde_json::json!("openai")),
2661 (
2662 reasoning_updates::STATE_KEY.into(),
2663 serde_json::json!(state),
2664 ),
2665 (
2666 "reasoning_effort".into(),
2667 serde_json::json!(state.effective),
2668 ),
2669 ]));
2670 }
2671 let output_message_id = message_id;
2672 self.event_emitter
2673 .emit(EventRequest::new(
2674 session_id,
2675 event_context.clone(),
2676 OutputMessageCompletedData::new(assistant_message),
2677 ))
2678 .await?;
2679
2680 let accumulated_len = accumulated.len();
2682 let _ = self
2683 .event_emitter
2684 .emit(EventRequest::new(
2685 session_id,
2686 event_context.clone(),
2687 ReasonRecoveredData {
2688 turn_id,
2689 mode: RecoveryMode::Finalize,
2690 accumulated_len,
2691 },
2692 ))
2693 .await;
2694
2695 tracing::info!(
2696 session_id = %session_id,
2697 turn_id = %turn_id,
2698 accumulated_len,
2699 "ReasonAtom: finalized partial stream from persisted accumulated text"
2700 );
2701
2702 Ok(ReasonResult {
2703 native_counts: None,
2704 success: true,
2705 text: accumulated,
2706 tool_calls: vec![],
2707 has_tool_calls: false,
2708 tool_definitions: runtime_agent.tools.clone(),
2709 max_iterations: runtime_agent.max_iterations,
2710 error: None,
2711 user_facing_error: None,
2712 error_disclosure: None,
2713 usage: None,
2714 output_message_id: Some(output_message_id),
2715 time_to_first_token_ms: None,
2716 response_id: None,
2717 finish_reason: Some("stop".to_string()),
2718 ..ReasonResult::default()
2720 })
2721 }
2722
2723 async fn resolve_images(&self, messages: &[RuntimeMessage]) -> HashMap<Uuid, ResolvedImage> {
2734 let mut resolved = HashMap::new();
2735
2736 let resolver = match &self.image_resolver {
2738 Some(r) => r,
2739 None => return resolved,
2740 };
2741
2742 let image_ids: Vec<Uuid> = messages
2744 .iter()
2745 .flat_map(crate::llm_conversions::extract_image_file_ids)
2746 .collect::<std::collections::HashSet<_>>()
2747 .into_iter()
2748 .collect();
2749
2750 if image_ids.is_empty() {
2751 return resolved;
2752 }
2753
2754 tracing::debug!(
2755 image_count = image_ids.len(),
2756 "ReasonAtom: resolving image_file references"
2757 );
2758
2759 for image_id in image_ids {
2761 match resolver.resolve_image(image_id).await {
2762 Ok(Some(image)) => {
2763 resolved.insert(image_id, image);
2764 }
2765 Ok(None) => {
2766 tracing::warn!(
2767 image_id = %image_id,
2768 "ReasonAtom: image not found during resolution"
2769 );
2770 }
2771 Err(e) => {
2772 tracing::warn!(
2773 image_id = %image_id,
2774 error = %e,
2775 "ReasonAtom: failed to resolve image"
2776 );
2777 }
2778 }
2779 }
2780
2781 tracing::debug!(
2782 resolved_count = resolved.len(),
2783 "ReasonAtom: image resolution complete"
2784 );
2785
2786 resolved
2787 }
2788
2789 async fn resolve_files(&self, messages: &[RuntimeMessage]) -> HashMap<Uuid, ResolvedFile> {
2790 let Some(resolver) = &self.file_resolver else {
2791 return HashMap::new();
2792 };
2793
2794 let file_ids: Vec<Uuid> = messages
2795 .iter()
2796 .flat_map(crate::llm_conversions::extract_file_ids)
2797 .collect::<std::collections::HashSet<_>>()
2798 .into_iter()
2799 .collect();
2800
2801 if file_ids.is_empty() {
2802 return HashMap::new();
2803 }
2804
2805 match resolver.resolve_files(&file_ids).await {
2806 Ok(map) => map,
2807 Err(e) => {
2808 tracing::warn!(
2809 target: "reason",
2810 "ReasonAtom: file resolution failed: {e}"
2811 );
2812 HashMap::new()
2813 }
2814 }
2815 }
2816}
2817
2818#[cfg(test)]
2823mod tests;