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