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