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