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