1use serde::{Deserialize, Serialize};
31use std::collections::HashSet;
32use std::future::Future;
33use std::pin::Pin;
34use std::sync::Arc;
35use std::task::{Context, Poll};
36use web_time::Instant;
37
38use super::ExecutionContext;
39use super::act_hooks::{self, PostActHook};
40use super::tool_scheduler;
41use crate::error::Result;
42use crate::events::{
43 ActCompletedData, ActStartedData, EventContext, EventRequest, ToolCompletedData,
44 ToolStartedData,
45};
46use crate::message::ContentPart;
47use crate::phase_effects::{PhaseEffectEmitter, PhaseEffectSink};
48use crate::tool_fingerprint::{
49 tool_call_fingerprint, tool_error_fingerprint, tool_result_fingerprint,
50};
51use crate::tool_narration::{
52 GroupHeadlineAction, ToolNarrationContext, ToolNarrationPhase,
53 render_tool_narration_with_locale, summarize_group_actions, tool_call_for_group_summary,
54};
55use crate::tool_types::ConnectionRequired;
56use crate::tool_types::{SideEffectClass, ToolCall, ToolDefinition, ToolResult};
57use crate::typed_id::{AgentId, HarnessId};
58use crate::{
59 durability::DurableToolResultStore, durability::ToolCallClaimResult,
60 event_emitter::EventEmitter, execution_loading::AgentStore, execution_loading::SessionStore,
61 session_files::SessionFileSystem, tool_context::ToolContext, tool_execution::ToolExecutor,
62};
63use uuid::Uuid;
64
65struct AbortOnDropJoinHandle<T> {
69 handle: everruns_contracts::rt::JoinHandle<T>,
70}
71
72impl<T> AbortOnDropJoinHandle<T> {
73 fn new(handle: everruns_contracts::rt::JoinHandle<T>) -> Self {
74 Self { handle }
75 }
76}
77
78impl<T> Future for AbortOnDropJoinHandle<T> {
79 type Output = std::result::Result<T, everruns_contracts::rt::JoinError>;
80
81 fn poll(mut self: Pin<&mut Self>, cx: &mut Context<'_>) -> Poll<Self::Output> {
82 Pin::new(&mut self.handle).poll(cx)
83 }
84}
85
86impl<T> Drop for AbortOnDropJoinHandle<T> {
87 fn drop(&mut self) {
88 if !self.handle.is_finished() {
89 self.handle.abort();
90 }
91 }
92}
93
94#[derive(Debug, Clone, Serialize, Deserialize)]
100pub struct ActInput {
101 #[serde(skip_serializing_if = "Option::is_none")]
103 pub org_id: Option<i64>,
104 pub context: ExecutionContext,
106 pub harness_id: HarnessId,
108 #[serde(skip_serializing_if = "Option::is_none")]
110 pub agent_id: Option<AgentId>,
111 pub tool_calls: Vec<ToolCall>,
113 pub tool_definitions: Vec<ToolDefinition>,
115 #[serde(skip_serializing_if = "Option::is_none")]
117 pub locale: Option<String>,
118 #[serde(skip_serializing_if = "Option::is_none")]
121 pub blueprint_id: Option<String>,
122 #[serde(default, skip_serializing_if = "Option::is_none")]
124 pub network_access: Option<crate::network_access::NetworkAccessList>,
125 #[serde(default, skip_serializing_if = "Option::is_none")]
129 pub parallel_tool_calls: Option<bool>,
130}
131
132#[derive(Debug, Clone, Serialize, Deserialize)]
134pub struct ToolCallResult {
135 pub tool_call: ToolCall,
137 pub result: ToolResult,
139 pub success: bool,
141 pub status: String,
143 #[serde(default, skip_serializing_if = "Option::is_none")]
145 pub connection_required: Option<ConnectionRequired>,
146 #[serde(default, skip_serializing_if = "Option::is_none")]
149 pub determinism_fatal: Option<String>,
150}
151
152#[derive(Debug, Clone, Serialize, Deserialize)]
154pub struct ActResult {
155 pub results: Vec<ToolCallResult>,
157 pub completed: bool,
159 pub success_count: u32,
161 pub error_count: u32,
163 #[serde(default)]
168 pub waiting_for_tool_results: bool,
169 #[serde(default, skip_serializing_if = "is_false")]
174 pub waiting_for_url_elicitation: bool,
175 #[serde(default, skip_serializing_if = "is_false")]
177 pub blocked: bool,
178 #[serde(default, skip_serializing_if = "Vec::is_empty")]
182 pub client_tool_calls: Vec<ToolCall>,
183 #[serde(default, skip_serializing_if = "Vec::is_empty")]
185 pub client_tool_definitions: Vec<ToolDefinition>,
186}
187
188fn is_false(value: &bool) -> bool {
189 !*value
190}
191
192pub struct ActAtom<T, E>
210where
211 T: ToolExecutor,
212 E: PhaseEffectSink,
213{
214 tool_executor: Arc<T>,
217 event_emitter: PhaseEffectEmitter<E>,
218 context_services: crate::tool_context::ToolContextServices,
220 outbound_tool_rate_limiter: Option<Arc<dyn crate::tool_execution::OutboundToolRateLimiter>>,
224 durable_tool_result_store: Option<Arc<dyn DurableToolResultStore>>,
229 hooks: Vec<Box<dyn PostActHook>>,
232 post_tool_hooks: Vec<Arc<dyn act_hooks::PostToolExecHook>>,
235 pre_tool_hooks: Vec<Arc<dyn act_hooks::PreToolUseHook>>,
241 tool_call_hooks: Vec<Arc<dyn crate::capabilities::ToolCallHook>>,
244 final_post_tool_hooks: Vec<Arc<dyn act_hooks::PostToolExecHook>>,
247}
248
249impl<T, E> ActAtom<T, E>
250where
251 T: ToolExecutor,
252 E: PhaseEffectSink,
253{
254 pub fn new(tool_executor: T, event_emitter: E) -> Self {
256 Self {
257 tool_executor: Arc::new(tool_executor),
258 event_emitter: PhaseEffectEmitter::new(Arc::new(event_emitter)),
259 context_services: crate::tool_context::ToolContextServices::default(),
260 outbound_tool_rate_limiter: None,
261 durable_tool_result_store: None,
262 hooks: Self::default_hooks(),
263 post_tool_hooks: Vec::new(),
264 pre_tool_hooks: Vec::new(),
265 tool_call_hooks: Vec::new(),
266 final_post_tool_hooks: Self::default_final_hooks(),
267 }
268 }
269
270 pub fn with_file_store(
272 tool_executor: T,
273 event_emitter: E,
274 file_store: Arc<dyn SessionFileSystem>,
275 ) -> Self {
276 Self {
277 tool_executor: Arc::new(tool_executor),
278 event_emitter: PhaseEffectEmitter::new(Arc::new(event_emitter)),
279 context_services: crate::tool_context::ToolContextServices {
280 file_store: Some(file_store),
281 ..Default::default()
282 },
283 outbound_tool_rate_limiter: None,
284 durable_tool_result_store: None,
285 hooks: Self::default_hooks(),
286 post_tool_hooks: Vec::new(),
287 pre_tool_hooks: Vec::new(),
288 tool_call_hooks: Vec::new(),
289 final_post_tool_hooks: Self::default_final_hooks(),
290 }
291 }
292
293 pub fn with_context_services(
297 mut self,
298 services: crate::tool_context::ToolContextServices,
299 ) -> Self {
300 self.context_services = services;
301 self
302 }
303
304 pub fn with_hook(mut self, hook: Box<dyn PostActHook>) -> Self {
306 self.hooks.push(hook);
307 self
308 }
309
310 pub fn with_final_post_tool_hook(mut self, hook: Arc<dyn act_hooks::PostToolExecHook>) -> Self {
314 let hard_limit_index = self.final_post_tool_hooks.len().saturating_sub(1);
315 self.final_post_tool_hooks.insert(hard_limit_index, hook);
316 self
317 }
318
319 fn default_hooks() -> Vec<Box<dyn PostActHook>> {
322 vec![
323 Box::new(act_hooks::ConnectionSetupHook),
324 Box::new(act_hooks::UrlElicitationHook),
325 Box::new(act_hooks::FormElicitationHook),
326 Box::new(act_hooks::ToolApprovalPauseHook),
327 Box::new(act_hooks::ClientSideToolHook),
328 ]
329 }
330
331 fn default_final_hooks() -> Vec<Arc<dyn act_hooks::PostToolExecHook>> {
334 vec![Arc::new(act_hooks::OutputHardLimitHook)]
335 }
336
337 pub fn with_storage_store(
339 mut self,
340 store: Arc<dyn crate::session_services::SessionStorageStore>,
341 ) -> Self {
342 self.context_services.storage_store = Some(store);
343 self
344 }
345
346 pub fn with_image_store(
348 mut self,
349 store: Arc<dyn crate::image_services::ImageArtifactStore>,
350 ) -> Self {
351 self.context_services.image_store = Some(store);
352 self
353 }
354
355 pub fn with_provider_credential_store(
357 mut self,
358 store: Arc<dyn crate::connection_services::ProviderCredentialStore>,
359 ) -> Self {
360 self.context_services.provider_credential_store = Some(store);
361 self
362 }
363
364 pub fn with_utility_llm_service(mut self, service: Arc<dyn crate::UtilityLlmService>) -> Self {
366 self.context_services.utility_llm_service = Some(service);
367 self
368 }
369
370 pub fn with_mcp_invoker(mut self, invoker: Arc<dyn crate::McpToolInvoker>) -> Self {
372 self.context_services.mcp_invoker = Some(invoker);
373 self
374 }
375
376 pub fn with_egress_service(mut self, service: Arc<dyn crate::EgressService>) -> Self {
378 self.context_services.egress_service = Some(service);
379 self
380 }
381
382 pub fn with_connection_resolver(
384 mut self,
385 resolver: Arc<dyn crate::connection_services::UserConnectionResolver>,
386 ) -> Self {
387 self.context_services.connection_resolver = Some(resolver);
388 self
389 }
390
391 pub fn with_session_store(mut self, store: Arc<dyn SessionStore>) -> Self {
393 self.context_services.session_store = Some(store);
394 self
395 }
396
397 pub fn with_agent_store(mut self, store: Arc<dyn AgentStore>) -> Self {
399 self.context_services.agent_store = Some(store);
400 self
401 }
402
403 pub fn with_schedule_store(
405 mut self,
406 store: Arc<dyn crate::session_services::SessionScheduleStore>,
407 ) -> Self {
408 self.context_services.schedule_store = Some(store);
409 self
410 }
411
412 pub fn with_subagent_delegate(
414 mut self,
415 store: Arc<dyn crate::subagent_delegation::SubagentSessionDelegate>,
416 ) -> Self {
417 self.context_services.subagent_delegate = Some(store);
418 self
419 }
420
421 pub fn with_leased_resource_store(
423 mut self,
424 store: Arc<dyn crate::session_services::LeasedResourceStore>,
425 ) -> Self {
426 self.context_services.leased_resource_store = Some(store);
427 self
428 }
429
430 pub fn with_session_resource_registry(
432 mut self,
433 registry: Arc<dyn crate::session_services::SessionResourceRegistry>,
434 ) -> Self {
435 self.context_services.session_resource_registry = Some(registry);
436 self
437 }
438
439 pub fn with_session_task_registry(
441 mut self,
442 registry: Arc<dyn crate::session_task::SessionTaskRegistry>,
443 ) -> Self {
444 self.context_services.session_task_registry = Some(registry);
445 self
446 }
447
448 pub fn with_capability_registry(
449 mut self,
450 registry: crate::capabilities::CapabilityRegistry,
451 ) -> Self {
452 self.context_services.capability_registry = Some(registry);
453 self
454 }
455
456 pub fn with_tool_registry(mut self, registry: Arc<crate::tools::ToolRegistry>) -> Self {
458 self.context_services.tool_registry = Some(registry);
459 self
460 }
461
462 pub fn with_post_tool_hooks(
466 mut self,
467 hooks: Vec<Arc<dyn act_hooks::PostToolExecHook>>,
468 ) -> Self {
469 self.post_tool_hooks.extend(hooks);
470 self
471 }
472
473 pub fn with_pre_tool_hooks(mut self, hooks: Vec<Arc<dyn act_hooks::PreToolUseHook>>) -> Self {
477 self.pre_tool_hooks.extend(hooks);
478 self
479 }
480
481 pub fn with_tool_call_hooks(
482 mut self,
483 hooks: Vec<Arc<dyn crate::capabilities::ToolCallHook>>,
484 ) -> Self {
485 self.tool_call_hooks.extend(hooks);
486 self
487 }
488
489 pub fn with_org_id(mut self, org_id: crate::typed_id::OrgId) -> Self {
491 self.context_services.org_id = Some(org_id);
492 self
493 }
494
495 pub fn with_network_access(
497 mut self,
498 network_access: Option<crate::network_access::NetworkAccessList>,
499 ) -> Self {
500 self.context_services.network_access = network_access;
501 self
502 }
503
504 pub fn with_budget_checker(
506 mut self,
507 checker: Arc<dyn crate::tool_execution::BudgetChecker>,
508 ) -> Self {
509 self.context_services.budget_checker = Some(checker);
510 self
511 }
512
513 pub fn with_payment_authority(
515 mut self,
516 authority: Arc<dyn crate::tool_execution::PaymentAuthority>,
517 ) -> Self {
518 self.context_services.payment_authority = Some(authority);
519 self
520 }
521
522 pub fn with_session_creation_authority(
524 mut self,
525 authority: Arc<dyn crate::delegation_services::SessionCreationAuthority>,
526 ) -> Self {
527 self.context_services.session_creation_authority = Some(authority);
528 self
529 }
530
531 pub fn with_outbound_tool_rate_limiter(
533 mut self,
534 limiter: Arc<dyn crate::tool_execution::OutboundToolRateLimiter>,
535 ) -> Self {
536 self.outbound_tool_rate_limiter = Some(limiter);
537 self
538 }
539
540 pub fn with_durable_tool_result_store(
542 mut self,
543 store: Arc<dyn DurableToolResultStore>,
544 ) -> Self {
545 self.durable_tool_result_store = Some(store);
546 self
547 }
548
549 pub fn with_subagent_spawn_store(
551 mut self,
552 store: Arc<dyn crate::delegation_services::SubagentSpawnStore>,
553 ) -> Self {
554 self.context_services.subagent_spawn_store = Some(store);
555 self
556 }
557
558 pub fn with_subagent_nesting_policy(
560 mut self,
561 policy: crate::delegation_services::SubagentNestingPolicy,
562 ) -> Self {
563 self.context_services.subagent_nesting_policy = policy;
564 self
565 }
566
567 pub fn with_reasoning_effort_handle(
571 mut self,
572 handle: crate::tool_context::ReasoningEffortHandle,
573 ) -> Self {
574 self.context_services.reasoning_effort_handle = Some(handle);
575 self
576 }
577}
578
579impl<T, E> ActAtom<T, E>
580where
581 T: ToolExecutor + Send + Sync + 'static,
582 E: EventEmitter + Send + Sync + 'static,
583{
584 pub fn name(&self) -> &'static str {
586 "act"
587 }
588
589 pub async fn execute(&self, input: ActInput) -> Result<ActResult> {
591 let ActInput {
592 context,
593 tool_calls,
594 tool_definitions,
595 locale,
596 network_access,
597 parallel_tool_calls,
598 .. } = input;
600
601 let (server_tool_calls, client_tool_calls): (Vec<_>, Vec<_>) = tool_calls
604 .into_iter()
605 .partition(|tc| act_hooks::runs_on_server(tc, &tool_definitions));
606
607 let client_tool_calls: Vec<_> = client_tool_calls
608 .into_iter()
609 .map(|tool_call| self.transform_tool_call_for_execution(tool_call))
610 .collect();
611
612 let client_tool_definitions: Vec<_> = if client_tool_calls.is_empty() {
613 vec![]
614 } else {
615 tool_definitions
616 .iter()
617 .filter(|td| {
618 if let ToolDefinition::ClientSide(ct) = td {
619 client_tool_calls.iter().any(|tc| tc.name == ct.name)
620 } else {
621 false
622 }
623 })
624 .cloned()
625 .collect()
626 };
627
628 if server_tool_calls.is_empty() && client_tool_calls.is_empty() {
629 return Ok(ActResult {
630 results: vec![],
631 completed: true,
632 success_count: 0,
633 error_count: 0,
634 waiting_for_tool_results: false,
635 waiting_for_url_elicitation: false,
636 blocked: false,
637 client_tool_calls: vec![],
638 client_tool_definitions: vec![],
639 });
640 }
641
642 if server_tool_calls.is_empty() {
645 let mut result = ActResult {
646 results: vec![],
647 completed: true,
648 success_count: 0,
649 error_count: 0,
650 waiting_for_tool_results: false,
651 waiting_for_url_elicitation: false,
652 blocked: false,
653 client_tool_calls,
654 client_tool_definitions,
655 };
656 act_hooks::run_post_act_hooks(
657 &self.hooks,
658 &context,
659 &mut result,
660 &tool_definitions,
661 &self.event_emitter,
662 locale.as_deref(),
663 )
664 .await;
665 return Ok(result);
666 }
667
668 let tool_calls = server_tool_calls;
670
671 tracing::info!(
672 session_id = %context.session_id,
673 turn_id = %context.turn_id,
674 exec_id = %context.exec_id,
675 tool_count = %tool_calls.len(),
676 "ActAtom: executing tools in parallel"
677 );
678
679 let trace_id = context.turn_id.to_string();
687 let act_span_id = Uuid::now_v7().to_string();
688 let parent_span_id = trace_id.clone(); let event_context = EventContext::from_execution_context(&context).with_span(
692 trace_id.clone(),
693 act_span_id.clone(),
694 Some(parent_span_id.clone()),
695 );
696
697 let act_start = Instant::now();
699
700 let visible_tool_names = Arc::new(
701 tool_definitions
702 .iter()
703 .map(|def| def.name().to_string())
704 .collect::<HashSet<_>>(),
705 );
706
707 let tool_map: std::collections::HashMap<&str, &ToolDefinition> = tool_definitions
709 .iter()
710 .map(|def| {
711 let name = def.name();
712 (name, def)
713 })
714 .collect();
715
716 let mut started_data = ActStartedData::with_definitions_and_locale(
717 &tool_calls,
718 &tool_definitions,
719 locale.as_deref(),
720 );
721 for summary in &mut started_data.tool_calls {
722 if let Some(tool_call) = tool_calls.iter().find(|tc| tc.id == summary.id) {
723 let tool_def = tool_map.get(tool_call.name.as_str()).copied();
724 summary.narration = Some(self.render_tool_narration(
725 &context,
726 tool_def,
727 tool_call,
728 ToolNarrationPhase::Started,
729 locale.as_deref(),
730 ));
731 summary.completed_narration = Some(self.render_tool_narration(
732 &context,
733 tool_def,
734 tool_call,
735 ToolNarrationPhase::Completed,
736 locale.as_deref(),
737 ));
738 }
739 }
740 started_data.headline = self.render_group_headline(
741 &context,
742 &tool_calls,
743 &tool_map,
744 ToolNarrationPhase::Started,
745 locale.as_deref(),
746 );
747
748 if let Err(e) = self
750 .event_emitter
751 .emit(EventRequest::new(
752 context.session_id,
753 event_context.clone(),
754 started_data,
755 ))
756 .await
757 {
758 tracing::warn!(
759 session_id = %context.session_id,
760 error = %e,
761 "ActAtom: failed to emit act.started event"
762 );
763 }
764
765 let classes: Vec<Option<String>> = tool_calls
772 .iter()
773 .map(|tool_call| {
774 tool_map
775 .get(tool_call.name.as_str())
776 .and_then(|def| def.concurrency_class())
777 .map(|class| class.to_string())
778 })
779 .collect();
780 let schedule_config = tool_scheduler::ScheduleConfig {
781 serialize_all: parallel_tool_calls == Some(false),
782 ..tool_scheduler::ScheduleConfig::default()
783 };
784 let results =
785 tool_scheduler::schedule(tool_calls.len(), &classes, schedule_config, |index| {
786 let tool_call = &tool_calls[index];
787 let tool_def = tool_map.get(tool_call.name.as_str()).cloned();
788 self.execute_single_tool(
789 &context,
790 tool_call.clone(),
791 tool_def,
792 &trace_id,
793 &act_span_id,
794 locale.as_deref(),
795 network_access.as_ref(),
796 visible_tool_names.clone(),
797 )
798 })
799 .await;
800
801 let success_count = results.iter().filter(|r| r.success).count() as u32;
803 let error_count = results.iter().filter(|r| !r.success).count() as u32;
804
805 let act_duration_ms = act_start.elapsed().as_millis() as u64;
807
808 let completed_context = EventContext::from_execution_context(&context).with_span(
810 trace_id.clone(),
811 act_span_id.clone(), Some(parent_span_id.clone()),
813 );
814 let mut completed_headline = self.render_group_headline(
815 &context,
816 &tool_calls,
817 &tool_map,
818 ToolNarrationPhase::Completed,
819 locale.as_deref(),
820 );
821 if error_count > 0 {
822 let suffix = crate::localization::format_error_suffix(locale.as_deref(), error_count);
823 completed_headline = Some(match completed_headline {
824 Some(text) => format!("{text}{suffix}"),
825 None => {
826 crate::localization::format_completed_tool_batch(locale.as_deref(), error_count)
827 }
828 });
829 }
830
831 if let Err(e) = self
832 .event_emitter
833 .emit(EventRequest::new(
834 context.session_id,
835 completed_context,
836 ActCompletedData {
837 completed: true,
838 success_count,
839 error_count,
840 duration_ms: Some(act_duration_ms),
841 headline: completed_headline,
842 },
843 ))
844 .await
845 {
846 tracing::warn!(
847 session_id = %context.session_id,
848 error = %e,
849 "ActAtom: failed to emit act.completed event"
850 );
851 }
852
853 tracing::info!(
854 session_id = %context.session_id,
855 turn_id = %context.turn_id,
856 success_count = %success_count,
857 error_count = %error_count,
858 "ActAtom: all tools completed"
859 );
860
861 if let Some(fatal_msg) = results.iter().find_map(|r| r.determinism_fatal.as_deref()) {
864 return Err(crate::error::AgentLoopError::tool(format!(
865 "act activity aborted due to determinism violation: {fatal_msg}"
866 )));
867 }
868
869 let mut act_result = ActResult {
870 results,
871 completed: true,
872 success_count,
873 error_count,
874 waiting_for_tool_results: false,
875 waiting_for_url_elicitation: false,
876 blocked: false,
877 client_tool_calls,
878 client_tool_definitions,
879 };
880
881 act_hooks::run_post_act_hooks(
883 &self.hooks,
884 &context,
885 &mut act_result,
886 &tool_definitions,
887 &self.event_emitter,
888 locale.as_deref(),
889 )
890 .await;
891
892 Ok(act_result)
893 }
894}
895
896impl<T, E> ActAtom<T, E>
897where
898 T: ToolExecutor + Send + Sync + 'static,
899 E: EventEmitter + Send + Sync + 'static,
900{
901 fn render_tool_narration(
902 &self,
903 execution_context: &ExecutionContext,
904 tool_def: Option<&ToolDefinition>,
905 tool_call: &ToolCall,
906 phase: ToolNarrationPhase,
907 locale: Option<&str>,
908 ) -> String {
909 let wrapped_store = self.wrap_file_store_for_narration(execution_context);
910 let ctx = ToolNarrationContext::new(wrapped_store.as_deref());
911 for hook in &self.tool_call_hooks {
912 if let Some(narration) = hook.narration(tool_def, tool_call, phase, locale, ctx) {
913 return narration;
914 }
915 }
916 if let Some(narration) = self
922 .context_services
923 .tool_registry
924 .as_ref()
925 .and_then(|registry| registry.get(&tool_call.name))
926 .and_then(|tool| tool.narrate(tool_call, phase, locale, ctx))
927 {
928 return narration;
929 }
930 render_tool_narration_with_locale(tool_def, tool_call, phase, locale)
931 }
932
933 fn render_group_headline(
934 &self,
935 execution_context: &ExecutionContext,
936 tool_calls: &[ToolCall],
937 tool_map: &std::collections::HashMap<&str, &ToolDefinition>,
938 phase: ToolNarrationPhase,
939 locale: Option<&str>,
940 ) -> Option<String> {
941 if tool_calls.is_empty() {
942 return None;
943 }
944 if let [tool_call] = tool_calls {
945 return Some(self.render_tool_narration(
946 execution_context,
947 tool_map.get(tool_call.name.as_str()).copied(),
948 tool_call,
949 phase,
950 locale,
951 ));
952 }
953
954 let actions = tool_calls
955 .iter()
956 .map(|tool_call| {
957 let tool_def = tool_map.get(tool_call.name.as_str()).copied();
958 let narration = self.render_tool_narration(
959 execution_context,
960 tool_def,
961 tool_call,
962 phase,
963 locale,
964 );
965 let repeated_narration = self.render_tool_narration(
966 execution_context,
967 tool_def,
968 &tool_call_for_group_summary(tool_call),
969 phase,
970 locale,
971 );
972 GroupHeadlineAction::new(tool_call, narration, repeated_narration)
973 })
974 .collect::<Vec<_>>();
975
976 Some(summarize_group_actions(&actions, locale))
977 }
978
979 fn wrap_file_store_for_narration(
982 &self,
983 execution_context: &ExecutionContext,
984 ) -> Option<Arc<dyn SessionFileSystem>> {
985 let store = self.context_services.file_store.as_ref()?.clone();
986 let store = if let Some(workspace_id) = execution_context.workspace_id {
987 crate::session_files::WorkspaceScopedFileSystem::wrap(store, workspace_id)
988 } else {
989 store
990 };
991 Some(crate::mount_fs::MountFs::wrap_if_needed(store))
992 }
993
994 fn transform_tool_call_for_execution(&self, tool_call: ToolCall) -> ToolCall {
995 self.tool_call_hooks
996 .iter()
997 .fold(tool_call, |tool_call, hook| {
998 hook.transform_for_execution(tool_call)
999 })
1000 }
1001
1002 #[allow(clippy::too_many_arguments)]
1008 async fn execute_single_tool(
1009 &self,
1010 context: &ExecutionContext,
1011 tool_call: ToolCall,
1012 tool_def: Option<&ToolDefinition>,
1013 trace_id: &str,
1014 act_span_id: &str,
1015 locale: Option<&str>,
1016 network_access: Option<&crate::network_access::NetworkAccessList>,
1017 visible_tool_names: Arc<HashSet<String>>,
1018 ) -> ToolCallResult {
1019 tracing::debug!(
1020 session_id = %context.session_id,
1021 turn_id = %context.turn_id,
1022 tool_name = %tool_call.name,
1023 tool_call_id = %tool_call.id,
1024 "ActAtom: executing tool"
1025 );
1026
1027 let tool_span_id = Uuid::now_v7().to_string();
1029
1030 let event_context = EventContext::from_execution_context(context).with_span(
1032 trace_id.to_string(),
1033 tool_span_id.clone(),
1034 Some(act_span_id.to_string()),
1035 );
1036
1037 let tool_start = Instant::now();
1039 let tool_call_fingerprint = tool_call_fingerprint(&tool_call);
1040
1041 let display_name = crate::localization::localized_tool_display_name(
1043 &tool_call.name,
1044 tool_def.and_then(|d| d.display_name()),
1045 locale,
1046 );
1047 let capability_attribution = tool_def.and_then(|def| {
1048 def.capability_attribution()
1049 .map(|(id, name)| (id.to_string(), name.map(str::to_string)))
1050 });
1051
1052 if let (Some(limiter), Some(ref org_id)) = (
1056 &self.outbound_tool_rate_limiter,
1057 self.context_services.org_id,
1058 ) && !limiter.check_org(org_id).await
1059 {
1060 tracing::warn!(
1061 session_id = %context.session_id,
1062 tool_name = %tool_call.name,
1063 "ActAtom: outbound tool rate limit exceeded for org"
1064 );
1065 return ToolCallResult {
1066 tool_call: tool_call.clone(),
1067 result: ToolResult {
1068 tool_call_id: tool_call.id.clone(),
1069 result: None,
1070 images: None,
1071 error: Some(
1072 "Outbound tool rate limit exceeded for this organization; back off and retry later.".to_string(),
1073 ),
1074 connection_required: None,
1075 raw_output: None,
1076 },
1077 success: false,
1078 status: "error".to_string(),
1079 connection_required: None,
1080 determinism_fatal: None,
1081 };
1082 }
1083
1084 let claim_token = if let Some(ref store) = self.durable_tool_result_store {
1087 let turn_id = context.turn_id.to_string();
1088 match store
1089 .try_claim_tool_call(
1090 &turn_id,
1091 &tool_call.id,
1092 &tool_call.name,
1093 &tool_call_fingerprint,
1094 )
1095 .await
1096 {
1097 Ok(ToolCallClaimResult::Claimed { claim_token }) => Some(claim_token),
1098
1099 Ok(ToolCallClaimResult::AlreadySettled {
1100 result_json,
1101 args_fingerprint: stored_fp,
1102 }) => {
1103 if stored_fp != tool_call_fingerprint {
1105 let err_msg = format!(
1106 "determinism violation: tool '{}' replay args fingerprint \
1107 does not match prior execution (stored={stored_fp}, \
1108 current={})",
1109 tool_call.name, tool_call_fingerprint
1110 );
1111 tracing::error!(
1112 session_id = %context.session_id,
1113 turn_id = %context.turn_id,
1114 tool_call_id = %tool_call.id,
1115 stored_fp = %stored_fp,
1116 current_fp = %tool_call_fingerprint,
1117 "ActAtom: determinism violation — replay args fingerprint mismatch"
1118 );
1119 let result_fp =
1120 tool_result_fingerprint(&tool_call.name, &ToolResult::error(&err_msg));
1121 let _ = self
1122 .event_emitter
1123 .emit(EventRequest::new(
1124 context.session_id,
1125 event_context,
1126 ToolCompletedData::failure(
1127 tool_call.id.clone(),
1128 tool_call.name.clone(),
1129 "error".to_string(),
1130 err_msg.clone(),
1131 None,
1132 )
1133 .with_fingerprints(tool_call_fingerprint.clone(), result_fp)
1134 .with_display_name(display_name.clone()),
1135 ))
1136 .await;
1137 return ToolCallResult {
1138 tool_call: tool_call.clone(),
1139 result: ToolResult {
1140 tool_call_id: tool_call.id.clone(),
1141 result: None,
1142 images: None,
1143 error: Some(err_msg.clone()),
1144 connection_required: None,
1145 raw_output: None,
1146 },
1147 success: false,
1148 status: "error".to_string(),
1149 connection_required: None,
1150 determinism_fatal: Some(err_msg),
1151 };
1152 }
1153 tracing::debug!(
1154 session_id = %context.session_id,
1155 turn_id = %context.turn_id,
1156 tool_call_id = %tool_call.id,
1157 "ActAtom: replaying already-settled tool call"
1158 );
1159 let replayed_result: ToolResult = serde_json::from_value(result_json.clone())
1161 .unwrap_or(ToolResult {
1162 tool_call_id: tool_call.id.clone(),
1163 result: Some(result_json),
1164 images: None,
1165 error: None,
1166 connection_required: None,
1167 raw_output: None,
1168 });
1169 let success = replayed_result.error.is_none();
1170 let status = if success { "success" } else { "error" };
1171 let result_fp = tool_result_fingerprint(&tool_call.name, &replayed_result);
1172 let completed_data = if success {
1173 let mut content = replayed_result
1175 .result
1176 .as_ref()
1177 .map(|r| vec![ContentPart::tool_result_text(r)])
1178 .unwrap_or_default();
1179 if let Some(ref images) = replayed_result.images {
1180 for img in images {
1181 content.push(ContentPart::Image(
1182 crate::message::ImageContentPart::from_base64(
1183 &img.base64,
1184 &img.media_type,
1185 ),
1186 ));
1187 }
1188 }
1189 ToolCompletedData::success(
1190 tool_call.id.clone(),
1191 tool_call.name.clone(),
1192 content,
1193 None,
1194 )
1195 .with_fingerprints(tool_call_fingerprint.clone(), result_fp)
1196 .with_display_name(display_name.clone())
1197 } else {
1198 ToolCompletedData::failure(
1199 tool_call.id.clone(),
1200 tool_call.name.clone(),
1201 status.to_string(),
1202 replayed_result.error.clone().unwrap_or_default(),
1203 None,
1204 )
1205 .with_fingerprints(tool_call_fingerprint.clone(), result_fp)
1206 .with_display_name(display_name.clone())
1207 };
1208 let _ = self
1209 .event_emitter
1210 .emit(EventRequest::new(
1211 context.session_id,
1212 event_context,
1213 completed_data,
1214 ))
1215 .await;
1216 let conn_req = replayed_result.connection_required.clone();
1217 return ToolCallResult {
1218 tool_call,
1219 result: replayed_result,
1220 success,
1221 status: status.to_string(),
1222 connection_required: conn_req,
1223 determinism_fatal: None,
1224 };
1225 }
1226
1227 Ok(ToolCallClaimResult::AlreadyRunning {
1228 args_fingerprint: stored_fp,
1229 }) => {
1230 if stored_fp != tool_call_fingerprint {
1233 let err_msg = format!(
1234 "determinism violation: tool '{}' args fingerprint changed \
1235 while prior claim is still running (stored={stored_fp}, \
1236 current={tool_call_fingerprint})",
1237 tool_call.name
1238 );
1239 tracing::error!(
1240 session_id = %context.session_id,
1241 turn_id = %context.turn_id,
1242 tool_call_id = %tool_call.id,
1243 stored = %stored_fp,
1244 current = %tool_call_fingerprint,
1245 "ActAtom: determinism violation — running claim fingerprint mismatch"
1246 );
1247 let result_fp =
1248 tool_result_fingerprint(&tool_call.name, &ToolResult::error(&err_msg));
1249 let _ = self
1250 .event_emitter
1251 .emit(EventRequest::new(
1252 context.session_id,
1253 event_context,
1254 ToolCompletedData::failure(
1255 tool_call.id.clone(),
1256 tool_call.name.clone(),
1257 "error".to_string(),
1258 err_msg.clone(),
1259 None,
1260 )
1261 .with_fingerprints(tool_call_fingerprint.clone(), result_fp)
1262 .with_display_name(display_name.clone()),
1263 ))
1264 .await;
1265 return ToolCallResult {
1266 tool_call: tool_call.clone(),
1267 result: ToolResult {
1268 tool_call_id: tool_call.id.clone(),
1269 result: None,
1270 images: None,
1271 error: Some(err_msg.clone()),
1272 connection_required: None,
1273 raw_output: None,
1274 },
1275 success: false,
1276 status: "error".to_string(),
1277 connection_required: None,
1278 determinism_fatal: Some(err_msg),
1279 };
1280 }
1281
1282 let sec = tool_def
1283 .map(|d| d.side_effect_class())
1284 .unwrap_or(SideEffectClass::AtMostOnce);
1285 match sec {
1286 SideEffectClass::Pure | SideEffectClass::Idempotent => {
1287 tracing::debug!(
1289 session_id = %context.session_id,
1290 tool_call_id = %tool_call.id,
1291 "ActAtom: stale running claim for idempotent tool, re-executing"
1292 );
1293 None
1294 }
1295 SideEffectClass::AtMostOnce => {
1296 tracing::warn!(
1297 session_id = %context.session_id,
1298 turn_id = %context.turn_id,
1299 tool_call_id = %tool_call.id,
1300 "ActAtom: AtMostOnce tool has stale running claim; returning interrupted result"
1301 );
1302 let _ = store
1304 .settle_tool_call(
1305 &turn_id,
1306 &tool_call.id,
1307 serde_json::Value::Null,
1308 "interrupted",
1309 Uuid::nil(), )
1311 .await;
1312 let err_msg = format!(
1313 "tool '{}' was interrupted mid-execution during a prior \
1314 worker failure; result is uncertain and was not re-run \
1315 (AtMostOnce safety)",
1316 tool_call.name
1317 );
1318 let result_fp = tool_result_fingerprint(
1319 &tool_call.name,
1320 &ToolResult::error(&err_msg),
1321 );
1322 let _ = self
1323 .event_emitter
1324 .emit(EventRequest::new(
1325 context.session_id,
1326 event_context,
1327 ToolCompletedData::failure(
1328 tool_call.id.clone(),
1329 tool_call.name.clone(),
1330 "interrupted".to_string(),
1331 err_msg.clone(),
1332 None,
1333 )
1334 .with_fingerprints(tool_call_fingerprint.clone(), result_fp)
1335 .with_display_name(display_name.clone()),
1336 ))
1337 .await;
1338 return ToolCallResult {
1339 tool_call: tool_call.clone(),
1340 result: ToolResult {
1341 tool_call_id: tool_call.id.clone(),
1342 result: None,
1343 images: None,
1344 error: Some(err_msg),
1345 connection_required: None,
1346 raw_output: None,
1347 },
1348 success: false,
1349 status: "error".to_string(),
1350 connection_required: None,
1351 determinism_fatal: None,
1352 };
1353 }
1354 }
1355 }
1356
1357 Ok(ToolCallClaimResult::DeterminismViolation {
1358 stored_fingerprint,
1359 current_fingerprint,
1360 }) => {
1361 let err_msg = format!(
1362 "determinism violation: tool '{}' args fingerprint changed \
1363 on replay (stored={stored_fingerprint}, \
1364 current={current_fingerprint})",
1365 tool_call.name
1366 );
1367 tracing::error!(
1368 session_id = %context.session_id,
1369 turn_id = %context.turn_id,
1370 tool_call_id = %tool_call.id,
1371 stored = %stored_fingerprint,
1372 current = %current_fingerprint,
1373 "ActAtom: determinism violation on claim"
1374 );
1375 let result_fp =
1376 tool_result_fingerprint(&tool_call.name, &ToolResult::error(&err_msg));
1377 let _ = self
1378 .event_emitter
1379 .emit(EventRequest::new(
1380 context.session_id,
1381 event_context,
1382 ToolCompletedData::failure(
1383 tool_call.id.clone(),
1384 tool_call.name.clone(),
1385 "error".to_string(),
1386 err_msg.clone(),
1387 None,
1388 )
1389 .with_fingerprints(tool_call_fingerprint.clone(), result_fp)
1390 .with_display_name(display_name.clone()),
1391 ))
1392 .await;
1393 return ToolCallResult {
1394 tool_call: tool_call.clone(),
1395 result: ToolResult {
1396 tool_call_id: tool_call.id.clone(),
1397 result: None,
1398 images: None,
1399 error: Some(err_msg.clone()),
1400 connection_required: None,
1401 raw_output: None,
1402 },
1403 success: false,
1404 status: "error".to_string(),
1405 connection_required: None,
1406 determinism_fatal: Some(err_msg),
1407 };
1408 }
1409
1410 Err(e) => {
1411 tracing::warn!(
1412 session_id = %context.session_id,
1413 tool_call_id = %tool_call.id,
1414 error = %e,
1415 "ActAtom: durable claim failed; proceeding without idempotency"
1416 );
1417 None
1418 }
1419 }
1420 } else {
1421 None
1422 };
1423
1424 if let Err(e) = self
1426 .event_emitter
1427 .emit(EventRequest::new(
1428 context.session_id,
1429 event_context.clone(),
1430 ToolStartedData {
1431 tool_call: tool_call.clone(),
1432 tool_call_fingerprint: Some(tool_call_fingerprint.clone()),
1433 display_name: display_name.clone(),
1434 narration: Some(self.render_tool_narration(
1435 context,
1436 tool_def,
1437 &tool_call,
1438 ToolNarrationPhase::Started,
1439 locale,
1440 )),
1441 },
1442 ))
1443 .await
1444 {
1445 tracing::warn!(
1446 session_id = %context.session_id,
1447 tool_call_id = %tool_call.id,
1448 error = %e,
1449 "ActAtom: failed to emit tool.started event"
1450 );
1451 }
1452
1453 let Some(tool_def) = tool_def else {
1455 let error_msg = format!("Tool definition not found: {}", tool_call.name);
1456 let tool_duration_ms = tool_start.elapsed().as_millis() as u64;
1457
1458 if let Err(e) = self
1460 .event_emitter
1461 .emit(EventRequest::new(
1462 context.session_id,
1463 event_context,
1464 ToolCompletedData::failure(
1465 tool_call.id.clone(),
1466 tool_call.name.clone(),
1467 "error".to_string(),
1468 error_msg.clone(),
1469 Some(tool_duration_ms),
1470 )
1471 .with_fingerprints(
1472 tool_call_fingerprint.clone(),
1473 tool_error_fingerprint(&tool_call.name, "error", &error_msg),
1474 )
1475 .with_narration(Some(self.render_tool_narration(
1476 context,
1477 None,
1478 &tool_call,
1479 ToolNarrationPhase::Failed,
1480 locale,
1481 ))),
1482 ))
1483 .await
1484 {
1485 tracing::warn!(
1486 session_id = %context.session_id,
1487 tool_call_id = %tool_call.id,
1488 error = %e,
1489 "ActAtom: failed to emit tool.completed event"
1490 );
1491 }
1492
1493 return ToolCallResult {
1494 tool_call: tool_call.clone(),
1495 result: ToolResult {
1496 tool_call_id: tool_call.id.clone(),
1497 result: None,
1498 images: None,
1499 error: Some(error_msg),
1500 connection_required: None,
1501 raw_output: None,
1502 },
1503 success: false,
1504 status: "error".to_string(),
1505 connection_required: None,
1506 determinism_fatal: None,
1507 };
1508 };
1509
1510 let mut tool_context =
1512 ToolContext::from_services(context.session_id, &self.context_services);
1513 if let Some(resolver) = tool_context.connection_resolver.as_ref()
1514 && let Some(bound) = resolver.for_execution(context.input_message_id.uuid())
1515 {
1516 tool_context.connection_resolver = Some(bound);
1517 }
1518
1519 if let Some(scope) =
1520 tool_context.extension::<everruns_core::tool_context::ExecutionServicesExt>()
1521 {
1522 scope
1523 .0
1524 .bind(&mut tool_context, context.input_message_id.uuid());
1525 }
1526 if let Some(invoker) = tool_context.mcp_invoker.as_ref()
1527 && let Some(bound) = invoker.for_execution(context.input_message_id.uuid())
1528 {
1529 tool_context.mcp_invoker = Some(bound);
1530 }
1531 if let Some(authority) = tool_context.session_creation_authority.as_ref()
1532 && let Some(bound) = authority.for_execution(context.input_message_id.uuid())
1533 {
1534 tool_context.session_creation_authority = Some(bound);
1535 }
1536 if let Some(workspace_id) = context.workspace_id {
1541 tool_context.workspace_id = workspace_id;
1542 if let Some(store) = tool_context.file_store.take() {
1543 tool_context.file_store = Some(
1544 crate::session_files::WorkspaceScopedFileSystem::wrap(store, workspace_id),
1545 );
1546 }
1547 }
1548 if let Some(store) = tool_context.file_store.take() {
1552 tool_context.file_store = Some(crate::mount_fs::MountFs::wrap_if_needed(store));
1553 }
1554 tool_context.visible_tool_names = Some(visible_tool_names.clone());
1555 tool_context.network_access = network_access
1557 .cloned()
1558 .or_else(|| self.context_services.network_access.clone());
1559 if tool_context.event_emitter.is_none() {
1561 tool_context.event_emitter =
1562 Some(Arc::new(self.event_emitter.clone()) as Arc<dyn EventEmitter>);
1563 }
1564 tool_context.bind_to_turn(event_context.clone());
1565 tool_context.tool_call_id = Some(tool_call.id.clone());
1566
1567 let call_cancellation = tokio_util::sync::CancellationToken::new();
1575 tool_context.cancellation = Some(call_cancellation.clone());
1576 let _cancel_on_call_end = call_cancellation.drop_guard();
1577
1578 let execution_tool_call = self.transform_tool_call_for_execution(tool_call.clone());
1579
1580 let (execution_tool_call, pre_hook_result) = act_hooks::pre_tool_use_outcome(
1584 &self.pre_tool_hooks,
1585 execution_tool_call,
1586 tool_def,
1587 &tool_context,
1588 )
1589 .await;
1590
1591 let result = if let Some(pre_hook_result) = pre_hook_result {
1592 Ok(pre_hook_result)
1593 } else if tool_def.is_cpu_bound() {
1594 let executor = self.tool_executor.clone();
1600 let call = execution_tool_call.clone();
1601 let def = tool_def.clone();
1602 let ctx = tool_context.clone();
1603 match AbortOnDropJoinHandle::new(everruns_contracts::rt::spawn(async move {
1604 executor.execute_with_context(&call, &def, &ctx).await
1605 }))
1606 .await
1607 {
1608 Ok(result) => result,
1609 Err(join_err) => Err(crate::error::AgentLoopError::tool(format!(
1610 "tool task failed to complete: {join_err}"
1611 ))),
1612 }
1613 } else {
1614 self.tool_executor
1615 .execute_with_context(&execution_tool_call, tool_def, &tool_context)
1616 .await
1617 };
1618
1619 match result {
1620 Ok(mut tool_result) => {
1621 act_hooks::run_post_tool_exec_hooks(
1623 &self.post_tool_hooks,
1624 &self.final_post_tool_hooks,
1625 &execution_tool_call,
1626 tool_def,
1627 &mut tool_result,
1628 &tool_context,
1629 )
1630 .await;
1631
1632 let tool_duration_ms = tool_start.elapsed().as_millis() as u64;
1633 let success = tool_result.error.is_none();
1634 let status = if success { "success" } else { "error" };
1635
1636 let completed_data = if success {
1638 let result_fingerprint = tool_result_fingerprint(&tool_call.name, &tool_result);
1639 let mut result_content = tool_result
1641 .result
1642 .as_ref()
1643 .map(|r| vec![ContentPart::tool_result_text(r)])
1644 .unwrap_or_default();
1645 if let Some(ref images) = tool_result.images {
1647 for img in images {
1648 result_content.push(ContentPart::Image(
1649 crate::message::ImageContentPart::from_base64(
1650 &img.base64,
1651 &img.media_type,
1652 ),
1653 ));
1654 }
1655 }
1656 ToolCompletedData::success(
1657 tool_call.id.clone(),
1658 tool_call.name.clone(),
1659 result_content,
1660 Some(tool_duration_ms),
1661 )
1662 .with_fingerprints(tool_call_fingerprint.clone(), result_fingerprint)
1663 .with_display_name(display_name.clone())
1664 .with_capability_attribution(
1665 capability_attribution.as_ref().map(|(id, _)| id.clone()),
1666 capability_attribution
1667 .as_ref()
1668 .and_then(|(_, name)| name.clone()),
1669 )
1670 .with_narration(Some(self.render_tool_narration(
1671 context,
1672 Some(tool_def),
1673 &tool_call,
1674 ToolNarrationPhase::Completed,
1675 locale,
1676 )))
1677 } else {
1678 let result_fingerprint = tool_result_fingerprint(&tool_call.name, &tool_result);
1679 ToolCompletedData::failure(
1680 tool_call.id.clone(),
1681 tool_call.name.clone(),
1682 status.to_string(),
1683 tool_result.error.clone().unwrap_or_default(),
1684 Some(tool_duration_ms),
1685 )
1686 .with_fingerprints(tool_call_fingerprint.clone(), result_fingerprint)
1687 .with_display_name(display_name.clone())
1688 .with_capability_attribution(
1689 capability_attribution.as_ref().map(|(id, _)| id.clone()),
1690 capability_attribution
1691 .as_ref()
1692 .and_then(|(_, name)| name.clone()),
1693 )
1694 .with_narration(Some(self.render_tool_narration(
1695 context,
1696 Some(tool_def),
1697 &tool_call,
1698 ToolNarrationPhase::Failed,
1699 locale,
1700 )))
1701 };
1702
1703 if let Err(e) = self
1704 .event_emitter
1705 .emit(EventRequest::new(
1706 context.session_id,
1707 event_context.clone(),
1708 completed_data,
1709 ))
1710 .await
1711 {
1712 tracing::warn!(
1713 session_id = %context.session_id,
1714 tool_call_id = %tool_call.id,
1715 error = %e,
1716 "ActAtom: failed to emit tool.completed event"
1717 );
1718 }
1719
1720 tracing::debug!(
1721 session_id = %context.session_id,
1722 tool_name = %tool_call.name,
1723 tool_call_id = %tool_call.id,
1724 success = %success,
1725 "ActAtom: tool execution completed"
1726 );
1727
1728 if let (Some(store), Some(token)) = (&self.durable_tool_result_store, claim_token) {
1730 let result_snapshot =
1731 serde_json::to_value(&tool_result).unwrap_or(serde_json::Value::Null);
1732 match store
1733 .settle_tool_call(
1734 &context.turn_id.to_string(),
1735 &tool_call.id,
1736 result_snapshot,
1737 "settled",
1738 token,
1739 )
1740 .await
1741 {
1742 Ok(false) => {
1743 tracing::warn!(
1744 session_id = %context.session_id,
1745 tool_call_id = %tool_call.id,
1746 "ActAtom: settle ownership check failed (task reclaimed)"
1747 );
1748 }
1749 Err(e) => {
1750 tracing::warn!(
1751 session_id = %context.session_id,
1752 tool_call_id = %tool_call.id,
1753 error = %e,
1754 "ActAtom: settle_tool_call failed"
1755 );
1756 }
1757 Ok(true) => {}
1758 }
1759 }
1760
1761 let conn_req = tool_result.connection_required.clone();
1762 ToolCallResult {
1763 tool_call,
1764 result: tool_result,
1765 success,
1766 status: status.to_string(),
1767 connection_required: conn_req,
1768 determinism_fatal: None,
1769 }
1770 }
1771 Err(e) => {
1772 let tool_duration_ms = tool_start.elapsed().as_millis() as u64;
1773 let error_msg = e.to_string();
1774
1775 if let Err(emit_err) = self
1777 .event_emitter
1778 .emit(EventRequest::new(
1779 context.session_id,
1780 event_context,
1781 ToolCompletedData::failure(
1782 tool_call.id.clone(),
1783 tool_call.name.clone(),
1784 "error".to_string(),
1785 error_msg.clone(),
1786 Some(tool_duration_ms),
1787 )
1788 .with_fingerprints(
1789 tool_call_fingerprint.clone(),
1790 tool_error_fingerprint(&tool_call.name, "error", &error_msg),
1791 )
1792 .with_display_name(display_name.clone())
1793 .with_capability_attribution(
1794 capability_attribution.as_ref().map(|(id, _)| id.clone()),
1795 capability_attribution
1796 .as_ref()
1797 .and_then(|(_, name)| name.clone()),
1798 )
1799 .with_narration(Some(self.render_tool_narration(
1800 context,
1801 Some(tool_def),
1802 &tool_call,
1803 ToolNarrationPhase::Failed,
1804 locale,
1805 ))),
1806 ))
1807 .await
1808 {
1809 tracing::warn!(
1810 session_id = %context.session_id,
1811 tool_call_id = %tool_call.id,
1812 error = %emit_err,
1813 "ActAtom: failed to emit tool.completed event"
1814 );
1815 }
1816
1817 tracing::warn!(
1818 session_id = %context.session_id,
1819 tool_name = %tool_call.name,
1820 tool_call_id = %tool_call.id,
1821 error = %e,
1822 "ActAtom: tool execution failed"
1823 );
1824
1825 ToolCallResult {
1826 tool_call: tool_call.clone(),
1827 result: ToolResult {
1828 tool_call_id: tool_call.id.clone(),
1829 result: None,
1830 images: None,
1831 error: Some(error_msg),
1832 connection_required: None,
1833 raw_output: None,
1834 },
1835 success: false,
1836 status: "error".to_string(),
1837 connection_required: None,
1838 determinism_fatal: None,
1839 }
1840 }
1841 }
1842 }
1843}
1844
1845#[cfg(test)]
1850#[path = "act_tests.rs"]
1851mod tests;
1852
1853#[cfg(test)]
1854#[path = "act_approval_tests.rs"]
1855mod approval_tests;