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>> {
323 vec![
324 Box::new(act_hooks::ConnectionSetupHook),
325 Box::new(act_hooks::UrlElicitationHook),
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<_>) =
603 tool_calls.into_iter().partition(|tc| {
604 tool_definitions
605 .iter()
606 .find(|td| td.name() == tc.name)
607 .map(|td| !matches!(td, ToolDefinition::ClientSide(_)))
608 .unwrap_or(true) });
610
611 let client_tool_calls: Vec<_> = client_tool_calls
612 .into_iter()
613 .map(|tool_call| self.transform_tool_call_for_execution(tool_call))
614 .collect();
615
616 let client_tool_definitions: Vec<_> = if client_tool_calls.is_empty() {
617 vec![]
618 } else {
619 tool_definitions
620 .iter()
621 .filter(|td| {
622 if let ToolDefinition::ClientSide(ct) = td {
623 client_tool_calls.iter().any(|tc| tc.name == ct.name)
624 } else {
625 false
626 }
627 })
628 .cloned()
629 .collect()
630 };
631
632 if server_tool_calls.is_empty() && client_tool_calls.is_empty() {
633 return Ok(ActResult {
634 results: vec![],
635 completed: true,
636 success_count: 0,
637 error_count: 0,
638 waiting_for_tool_results: false,
639 waiting_for_url_elicitation: false,
640 blocked: false,
641 client_tool_calls: vec![],
642 client_tool_definitions: vec![],
643 });
644 }
645
646 if server_tool_calls.is_empty() {
649 let mut result = ActResult {
650 results: vec![],
651 completed: true,
652 success_count: 0,
653 error_count: 0,
654 waiting_for_tool_results: false,
655 waiting_for_url_elicitation: false,
656 blocked: false,
657 client_tool_calls,
658 client_tool_definitions,
659 };
660 act_hooks::run_post_act_hooks(
661 &self.hooks,
662 &context,
663 &mut result,
664 &tool_definitions,
665 &self.event_emitter,
666 locale.as_deref(),
667 )
668 .await;
669 return Ok(result);
670 }
671
672 let tool_calls = server_tool_calls;
674
675 tracing::info!(
676 session_id = %context.session_id,
677 turn_id = %context.turn_id,
678 exec_id = %context.exec_id,
679 tool_count = %tool_calls.len(),
680 "ActAtom: executing tools in parallel"
681 );
682
683 let trace_id = context.turn_id.to_string();
691 let act_span_id = Uuid::now_v7().to_string();
692 let parent_span_id = trace_id.clone(); let event_context = EventContext::from_execution_context(&context).with_span(
696 trace_id.clone(),
697 act_span_id.clone(),
698 Some(parent_span_id.clone()),
699 );
700
701 let act_start = Instant::now();
703
704 let visible_tool_names = Arc::new(
705 tool_definitions
706 .iter()
707 .map(|def| def.name().to_string())
708 .collect::<HashSet<_>>(),
709 );
710
711 let tool_map: std::collections::HashMap<&str, &ToolDefinition> = tool_definitions
713 .iter()
714 .map(|def| {
715 let name = def.name();
716 (name, def)
717 })
718 .collect();
719
720 let mut started_data = ActStartedData::with_definitions_and_locale(
721 &tool_calls,
722 &tool_definitions,
723 locale.as_deref(),
724 );
725 for summary in &mut started_data.tool_calls {
726 if let Some(tool_call) = tool_calls.iter().find(|tc| tc.id == summary.id) {
727 let tool_def = tool_map.get(tool_call.name.as_str()).copied();
728 summary.narration = Some(self.render_tool_narration(
729 &context,
730 tool_def,
731 tool_call,
732 ToolNarrationPhase::Started,
733 locale.as_deref(),
734 ));
735 summary.completed_narration = Some(self.render_tool_narration(
736 &context,
737 tool_def,
738 tool_call,
739 ToolNarrationPhase::Completed,
740 locale.as_deref(),
741 ));
742 }
743 }
744 started_data.headline = self.render_group_headline(
745 &context,
746 &tool_calls,
747 &tool_map,
748 ToolNarrationPhase::Started,
749 locale.as_deref(),
750 );
751
752 if let Err(e) = self
754 .event_emitter
755 .emit(EventRequest::new(
756 context.session_id,
757 event_context.clone(),
758 started_data,
759 ))
760 .await
761 {
762 tracing::warn!(
763 session_id = %context.session_id,
764 error = %e,
765 "ActAtom: failed to emit act.started event"
766 );
767 }
768
769 let classes: Vec<Option<String>> = tool_calls
776 .iter()
777 .map(|tool_call| {
778 tool_map
779 .get(tool_call.name.as_str())
780 .and_then(|def| def.concurrency_class())
781 .map(|class| class.to_string())
782 })
783 .collect();
784 let schedule_config = tool_scheduler::ScheduleConfig {
785 serialize_all: parallel_tool_calls == Some(false),
786 ..tool_scheduler::ScheduleConfig::default()
787 };
788 let results =
789 tool_scheduler::schedule(tool_calls.len(), &classes, schedule_config, |index| {
790 let tool_call = &tool_calls[index];
791 let tool_def = tool_map.get(tool_call.name.as_str()).cloned();
792 self.execute_single_tool(
793 &context,
794 tool_call.clone(),
795 tool_def,
796 &trace_id,
797 &act_span_id,
798 locale.as_deref(),
799 network_access.as_ref(),
800 visible_tool_names.clone(),
801 )
802 })
803 .await;
804
805 let success_count = results.iter().filter(|r| r.success).count() as u32;
807 let error_count = results.iter().filter(|r| !r.success).count() as u32;
808
809 let act_duration_ms = act_start.elapsed().as_millis() as u64;
811
812 let completed_context = EventContext::from_execution_context(&context).with_span(
814 trace_id.clone(),
815 act_span_id.clone(), Some(parent_span_id.clone()),
817 );
818 let mut completed_headline = self.render_group_headline(
819 &context,
820 &tool_calls,
821 &tool_map,
822 ToolNarrationPhase::Completed,
823 locale.as_deref(),
824 );
825 if error_count > 0 {
826 let suffix = crate::localization::format_error_suffix(locale.as_deref(), error_count);
827 completed_headline = Some(match completed_headline {
828 Some(text) => format!("{text}{suffix}"),
829 None => {
830 crate::localization::format_completed_tool_batch(locale.as_deref(), error_count)
831 }
832 });
833 }
834
835 if let Err(e) = self
836 .event_emitter
837 .emit(EventRequest::new(
838 context.session_id,
839 completed_context,
840 ActCompletedData {
841 completed: true,
842 success_count,
843 error_count,
844 duration_ms: Some(act_duration_ms),
845 headline: completed_headline,
846 },
847 ))
848 .await
849 {
850 tracing::warn!(
851 session_id = %context.session_id,
852 error = %e,
853 "ActAtom: failed to emit act.completed event"
854 );
855 }
856
857 tracing::info!(
858 session_id = %context.session_id,
859 turn_id = %context.turn_id,
860 success_count = %success_count,
861 error_count = %error_count,
862 "ActAtom: all tools completed"
863 );
864
865 if let Some(fatal_msg) = results.iter().find_map(|r| r.determinism_fatal.as_deref()) {
868 return Err(crate::error::AgentLoopError::tool(format!(
869 "act activity aborted due to determinism violation: {fatal_msg}"
870 )));
871 }
872
873 let mut act_result = ActResult {
874 results,
875 completed: true,
876 success_count,
877 error_count,
878 waiting_for_tool_results: false,
879 waiting_for_url_elicitation: false,
880 blocked: false,
881 client_tool_calls,
882 client_tool_definitions,
883 };
884
885 act_hooks::run_post_act_hooks(
887 &self.hooks,
888 &context,
889 &mut act_result,
890 &tool_definitions,
891 &self.event_emitter,
892 locale.as_deref(),
893 )
894 .await;
895
896 Ok(act_result)
897 }
898}
899
900impl<T, E> ActAtom<T, E>
901where
902 T: ToolExecutor + Send + Sync + 'static,
903 E: EventEmitter + Send + Sync + 'static,
904{
905 fn render_tool_narration(
906 &self,
907 execution_context: &ExecutionContext,
908 tool_def: Option<&ToolDefinition>,
909 tool_call: &ToolCall,
910 phase: ToolNarrationPhase,
911 locale: Option<&str>,
912 ) -> String {
913 let wrapped_store = self.wrap_file_store_for_narration(execution_context);
914 let ctx = ToolNarrationContext::new(wrapped_store.as_deref());
915 for hook in &self.tool_call_hooks {
916 if let Some(narration) = hook.narration(tool_def, tool_call, phase, locale, ctx) {
917 return narration;
918 }
919 }
920 if let Some(narration) = self
926 .context_services
927 .tool_registry
928 .as_ref()
929 .and_then(|registry| registry.get(&tool_call.name))
930 .and_then(|tool| tool.narrate(tool_call, phase, locale, ctx))
931 {
932 return narration;
933 }
934 render_tool_narration_with_locale(tool_def, tool_call, phase, locale)
935 }
936
937 fn render_group_headline(
938 &self,
939 execution_context: &ExecutionContext,
940 tool_calls: &[ToolCall],
941 tool_map: &std::collections::HashMap<&str, &ToolDefinition>,
942 phase: ToolNarrationPhase,
943 locale: Option<&str>,
944 ) -> Option<String> {
945 if tool_calls.is_empty() {
946 return None;
947 }
948 if let [tool_call] = tool_calls {
949 return Some(self.render_tool_narration(
950 execution_context,
951 tool_map.get(tool_call.name.as_str()).copied(),
952 tool_call,
953 phase,
954 locale,
955 ));
956 }
957
958 let actions = tool_calls
959 .iter()
960 .map(|tool_call| {
961 let tool_def = tool_map.get(tool_call.name.as_str()).copied();
962 let narration = self.render_tool_narration(
963 execution_context,
964 tool_def,
965 tool_call,
966 phase,
967 locale,
968 );
969 let repeated_narration = self.render_tool_narration(
970 execution_context,
971 tool_def,
972 &tool_call_for_group_summary(tool_call),
973 phase,
974 locale,
975 );
976 GroupHeadlineAction::new(tool_call, narration, repeated_narration)
977 })
978 .collect::<Vec<_>>();
979
980 Some(summarize_group_actions(&actions, locale))
981 }
982
983 fn wrap_file_store_for_narration(
986 &self,
987 execution_context: &ExecutionContext,
988 ) -> Option<Arc<dyn SessionFileSystem>> {
989 let store = self.context_services.file_store.as_ref()?.clone();
990 let store = if let Some(workspace_id) = execution_context.workspace_id {
991 crate::session_files::WorkspaceScopedFileSystem::wrap(store, workspace_id)
992 } else {
993 store
994 };
995 Some(crate::mount_fs::MountFs::wrap_if_needed(store))
996 }
997
998 fn transform_tool_call_for_execution(&self, tool_call: ToolCall) -> ToolCall {
999 self.tool_call_hooks
1000 .iter()
1001 .fold(tool_call, |tool_call, hook| {
1002 hook.transform_for_execution(tool_call)
1003 })
1004 }
1005
1006 #[allow(clippy::too_many_arguments)]
1012 async fn execute_single_tool(
1013 &self,
1014 context: &ExecutionContext,
1015 tool_call: ToolCall,
1016 tool_def: Option<&ToolDefinition>,
1017 trace_id: &str,
1018 act_span_id: &str,
1019 locale: Option<&str>,
1020 network_access: Option<&crate::network_access::NetworkAccessList>,
1021 visible_tool_names: Arc<HashSet<String>>,
1022 ) -> ToolCallResult {
1023 tracing::debug!(
1024 session_id = %context.session_id,
1025 turn_id = %context.turn_id,
1026 tool_name = %tool_call.name,
1027 tool_call_id = %tool_call.id,
1028 "ActAtom: executing tool"
1029 );
1030
1031 let tool_span_id = Uuid::now_v7().to_string();
1033
1034 let event_context = EventContext::from_execution_context(context).with_span(
1036 trace_id.to_string(),
1037 tool_span_id.clone(),
1038 Some(act_span_id.to_string()),
1039 );
1040
1041 let tool_start = Instant::now();
1043 let tool_call_fingerprint = tool_call_fingerprint(&tool_call);
1044
1045 let display_name = crate::localization::localized_tool_display_name(
1047 &tool_call.name,
1048 tool_def.and_then(|d| d.display_name()),
1049 locale,
1050 );
1051 let capability_attribution = tool_def.and_then(|def| {
1052 def.capability_attribution()
1053 .map(|(id, name)| (id.to_string(), name.map(str::to_string)))
1054 });
1055
1056 if let (Some(limiter), Some(ref org_id)) = (
1060 &self.outbound_tool_rate_limiter,
1061 self.context_services.org_id,
1062 ) && !limiter.check_org(org_id).await
1063 {
1064 tracing::warn!(
1065 session_id = %context.session_id,
1066 tool_name = %tool_call.name,
1067 "ActAtom: outbound tool rate limit exceeded for org"
1068 );
1069 return ToolCallResult {
1070 tool_call: tool_call.clone(),
1071 result: ToolResult {
1072 tool_call_id: tool_call.id.clone(),
1073 result: None,
1074 images: None,
1075 error: Some(
1076 "Outbound tool rate limit exceeded for this organization; back off and retry later.".to_string(),
1077 ),
1078 connection_required: None,
1079 raw_output: None,
1080 },
1081 success: false,
1082 status: "error".to_string(),
1083 connection_required: None,
1084 determinism_fatal: None,
1085 };
1086 }
1087
1088 let claim_token = if let Some(ref store) = self.durable_tool_result_store {
1091 let turn_id = context.turn_id.to_string();
1092 match store
1093 .try_claim_tool_call(
1094 &turn_id,
1095 &tool_call.id,
1096 &tool_call.name,
1097 &tool_call_fingerprint,
1098 )
1099 .await
1100 {
1101 Ok(ToolCallClaimResult::Claimed { claim_token }) => Some(claim_token),
1102
1103 Ok(ToolCallClaimResult::AlreadySettled {
1104 result_json,
1105 args_fingerprint: stored_fp,
1106 }) => {
1107 if stored_fp != tool_call_fingerprint {
1109 let err_msg = format!(
1110 "determinism violation: tool '{}' replay args fingerprint \
1111 does not match prior execution (stored={stored_fp}, \
1112 current={})",
1113 tool_call.name, tool_call_fingerprint
1114 );
1115 tracing::error!(
1116 session_id = %context.session_id,
1117 turn_id = %context.turn_id,
1118 tool_call_id = %tool_call.id,
1119 stored_fp = %stored_fp,
1120 current_fp = %tool_call_fingerprint,
1121 "ActAtom: determinism violation — replay args fingerprint mismatch"
1122 );
1123 let result_fp =
1124 tool_result_fingerprint(&tool_call.name, &ToolResult::error(&err_msg));
1125 let _ = self
1126 .event_emitter
1127 .emit(EventRequest::new(
1128 context.session_id,
1129 event_context,
1130 ToolCompletedData::failure(
1131 tool_call.id.clone(),
1132 tool_call.name.clone(),
1133 "error".to_string(),
1134 err_msg.clone(),
1135 None,
1136 )
1137 .with_fingerprints(tool_call_fingerprint.clone(), result_fp)
1138 .with_display_name(display_name.clone()),
1139 ))
1140 .await;
1141 return ToolCallResult {
1142 tool_call: tool_call.clone(),
1143 result: ToolResult {
1144 tool_call_id: tool_call.id.clone(),
1145 result: None,
1146 images: None,
1147 error: Some(err_msg.clone()),
1148 connection_required: None,
1149 raw_output: None,
1150 },
1151 success: false,
1152 status: "error".to_string(),
1153 connection_required: None,
1154 determinism_fatal: Some(err_msg),
1155 };
1156 }
1157 tracing::debug!(
1158 session_id = %context.session_id,
1159 turn_id = %context.turn_id,
1160 tool_call_id = %tool_call.id,
1161 "ActAtom: replaying already-settled tool call"
1162 );
1163 let replayed_result: ToolResult = serde_json::from_value(result_json.clone())
1165 .unwrap_or(ToolResult {
1166 tool_call_id: tool_call.id.clone(),
1167 result: Some(result_json),
1168 images: None,
1169 error: None,
1170 connection_required: None,
1171 raw_output: None,
1172 });
1173 let success = replayed_result.error.is_none();
1174 let status = if success { "success" } else { "error" };
1175 let result_fp = tool_result_fingerprint(&tool_call.name, &replayed_result);
1176 let completed_data = if success {
1177 let mut content = replayed_result
1179 .result
1180 .as_ref()
1181 .map(|r| vec![ContentPart::tool_result_text(r)])
1182 .unwrap_or_default();
1183 if let Some(ref images) = replayed_result.images {
1184 for img in images {
1185 content.push(ContentPart::Image(
1186 crate::message::ImageContentPart::from_base64(
1187 &img.base64,
1188 &img.media_type,
1189 ),
1190 ));
1191 }
1192 }
1193 ToolCompletedData::success(
1194 tool_call.id.clone(),
1195 tool_call.name.clone(),
1196 content,
1197 None,
1198 )
1199 .with_fingerprints(tool_call_fingerprint.clone(), result_fp)
1200 .with_display_name(display_name.clone())
1201 } else {
1202 ToolCompletedData::failure(
1203 tool_call.id.clone(),
1204 tool_call.name.clone(),
1205 status.to_string(),
1206 replayed_result.error.clone().unwrap_or_default(),
1207 None,
1208 )
1209 .with_fingerprints(tool_call_fingerprint.clone(), result_fp)
1210 .with_display_name(display_name.clone())
1211 };
1212 let _ = self
1213 .event_emitter
1214 .emit(EventRequest::new(
1215 context.session_id,
1216 event_context,
1217 completed_data,
1218 ))
1219 .await;
1220 let conn_req = replayed_result.connection_required.clone();
1221 return ToolCallResult {
1222 tool_call,
1223 result: replayed_result,
1224 success,
1225 status: status.to_string(),
1226 connection_required: conn_req,
1227 determinism_fatal: None,
1228 };
1229 }
1230
1231 Ok(ToolCallClaimResult::AlreadyRunning {
1232 args_fingerprint: stored_fp,
1233 }) => {
1234 if stored_fp != tool_call_fingerprint {
1237 let err_msg = format!(
1238 "determinism violation: tool '{}' args fingerprint changed \
1239 while prior claim is still running (stored={stored_fp}, \
1240 current={tool_call_fingerprint})",
1241 tool_call.name
1242 );
1243 tracing::error!(
1244 session_id = %context.session_id,
1245 turn_id = %context.turn_id,
1246 tool_call_id = %tool_call.id,
1247 stored = %stored_fp,
1248 current = %tool_call_fingerprint,
1249 "ActAtom: determinism violation — running claim fingerprint mismatch"
1250 );
1251 let result_fp =
1252 tool_result_fingerprint(&tool_call.name, &ToolResult::error(&err_msg));
1253 let _ = self
1254 .event_emitter
1255 .emit(EventRequest::new(
1256 context.session_id,
1257 event_context,
1258 ToolCompletedData::failure(
1259 tool_call.id.clone(),
1260 tool_call.name.clone(),
1261 "error".to_string(),
1262 err_msg.clone(),
1263 None,
1264 )
1265 .with_fingerprints(tool_call_fingerprint.clone(), result_fp)
1266 .with_display_name(display_name.clone()),
1267 ))
1268 .await;
1269 return ToolCallResult {
1270 tool_call: tool_call.clone(),
1271 result: ToolResult {
1272 tool_call_id: tool_call.id.clone(),
1273 result: None,
1274 images: None,
1275 error: Some(err_msg.clone()),
1276 connection_required: None,
1277 raw_output: None,
1278 },
1279 success: false,
1280 status: "error".to_string(),
1281 connection_required: None,
1282 determinism_fatal: Some(err_msg),
1283 };
1284 }
1285
1286 let sec = tool_def
1287 .map(|d| d.side_effect_class())
1288 .unwrap_or(SideEffectClass::AtMostOnce);
1289 match sec {
1290 SideEffectClass::Pure | SideEffectClass::Idempotent => {
1291 tracing::debug!(
1293 session_id = %context.session_id,
1294 tool_call_id = %tool_call.id,
1295 "ActAtom: stale running claim for idempotent tool, re-executing"
1296 );
1297 None
1298 }
1299 SideEffectClass::AtMostOnce => {
1300 tracing::warn!(
1301 session_id = %context.session_id,
1302 turn_id = %context.turn_id,
1303 tool_call_id = %tool_call.id,
1304 "ActAtom: AtMostOnce tool has stale running claim; returning interrupted result"
1305 );
1306 let _ = store
1308 .settle_tool_call(
1309 &turn_id,
1310 &tool_call.id,
1311 serde_json::Value::Null,
1312 "interrupted",
1313 Uuid::nil(), )
1315 .await;
1316 let err_msg = format!(
1317 "tool '{}' was interrupted mid-execution during a prior \
1318 worker failure; result is uncertain and was not re-run \
1319 (AtMostOnce safety)",
1320 tool_call.name
1321 );
1322 let result_fp = tool_result_fingerprint(
1323 &tool_call.name,
1324 &ToolResult::error(&err_msg),
1325 );
1326 let _ = self
1327 .event_emitter
1328 .emit(EventRequest::new(
1329 context.session_id,
1330 event_context,
1331 ToolCompletedData::failure(
1332 tool_call.id.clone(),
1333 tool_call.name.clone(),
1334 "interrupted".to_string(),
1335 err_msg.clone(),
1336 None,
1337 )
1338 .with_fingerprints(tool_call_fingerprint.clone(), result_fp)
1339 .with_display_name(display_name.clone()),
1340 ))
1341 .await;
1342 return ToolCallResult {
1343 tool_call: tool_call.clone(),
1344 result: ToolResult {
1345 tool_call_id: tool_call.id.clone(),
1346 result: None,
1347 images: None,
1348 error: Some(err_msg),
1349 connection_required: None,
1350 raw_output: None,
1351 },
1352 success: false,
1353 status: "error".to_string(),
1354 connection_required: None,
1355 determinism_fatal: None,
1356 };
1357 }
1358 }
1359 }
1360
1361 Ok(ToolCallClaimResult::DeterminismViolation {
1362 stored_fingerprint,
1363 current_fingerprint,
1364 }) => {
1365 let err_msg = format!(
1366 "determinism violation: tool '{}' args fingerprint changed \
1367 on replay (stored={stored_fingerprint}, \
1368 current={current_fingerprint})",
1369 tool_call.name
1370 );
1371 tracing::error!(
1372 session_id = %context.session_id,
1373 turn_id = %context.turn_id,
1374 tool_call_id = %tool_call.id,
1375 stored = %stored_fingerprint,
1376 current = %current_fingerprint,
1377 "ActAtom: determinism violation on claim"
1378 );
1379 let result_fp =
1380 tool_result_fingerprint(&tool_call.name, &ToolResult::error(&err_msg));
1381 let _ = self
1382 .event_emitter
1383 .emit(EventRequest::new(
1384 context.session_id,
1385 event_context,
1386 ToolCompletedData::failure(
1387 tool_call.id.clone(),
1388 tool_call.name.clone(),
1389 "error".to_string(),
1390 err_msg.clone(),
1391 None,
1392 )
1393 .with_fingerprints(tool_call_fingerprint.clone(), result_fp)
1394 .with_display_name(display_name.clone()),
1395 ))
1396 .await;
1397 return ToolCallResult {
1398 tool_call: tool_call.clone(),
1399 result: ToolResult {
1400 tool_call_id: tool_call.id.clone(),
1401 result: None,
1402 images: None,
1403 error: Some(err_msg.clone()),
1404 connection_required: None,
1405 raw_output: None,
1406 },
1407 success: false,
1408 status: "error".to_string(),
1409 connection_required: None,
1410 determinism_fatal: Some(err_msg),
1411 };
1412 }
1413
1414 Err(e) => {
1415 tracing::warn!(
1416 session_id = %context.session_id,
1417 tool_call_id = %tool_call.id,
1418 error = %e,
1419 "ActAtom: durable claim failed; proceeding without idempotency"
1420 );
1421 None
1422 }
1423 }
1424 } else {
1425 None
1426 };
1427
1428 if let Err(e) = self
1430 .event_emitter
1431 .emit(EventRequest::new(
1432 context.session_id,
1433 event_context.clone(),
1434 ToolStartedData {
1435 tool_call: tool_call.clone(),
1436 tool_call_fingerprint: Some(tool_call_fingerprint.clone()),
1437 display_name: display_name.clone(),
1438 narration: Some(self.render_tool_narration(
1439 context,
1440 tool_def,
1441 &tool_call,
1442 ToolNarrationPhase::Started,
1443 locale,
1444 )),
1445 },
1446 ))
1447 .await
1448 {
1449 tracing::warn!(
1450 session_id = %context.session_id,
1451 tool_call_id = %tool_call.id,
1452 error = %e,
1453 "ActAtom: failed to emit tool.started event"
1454 );
1455 }
1456
1457 let Some(tool_def) = tool_def else {
1459 let error_msg = format!("Tool definition not found: {}", tool_call.name);
1460 let tool_duration_ms = tool_start.elapsed().as_millis() as u64;
1461
1462 if let Err(e) = self
1464 .event_emitter
1465 .emit(EventRequest::new(
1466 context.session_id,
1467 event_context,
1468 ToolCompletedData::failure(
1469 tool_call.id.clone(),
1470 tool_call.name.clone(),
1471 "error".to_string(),
1472 error_msg.clone(),
1473 Some(tool_duration_ms),
1474 )
1475 .with_fingerprints(
1476 tool_call_fingerprint.clone(),
1477 tool_error_fingerprint(&tool_call.name, "error", &error_msg),
1478 )
1479 .with_narration(Some(self.render_tool_narration(
1480 context,
1481 None,
1482 &tool_call,
1483 ToolNarrationPhase::Failed,
1484 locale,
1485 ))),
1486 ))
1487 .await
1488 {
1489 tracing::warn!(
1490 session_id = %context.session_id,
1491 tool_call_id = %tool_call.id,
1492 error = %e,
1493 "ActAtom: failed to emit tool.completed event"
1494 );
1495 }
1496
1497 return ToolCallResult {
1498 tool_call: tool_call.clone(),
1499 result: ToolResult {
1500 tool_call_id: tool_call.id.clone(),
1501 result: None,
1502 images: None,
1503 error: Some(error_msg),
1504 connection_required: None,
1505 raw_output: None,
1506 },
1507 success: false,
1508 status: "error".to_string(),
1509 connection_required: None,
1510 determinism_fatal: None,
1511 };
1512 };
1513
1514 let mut tool_context =
1516 ToolContext::from_services(context.session_id, &self.context_services);
1517 if let Some(workspace_id) = context.workspace_id {
1522 tool_context.workspace_id = workspace_id;
1523 if let Some(store) = tool_context.file_store.take() {
1524 tool_context.file_store = Some(
1525 crate::session_files::WorkspaceScopedFileSystem::wrap(store, workspace_id),
1526 );
1527 }
1528 }
1529 if let Some(store) = tool_context.file_store.take() {
1533 tool_context.file_store = Some(crate::mount_fs::MountFs::wrap_if_needed(store));
1534 }
1535 tool_context.visible_tool_names = Some(visible_tool_names.clone());
1536 tool_context.network_access = network_access
1538 .cloned()
1539 .or_else(|| self.context_services.network_access.clone());
1540 if tool_context.event_emitter.is_none() {
1542 tool_context.event_emitter =
1543 Some(Arc::new(self.event_emitter.clone()) as Arc<dyn EventEmitter>);
1544 }
1545 tool_context.event_context = Some(event_context.clone());
1546 tool_context.tool_call_id = Some(tool_call.id.clone());
1547
1548 let call_cancellation = tokio_util::sync::CancellationToken::new();
1556 tool_context.cancellation = Some(call_cancellation.clone());
1557 let _cancel_on_call_end = call_cancellation.drop_guard();
1558
1559 let execution_tool_call = self.transform_tool_call_for_execution(tool_call.clone());
1560
1561 let (execution_tool_call, pre_block_reason) = if self.pre_tool_hooks.is_empty() {
1566 (execution_tool_call, None)
1567 } else {
1568 match act_hooks::run_pre_tool_use_hooks(
1569 &self.pre_tool_hooks,
1570 execution_tool_call.clone(),
1571 tool_def,
1572 &tool_context,
1573 )
1574 .await
1575 {
1576 act_hooks::PreToolUseDecision::Continue(updated) => (updated, None),
1577 act_hooks::PreToolUseDecision::Block {
1578 tool_call: blocked,
1579 reason,
1580 ..
1581 } => (blocked, Some(reason)),
1582 }
1583 };
1584
1585 let result = if let Some(reason) = pre_block_reason {
1586 tracing::warn!(
1587 session_id = %context.session_id,
1588 tool_call_id = %execution_tool_call.id,
1589 tool_name = %execution_tool_call.name,
1590 reason = %reason,
1591 "ActAtom: pre_tool_use hook blocked execution"
1592 );
1593 Ok(crate::tool_types::ToolResult {
1594 tool_call_id: execution_tool_call.id.clone(),
1595 result: None,
1596 images: None,
1597 error: Some(format!("blocked by pre_tool_use hook: {reason}")),
1598 connection_required: None,
1599 raw_output: None,
1600 })
1601 } else if tool_def.is_cpu_bound() {
1602 let executor = self.tool_executor.clone();
1608 let call = execution_tool_call.clone();
1609 let def = tool_def.clone();
1610 let ctx = tool_context.clone();
1611 match AbortOnDropJoinHandle::new(tokio::spawn(async move {
1612 executor.execute_with_context(&call, &def, &ctx).await
1613 }))
1614 .await
1615 {
1616 Ok(result) => result,
1617 Err(join_err) => Err(crate::error::AgentLoopError::tool(format!(
1618 "tool task failed to complete: {join_err}"
1619 ))),
1620 }
1621 } else {
1622 self.tool_executor
1623 .execute_with_context(&execution_tool_call, tool_def, &tool_context)
1624 .await
1625 };
1626
1627 match result {
1628 Ok(mut tool_result) => {
1629 act_hooks::run_post_tool_exec_hooks(
1631 &self.post_tool_hooks,
1632 &self.final_post_tool_hooks,
1633 &execution_tool_call,
1634 tool_def,
1635 &mut tool_result,
1636 &tool_context,
1637 )
1638 .await;
1639
1640 let tool_duration_ms = tool_start.elapsed().as_millis() as u64;
1641 let success = tool_result.error.is_none();
1642 let status = if success { "success" } else { "error" };
1643
1644 let completed_data = if success {
1646 let result_fingerprint = tool_result_fingerprint(&tool_call.name, &tool_result);
1647 let mut result_content = tool_result
1649 .result
1650 .as_ref()
1651 .map(|r| vec![ContentPart::tool_result_text(r)])
1652 .unwrap_or_default();
1653 if let Some(ref images) = tool_result.images {
1655 for img in images {
1656 result_content.push(ContentPart::Image(
1657 crate::message::ImageContentPart::from_base64(
1658 &img.base64,
1659 &img.media_type,
1660 ),
1661 ));
1662 }
1663 }
1664 ToolCompletedData::success(
1665 tool_call.id.clone(),
1666 tool_call.name.clone(),
1667 result_content,
1668 Some(tool_duration_ms),
1669 )
1670 .with_fingerprints(tool_call_fingerprint.clone(), result_fingerprint)
1671 .with_display_name(display_name.clone())
1672 .with_capability_attribution(
1673 capability_attribution.as_ref().map(|(id, _)| id.clone()),
1674 capability_attribution
1675 .as_ref()
1676 .and_then(|(_, name)| name.clone()),
1677 )
1678 .with_narration(Some(self.render_tool_narration(
1679 context,
1680 Some(tool_def),
1681 &tool_call,
1682 ToolNarrationPhase::Completed,
1683 locale,
1684 )))
1685 } else {
1686 let result_fingerprint = tool_result_fingerprint(&tool_call.name, &tool_result);
1687 ToolCompletedData::failure(
1688 tool_call.id.clone(),
1689 tool_call.name.clone(),
1690 status.to_string(),
1691 tool_result.error.clone().unwrap_or_default(),
1692 Some(tool_duration_ms),
1693 )
1694 .with_fingerprints(tool_call_fingerprint.clone(), result_fingerprint)
1695 .with_display_name(display_name.clone())
1696 .with_capability_attribution(
1697 capability_attribution.as_ref().map(|(id, _)| id.clone()),
1698 capability_attribution
1699 .as_ref()
1700 .and_then(|(_, name)| name.clone()),
1701 )
1702 .with_narration(Some(self.render_tool_narration(
1703 context,
1704 Some(tool_def),
1705 &tool_call,
1706 ToolNarrationPhase::Failed,
1707 locale,
1708 )))
1709 };
1710
1711 if let Err(e) = self
1712 .event_emitter
1713 .emit(EventRequest::new(
1714 context.session_id,
1715 event_context.clone(),
1716 completed_data,
1717 ))
1718 .await
1719 {
1720 tracing::warn!(
1721 session_id = %context.session_id,
1722 tool_call_id = %tool_call.id,
1723 error = %e,
1724 "ActAtom: failed to emit tool.completed event"
1725 );
1726 }
1727
1728 tracing::debug!(
1729 session_id = %context.session_id,
1730 tool_name = %tool_call.name,
1731 tool_call_id = %tool_call.id,
1732 success = %success,
1733 "ActAtom: tool execution completed"
1734 );
1735
1736 if let (Some(store), Some(token)) = (&self.durable_tool_result_store, claim_token) {
1738 let result_snapshot =
1739 serde_json::to_value(&tool_result).unwrap_or(serde_json::Value::Null);
1740 match store
1741 .settle_tool_call(
1742 &context.turn_id.to_string(),
1743 &tool_call.id,
1744 result_snapshot,
1745 "settled",
1746 token,
1747 )
1748 .await
1749 {
1750 Ok(false) => {
1751 tracing::warn!(
1752 session_id = %context.session_id,
1753 tool_call_id = %tool_call.id,
1754 "ActAtom: settle ownership check failed (task reclaimed)"
1755 );
1756 }
1757 Err(e) => {
1758 tracing::warn!(
1759 session_id = %context.session_id,
1760 tool_call_id = %tool_call.id,
1761 error = %e,
1762 "ActAtom: settle_tool_call failed"
1763 );
1764 }
1765 Ok(true) => {}
1766 }
1767 }
1768
1769 let conn_req = tool_result.connection_required.clone();
1770 ToolCallResult {
1771 tool_call,
1772 result: tool_result,
1773 success,
1774 status: status.to_string(),
1775 connection_required: conn_req,
1776 determinism_fatal: None,
1777 }
1778 }
1779 Err(e) => {
1780 let tool_duration_ms = tool_start.elapsed().as_millis() as u64;
1781 let error_msg = e.to_string();
1782
1783 if let Err(emit_err) = self
1785 .event_emitter
1786 .emit(EventRequest::new(
1787 context.session_id,
1788 event_context,
1789 ToolCompletedData::failure(
1790 tool_call.id.clone(),
1791 tool_call.name.clone(),
1792 "error".to_string(),
1793 error_msg.clone(),
1794 Some(tool_duration_ms),
1795 )
1796 .with_fingerprints(
1797 tool_call_fingerprint.clone(),
1798 tool_error_fingerprint(&tool_call.name, "error", &error_msg),
1799 )
1800 .with_display_name(display_name.clone())
1801 .with_capability_attribution(
1802 capability_attribution.as_ref().map(|(id, _)| id.clone()),
1803 capability_attribution
1804 .as_ref()
1805 .and_then(|(_, name)| name.clone()),
1806 )
1807 .with_narration(Some(self.render_tool_narration(
1808 context,
1809 Some(tool_def),
1810 &tool_call,
1811 ToolNarrationPhase::Failed,
1812 locale,
1813 ))),
1814 ))
1815 .await
1816 {
1817 tracing::warn!(
1818 session_id = %context.session_id,
1819 tool_call_id = %tool_call.id,
1820 error = %emit_err,
1821 "ActAtom: failed to emit tool.completed event"
1822 );
1823 }
1824
1825 tracing::warn!(
1826 session_id = %context.session_id,
1827 tool_name = %tool_call.name,
1828 tool_call_id = %tool_call.id,
1829 error = %e,
1830 "ActAtom: tool execution failed"
1831 );
1832
1833 ToolCallResult {
1834 tool_call: tool_call.clone(),
1835 result: ToolResult {
1836 tool_call_id: tool_call.id.clone(),
1837 result: None,
1838 images: None,
1839 error: Some(error_msg),
1840 connection_required: None,
1841 raw_output: None,
1842 },
1843 success: false,
1844 status: "error".to_string(),
1845 connection_required: None,
1846 determinism_fatal: None,
1847 }
1848 }
1849 }
1850 }
1851}
1852
1853#[cfg(test)]
1858mod tests {
1859 use super::*;
1860 use crate::test_fixtures::NoopEventEmitter;
1861 use crate::tools::ToolRegistry;
1862 use crate::typed_id::{AgentId, HarnessId, MessageId, SessionId, TurnId};
1863 use async_trait::async_trait;
1864 use everruns_core::{Capability, DisabledUtilityLlmService, Tool, ToolExecutionResult};
1865 use everruns_provider::{BuiltinTool, ClientSideTool};
1866 use serde_json::json;
1867
1868 struct ArgumentEchoTool;
1869
1870 struct NarratingGrepTool;
1871
1872 struct HumanIntentFixtureHook;
1873
1874 impl crate::capabilities::ToolCallHook for HumanIntentFixtureHook {
1875 fn narration(
1876 &self,
1877 _tool_def: Option<&ToolDefinition>,
1878 tool_call: &ToolCall,
1879 _phase: crate::tool_narration::ToolNarrationPhase,
1880 _locale: Option<&str>,
1881 _ctx: crate::tool_narration::ToolNarrationContext<'_>,
1882 ) -> Option<String> {
1883 crate::tool_types::human_intent(&tool_call.arguments).map(str::to_string)
1884 }
1885
1886 fn transform_for_execution(&self, mut tool_call: ToolCall) -> ToolCall {
1887 tool_call.arguments = tool_call.execution_arguments();
1888 tool_call
1889 }
1890 }
1891
1892 #[async_trait]
1893 impl crate::tools::Tool for NarratingGrepTool {
1894 fn name(&self) -> &str {
1895 "grep_files"
1896 }
1897
1898 fn description(&self) -> &str {
1899 "Search files"
1900 }
1901
1902 fn parameters_schema(&self) -> serde_json::Value {
1903 json!({"type": "object"})
1904 }
1905
1906 async fn execute(&self, _arguments: serde_json::Value) -> ToolExecutionResult {
1907 ToolExecutionResult::success(json!({}))
1908 }
1909
1910 fn narrate(
1911 &self,
1912 tool_call: &ToolCall,
1913 phase: crate::tool_narration::ToolNarrationPhase,
1914 locale: Option<&str>,
1915 _ctx: crate::tool_narration::ToolNarrationContext<'_>,
1916 ) -> Option<String> {
1917 Some(crate::tool_narration::narrate_grep_files(
1918 &tool_call.arguments,
1919 phase,
1920 locale,
1921 ))
1922 }
1923 }
1924
1925 struct NarratingCapability;
1926
1927 #[async_trait]
1928 impl Capability for NarratingCapability {
1929 fn id(&self) -> &str {
1930 "narrating_test"
1931 }
1932
1933 fn name(&self) -> &str {
1934 "Narrating test"
1935 }
1936
1937 fn description(&self) -> &str {
1938 "Test-only narration capability"
1939 }
1940
1941 fn tools(&self) -> Vec<Box<dyn Tool>> {
1942 vec![Box::new(NarratingGrepTool)]
1943 }
1944 }
1945
1946 #[async_trait]
1947 impl crate::tools::Tool for ArgumentEchoTool {
1948 fn name(&self) -> &str {
1949 "argument_echo"
1950 }
1951
1952 fn description(&self) -> &str {
1953 "returns the execution arguments"
1954 }
1955
1956 fn parameters_schema(&self) -> serde_json::Value {
1957 json!({
1958 "type": "object",
1959 "properties": {
1960 "value": { "type": "string" }
1961 }
1962 })
1963 }
1964
1965 async fn execute(&self, arguments: serde_json::Value) -> ToolExecutionResult {
1966 ToolExecutionResult::success(arguments)
1967 }
1968 }
1969
1970 #[test]
1971 fn grouped_headline_uses_tool_owned_narration_for_repeated_actions() {
1972 use crate::capabilities::{Capability, CapabilityNarrationHook};
1973
1974 let capability: Arc<dyn Capability> = Arc::new(NarratingCapability);
1975 let tool_definitions = capability
1976 .tools()
1977 .into_iter()
1978 .map(|tool| tool.to_definition())
1979 .collect::<Vec<_>>();
1980 let tool_map = tool_definitions
1981 .iter()
1982 .map(|tool_def| (tool_def.name(), tool_def))
1983 .collect::<std::collections::HashMap<_, _>>();
1984 let atom = ActAtom::new(ToolRegistry::new(), NoopEventEmitter)
1985 .with_tool_call_hooks(vec![Arc::new(CapabilityNarrationHook(capability))]);
1986 let context = ExecutionContext::new(SessionId::new(), TurnId::new(), MessageId::new());
1987 let tool_calls = vec![
1988 ToolCall {
1989 id: "grep-1".to_string(),
1990 name: "grep_files".to_string(),
1991 arguments: json!({ "pattern": "full_name" }),
1992 },
1993 ToolCall {
1994 id: "grep-2".to_string(),
1995 name: "grep_files".to_string(),
1996 arguments: json!({ "pattern": "login" }),
1997 },
1998 ];
1999
2000 assert_eq!(
2001 atom.render_group_headline(
2002 &context,
2003 &tool_calls,
2004 &tool_map,
2005 ToolNarrationPhase::Started,
2006 None,
2007 )
2008 .as_deref(),
2009 Some("Searching files twice")
2010 );
2011 assert_eq!(
2012 atom.render_group_headline(
2013 &context,
2014 &tool_calls,
2015 &tool_map,
2016 ToolNarrationPhase::Completed,
2017 None,
2018 )
2019 .as_deref(),
2020 Some("Searched files twice")
2021 );
2022 }
2023
2024 #[test]
2029 fn registry_owned_tool_narrates_itself_without_a_capability_hook() {
2030 let mut registry = ToolRegistry::new();
2031 registry.register_boxed(Box::new(NarratingGrepTool));
2032 let atom = ActAtom::new(ToolRegistry::new(), NoopEventEmitter)
2033 .with_tool_registry(Arc::new(registry));
2034 let context = ExecutionContext::new(SessionId::new(), TurnId::new(), MessageId::new());
2035 let tool_call = ToolCall {
2036 id: "grep-1".to_string(),
2037 name: "grep_files".to_string(),
2038 arguments: json!({ "pattern": "full_name" }),
2039 };
2040
2041 assert_eq!(
2042 atom.render_tool_narration(
2043 &context,
2044 None,
2045 &tool_call,
2046 ToolNarrationPhase::Started,
2047 None,
2048 ),
2049 "Searching files for full_name"
2050 );
2051 }
2052
2053 struct UtilityLlmContextProbeTool;
2054
2055 #[async_trait]
2056 impl crate::tools::Tool for UtilityLlmContextProbeTool {
2057 fn name(&self) -> &str {
2058 "utility_llm_context_probe"
2059 }
2060
2061 fn description(&self) -> &str {
2062 "checks whether the utility LLM service is present in tool context"
2063 }
2064
2065 fn parameters_schema(&self) -> serde_json::Value {
2066 json!({
2067 "type": "object",
2068 "properties": {}
2069 })
2070 }
2071
2072 async fn execute(&self, _arguments: serde_json::Value) -> ToolExecutionResult {
2073 ToolExecutionResult::tool_error("context required")
2074 }
2075
2076 async fn execute_with_context(
2077 &self,
2078 _arguments: serde_json::Value,
2079 context: &crate::tool_context::ToolContext,
2080 ) -> ToolExecutionResult {
2081 ToolExecutionResult::success(json!({
2082 "utility_llm_service": context.utility_llm_service.is_some(),
2083 "configured": context
2084 .utility_llm_service
2085 .as_ref()
2086 .is_some_and(|service| service.is_configured()),
2087 }))
2088 }
2089
2090 fn requires_context(&self) -> bool {
2091 true
2092 }
2093 }
2094
2095 #[derive(Default)]
2097 struct SchedObservations {
2098 class_inflight: std::collections::HashMap<String, usize>,
2100 class_max: std::collections::HashMap<String, usize>,
2102 global_inflight: usize,
2104 global_max: usize,
2106 }
2107
2108 struct RecordingTool {
2111 name: String,
2112 class: Option<String>,
2113 obs: Arc<std::sync::Mutex<SchedObservations>>,
2114 }
2115
2116 #[async_trait]
2117 impl crate::tools::Tool for RecordingTool {
2118 fn name(&self) -> &str {
2119 &self.name
2120 }
2121 fn description(&self) -> &str {
2122 "records scheduling order"
2123 }
2124 fn parameters_schema(&self) -> serde_json::Value {
2125 json!({ "type": "object", "properties": {} })
2126 }
2127 async fn execute(&self, _arguments: serde_json::Value) -> ToolExecutionResult {
2128 {
2130 let mut obs = self.obs.lock().unwrap();
2131 obs.global_inflight += 1;
2132 let g = obs.global_inflight;
2133 if g > obs.global_max {
2134 obs.global_max = g;
2135 }
2136 if let Some(class) = &self.class {
2137 let n = obs.class_inflight.entry(class.clone()).or_default();
2138 *n += 1;
2139 let cur = *n;
2140 let m = obs.class_max.entry(class.clone()).or_default();
2141 if cur > *m {
2142 *m = cur;
2143 }
2144 }
2145 }
2146 tokio::time::sleep(std::time::Duration::from_millis(20)).await;
2148 {
2150 let mut obs = self.obs.lock().unwrap();
2151 obs.global_inflight -= 1;
2152 if let Some(class) = &self.class
2153 && let Some(n) = obs.class_inflight.get_mut(class)
2154 {
2155 *n -= 1;
2156 }
2157 }
2158 ToolExecutionResult::success(json!({ "tool": self.name }))
2159 }
2160 }
2161
2162 struct CancellationProbeTool {
2163 started: Arc<tokio::sync::Notify>,
2164 dropped_tx: Arc<std::sync::Mutex<Option<tokio::sync::oneshot::Sender<()>>>>,
2165 }
2166
2167 impl CancellationProbeTool {
2168 fn new(
2169 started: Arc<tokio::sync::Notify>,
2170 dropped_tx: tokio::sync::oneshot::Sender<()>,
2171 ) -> Self {
2172 Self {
2173 started,
2174 dropped_tx: Arc::new(std::sync::Mutex::new(Some(dropped_tx))),
2175 }
2176 }
2177 }
2178
2179 #[async_trait]
2180 impl crate::tools::Tool for CancellationProbeTool {
2181 fn name(&self) -> &str {
2182 "cancellation_probe"
2183 }
2184
2185 fn description(&self) -> &str {
2186 "waits until cancelled"
2187 }
2188
2189 fn parameters_schema(&self) -> serde_json::Value {
2190 json!({ "type": "object", "properties": {} })
2191 }
2192
2193 async fn execute(&self, _arguments: serde_json::Value) -> ToolExecutionResult {
2194 struct DropSignal {
2195 tx: Arc<std::sync::Mutex<Option<tokio::sync::oneshot::Sender<()>>>>,
2196 }
2197
2198 impl Drop for DropSignal {
2199 fn drop(&mut self) {
2200 if let Ok(mut guard) = self.tx.lock()
2201 && let Some(tx) = guard.take()
2202 {
2203 let _ = tx.send(());
2204 }
2205 }
2206 }
2207
2208 let _drop_signal = DropSignal {
2209 tx: self.dropped_tx.clone(),
2210 };
2211 self.started.notify_one();
2212 std::future::pending::<()>().await;
2213 unreachable!("pending cancellation probe should only finish by cancellation")
2214 }
2215 }
2216
2217 fn recording_tool_def(name: &str, class: Option<&str>, cpu_bound: bool) -> ToolDefinition {
2219 let mut hints = crate::tool_types::ToolHints::default();
2220 if let Some(class) = class {
2221 hints = hints.with_concurrency_class(class);
2222 }
2223 if cpu_bound {
2224 hints = hints.with_cpu_bound(true);
2225 }
2226 ToolDefinition::Builtin(BuiltinTool {
2227 name: name.to_string(),
2228 display_name: None,
2229 description: "records scheduling order".to_string(),
2230 parameters: json!({ "type": "object", "properties": {} }),
2231 policy: Default::default(),
2232 category: None,
2233 deferrable: Default::default(),
2234 hints,
2235 full_parameters: None,
2236 })
2237 }
2238
2239 #[tokio::test]
2240 async fn test_act_atom_empty_tool_calls() {
2241 let executor = ToolRegistry::with_defaults();
2242 let event_emitter = NoopEventEmitter;
2243 let atom = ActAtom::new(executor, event_emitter);
2244
2245 let context = ExecutionContext::new(SessionId::new(), TurnId::new(), MessageId::new());
2246 let input = ActInput {
2247 org_id: Some(1),
2248 context,
2249 harness_id: HarnessId::from_seed(1),
2250 agent_id: Some(AgentId::new()),
2251 tool_calls: vec![],
2252 tool_definitions: vec![],
2253 locale: None,
2254 blueprint_id: None,
2255 network_access: None,
2256 parallel_tool_calls: None,
2257 };
2258
2259 let result = atom.execute(input).await.unwrap();
2260
2261 assert!(result.completed);
2262 assert!(result.results.is_empty());
2263 assert_eq!(result.success_count, 0);
2264 assert_eq!(result.error_count, 0);
2265 }
2266
2267 #[tokio::test]
2268 async fn test_act_atom_threads_utility_llm_service_to_tool_context() {
2269 let mut executor = ToolRegistry::with_defaults();
2270 executor.register(UtilityLlmContextProbeTool);
2271 let event_emitter = NoopEventEmitter;
2272 let atom = ActAtom::new(executor, event_emitter)
2273 .with_utility_llm_service(Arc::new(DisabledUtilityLlmService));
2274
2275 let context = ExecutionContext::new(SessionId::new(), TurnId::new(), MessageId::new());
2276 let input = ActInput {
2277 org_id: Some(1),
2278 context,
2279 harness_id: HarnessId::from_seed(1),
2280 agent_id: Some(AgentId::new()),
2281 tool_calls: vec![ToolCall {
2282 id: "call_1".to_string(),
2283 name: "utility_llm_context_probe".to_string(),
2284 arguments: json!({}),
2285 }],
2286 tool_definitions: vec![ToolDefinition::Builtin(BuiltinTool {
2287 name: "utility_llm_context_probe".to_string(),
2288 display_name: None,
2289 description: "checks context".to_string(),
2290 parameters: json!({
2291 "type": "object",
2292 "properties": {}
2293 }),
2294 policy: Default::default(),
2295 category: None,
2296 deferrable: Default::default(),
2297 hints: crate::tool_types::ToolHints::default(),
2298 full_parameters: None,
2299 })],
2300 locale: None,
2301 blueprint_id: None,
2302 network_access: None,
2303 parallel_tool_calls: None,
2304 };
2305
2306 let result = atom.execute(input).await.unwrap();
2307
2308 assert_eq!(result.success_count, 1);
2309 let payload = result.results[0].result.result.as_ref().unwrap();
2310 assert_eq!(payload["utility_llm_service"], true);
2311 assert_eq!(payload["configured"], false);
2312 }
2313
2314 #[tokio::test]
2319 async fn test_act_atom_schedules_batch_by_concurrency_class() {
2320 let obs = Arc::new(std::sync::Mutex::new(SchedObservations::default()));
2321
2322 let mut executor = ToolRegistry::new();
2323 executor.register(RecordingTool {
2324 name: "writer_a".to_string(),
2325 class: Some("ws".to_string()),
2326 obs: obs.clone(),
2327 });
2328 executor.register(RecordingTool {
2329 name: "writer_b".to_string(),
2330 class: Some("ws".to_string()),
2331 obs: obs.clone(),
2332 });
2333 executor.register(RecordingTool {
2334 name: "reader".to_string(),
2335 class: None,
2336 obs: obs.clone(),
2337 });
2338
2339 let atom = ActAtom::new(executor, NoopEventEmitter);
2340 let context = ExecutionContext::new(SessionId::new(), TurnId::new(), MessageId::new());
2341
2342 let input = ActInput {
2345 org_id: Some(1),
2346 context,
2347 harness_id: HarnessId::from_seed(1),
2348 agent_id: Some(AgentId::new()),
2349 tool_calls: vec![
2350 ToolCall {
2351 id: "call_a".to_string(),
2352 name: "writer_a".to_string(),
2353 arguments: json!({}),
2354 },
2355 ToolCall {
2356 id: "call_r".to_string(),
2357 name: "reader".to_string(),
2358 arguments: json!({}),
2359 },
2360 ToolCall {
2361 id: "call_b".to_string(),
2362 name: "writer_b".to_string(),
2363 arguments: json!({}),
2364 },
2365 ],
2366 tool_definitions: vec![
2367 recording_tool_def("writer_a", Some("ws"), false),
2368 recording_tool_def("reader", None, false),
2369 recording_tool_def("writer_b", Some("ws"), true),
2370 ],
2371 locale: None,
2372 blueprint_id: None,
2373 network_access: None,
2374 parallel_tool_calls: None,
2375 };
2376
2377 let result = atom.execute(input).await.unwrap();
2378
2379 assert_eq!(result.success_count, 3, "all three tools should succeed");
2381 let names: Vec<&str> = result
2383 .results
2384 .iter()
2385 .map(|r| r.tool_call.name.as_str())
2386 .collect();
2387 assert_eq!(names, vec!["writer_a", "reader", "writer_b"]);
2388
2389 let obs = obs.lock().unwrap();
2390 assert_eq!(
2393 obs.class_max.get("ws").copied(),
2394 Some(1),
2395 "same-class tools must serialize"
2396 );
2397 assert!(
2400 obs.global_max >= 2,
2401 "independent tool should run concurrently with the class group (global_max={})",
2402 obs.global_max
2403 );
2404 }
2405
2406 struct DetachedWorkTool {
2410 cancelled_tx: Arc<std::sync::Mutex<Option<tokio::sync::oneshot::Sender<()>>>>,
2411 }
2412
2413 impl DetachedWorkTool {
2414 fn new(cancelled_tx: tokio::sync::oneshot::Sender<()>) -> Self {
2415 Self {
2416 cancelled_tx: Arc::new(std::sync::Mutex::new(Some(cancelled_tx))),
2417 }
2418 }
2419 }
2420
2421 #[async_trait]
2422 impl crate::tools::Tool for DetachedWorkTool {
2423 fn name(&self) -> &str {
2424 "detached_work"
2425 }
2426
2427 fn description(&self) -> &str {
2428 "spawns work that outlives the call unless cancelled"
2429 }
2430
2431 fn parameters_schema(&self) -> serde_json::Value {
2432 json!({ "type": "object", "properties": {} })
2433 }
2434
2435 fn requires_context(&self) -> bool {
2436 true
2437 }
2438
2439 async fn execute(&self, _arguments: serde_json::Value) -> ToolExecutionResult {
2440 ToolExecutionResult::tool_error("requires context")
2441 }
2442
2443 async fn execute_with_context(
2444 &self,
2445 _arguments: serde_json::Value,
2446 context: &crate::tool_context::ToolContext,
2447 ) -> ToolExecutionResult {
2448 let token = context
2449 .cancellation
2450 .clone()
2451 .expect("act must supply a cancellation token");
2452 assert!(!token.is_cancelled(), "token is live during the call");
2453 let tx = self.cancelled_tx.clone();
2454 tokio::spawn(async move {
2455 token.cancelled().await;
2456 if let Ok(mut guard) = tx.lock()
2457 && let Some(tx) = guard.take()
2458 {
2459 let _ = tx.send(());
2460 }
2461 });
2462 ToolExecutionResult::success(json!({ "spawned": true }))
2463 }
2464 }
2465
2466 #[tokio::test]
2470 async fn test_act_atom_cancels_detached_tool_work_when_the_call_ends() {
2471 let (cancelled_tx, cancelled_rx) = tokio::sync::oneshot::channel();
2472
2473 let mut executor = ToolRegistry::new();
2474 executor.register(DetachedWorkTool::new(cancelled_tx));
2475
2476 let atom = ActAtom::new(executor, NoopEventEmitter);
2477 let context = ExecutionContext::new(SessionId::new(), TurnId::new(), MessageId::new());
2478 let input = ActInput {
2479 org_id: Some(1),
2480 context,
2481 harness_id: HarnessId::from_seed(1),
2482 agent_id: Some(AgentId::new()),
2483 tool_calls: vec![ToolCall {
2484 id: "call_1".to_string(),
2485 name: "detached_work".to_string(),
2486 arguments: json!({}),
2487 }],
2488 tool_definitions: vec![recording_tool_def("detached_work", None, false)],
2489 locale: None,
2490 blueprint_id: None,
2491 network_access: None,
2492 parallel_tool_calls: None,
2493 };
2494
2495 atom.execute(input).await.expect("act should succeed");
2496
2497 tokio::time::timeout(std::time::Duration::from_secs(1), cancelled_rx)
2498 .await
2499 .expect("detached work should be cancelled once the call ends")
2500 .expect("cancellation signal should be sent");
2501 }
2502
2503 #[tokio::test]
2504 async fn test_act_atom_cancels_detached_tool_work_when_the_turn_is_cancelled() {
2505 let started = Arc::new(tokio::sync::Notify::new());
2506 let (dropped_tx, dropped_rx) = tokio::sync::oneshot::channel();
2507
2508 let mut executor = ToolRegistry::new();
2509 executor.register(CancellationProbeTool::new(started.clone(), dropped_tx));
2510
2511 let atom = ActAtom::new(executor, NoopEventEmitter);
2512 let context = ExecutionContext::new(SessionId::new(), TurnId::new(), MessageId::new());
2513 let input = ActInput {
2514 org_id: Some(1),
2515 context,
2516 harness_id: HarnessId::from_seed(1),
2517 agent_id: Some(AgentId::new()),
2518 tool_calls: vec![ToolCall {
2519 id: "call_1".to_string(),
2520 name: "cancellation_probe".to_string(),
2521 arguments: json!({}),
2522 }],
2523 tool_definitions: vec![recording_tool_def("cancellation_probe", None, true)],
2524 locale: None,
2525 blueprint_id: None,
2526 network_access: None,
2527 parallel_tool_calls: None,
2528 };
2529
2530 let act_task = tokio::spawn(async move { atom.execute(input).await });
2531 started.notified().await;
2532 act_task.abort();
2533 assert!(act_task.await.unwrap_err().is_cancelled());
2534
2535 tokio::time::timeout(std::time::Duration::from_secs(1), dropped_rx)
2537 .await
2538 .expect("tool future should be dropped when the turn is cancelled")
2539 .expect("drop signal should be sent");
2540 }
2541
2542 #[tokio::test]
2543 async fn test_act_atom_aborts_cpu_bound_tool_task_on_cancellation() {
2544 let started = Arc::new(tokio::sync::Notify::new());
2545 let (dropped_tx, dropped_rx) = tokio::sync::oneshot::channel();
2546
2547 let mut executor = ToolRegistry::new();
2548 executor.register(CancellationProbeTool::new(started.clone(), dropped_tx));
2549
2550 let atom = ActAtom::new(executor, NoopEventEmitter);
2551 let context = ExecutionContext::new(SessionId::new(), TurnId::new(), MessageId::new());
2552 let input = ActInput {
2553 org_id: Some(1),
2554 context,
2555 harness_id: HarnessId::from_seed(1),
2556 agent_id: Some(AgentId::new()),
2557 tool_calls: vec![ToolCall {
2558 id: "call_1".to_string(),
2559 name: "cancellation_probe".to_string(),
2560 arguments: json!({}),
2561 }],
2562 tool_definitions: vec![recording_tool_def("cancellation_probe", None, true)],
2563 locale: None,
2564 blueprint_id: None,
2565 network_access: None,
2566 parallel_tool_calls: None,
2567 };
2568
2569 let act_task = tokio::spawn(async move { atom.execute(input).await });
2570 started.notified().await;
2571 act_task.abort();
2572 assert!(act_task.await.unwrap_err().is_cancelled());
2573
2574 tokio::time::timeout(std::time::Duration::from_secs(1), dropped_rx)
2575 .await
2576 .expect("cpu-bound tool task should be aborted when ActAtom is cancelled")
2577 .expect("drop signal should be sent by cancelled tool future");
2578 }
2579
2580 #[tokio::test]
2583 async fn test_act_atom_parallel_tool_calls_false_serializes_everything() {
2584 let obs = Arc::new(std::sync::Mutex::new(SchedObservations::default()));
2585 let mut executor = ToolRegistry::new();
2586 for name in ["t0", "t1", "t2"] {
2587 executor.register(RecordingTool {
2588 name: name.to_string(),
2589 class: None,
2590 obs: obs.clone(),
2591 });
2592 }
2593 let atom = ActAtom::new(executor, NoopEventEmitter);
2594 let context = ExecutionContext::new(SessionId::new(), TurnId::new(), MessageId::new());
2595 let input = ActInput {
2596 org_id: Some(1),
2597 context,
2598 harness_id: HarnessId::from_seed(1),
2599 agent_id: Some(AgentId::new()),
2600 tool_calls: vec![
2601 ToolCall {
2602 id: "c0".to_string(),
2603 name: "t0".to_string(),
2604 arguments: json!({}),
2605 },
2606 ToolCall {
2607 id: "c1".to_string(),
2608 name: "t1".to_string(),
2609 arguments: json!({}),
2610 },
2611 ToolCall {
2612 id: "c2".to_string(),
2613 name: "t2".to_string(),
2614 arguments: json!({}),
2615 },
2616 ],
2617 tool_definitions: vec![
2618 recording_tool_def("t0", None, false),
2619 recording_tool_def("t1", None, false),
2620 recording_tool_def("t2", None, false),
2621 ],
2622 locale: None,
2623 blueprint_id: None,
2624 network_access: None,
2625 parallel_tool_calls: Some(false),
2626 };
2627
2628 let result = atom.execute(input).await.unwrap();
2629 assert_eq!(result.success_count, 3);
2630 assert_eq!(
2631 obs.lock().unwrap().global_max,
2632 1,
2633 "parallel_tool_calls=false must serialize the whole batch"
2634 );
2635 }
2636
2637 #[tokio::test]
2638 async fn test_act_atom_tool_not_found() {
2639 let executor = ToolRegistry::with_defaults();
2640 let event_emitter = NoopEventEmitter;
2641 let atom = ActAtom::new(executor, event_emitter);
2642
2643 let context = ExecutionContext::new(SessionId::new(), TurnId::new(), MessageId::new());
2644 let input = ActInput {
2645 org_id: Some(1),
2646 context,
2647 harness_id: HarnessId::from_seed(1),
2648 agent_id: Some(AgentId::new()),
2649 tool_calls: vec![ToolCall {
2650 id: "call_1".to_string(),
2651 name: "nonexistent_tool".to_string(),
2652 arguments: json!({}),
2653 }],
2654 tool_definitions: vec![],
2655 locale: None,
2656 blueprint_id: None,
2657 network_access: None,
2658 parallel_tool_calls: None,
2659 };
2660
2661 let result = atom.execute(input).await.unwrap();
2662
2663 assert!(result.completed);
2664 assert_eq!(result.results.len(), 1);
2665 assert!(!result.results[0].success);
2666 assert_eq!(result.results[0].status, "error");
2667 assert!(
2668 result.results[0]
2669 .result
2670 .error
2671 .as_ref()
2672 .unwrap()
2673 .contains("not found")
2674 );
2675 }
2676
2677 #[tokio::test]
2678 async fn test_act_atom_uses_tool_call_hooks_for_execution_arguments() {
2679 let mut executor = ToolRegistry::new();
2680 executor.register(ArgumentEchoTool);
2681 let tool_def = executor.get("argument_echo").unwrap().to_definition();
2682 let emitter = crate::test_fixtures::TestEventEmitter::new();
2683 let atom = ActAtom::new(executor, emitter.clone())
2684 .with_tool_call_hooks(vec![std::sync::Arc::new(HumanIntentFixtureHook)]);
2685
2686 let context = ExecutionContext::new(SessionId::new(), TurnId::new(), MessageId::new());
2687 let input = ActInput {
2688 org_id: Some(1),
2689 context,
2690 harness_id: HarnessId::from_seed(1),
2691 agent_id: Some(AgentId::new()),
2692 tool_calls: vec![ToolCall {
2693 id: "call_1".to_string(),
2694 name: "argument_echo".to_string(),
2695 arguments: json!({
2696 "value": "visible",
2697 "human_intent": "Echoing test arguments"
2698 }),
2699 }],
2700 tool_definitions: vec![tool_def],
2701 locale: None,
2702 blueprint_id: None,
2703 network_access: None,
2704 parallel_tool_calls: None,
2705 };
2706
2707 let result = atom.execute(input).await.unwrap();
2708
2709 assert!(result.results[0].success);
2710 assert_eq!(
2711 result.results[0].result.result,
2712 Some(json!({ "value": "visible" }))
2713 );
2714
2715 let events = emitter.events().await;
2716 assert_eq!(
2717 events
2718 .iter()
2719 .map(|event| event.event_type.as_str())
2720 .collect::<Vec<_>>(),
2721 vec![
2722 "act.started",
2723 "tool.started",
2724 "tool.completed",
2725 "act.completed",
2726 ],
2727 "all hosts must observe the engine-owned phase order",
2728 );
2729 let act_started = events
2730 .iter()
2731 .find(|event| event.event_type == "act.started")
2732 .expect("act.started event");
2733 let crate::events::EventData::ActStarted(data) = &act_started.data else {
2734 panic!("expected act.started data");
2735 };
2736 assert_eq!(data.headline.as_deref(), Some("Echoing test arguments"));
2737 assert_eq!(
2738 data.tool_calls[0].narration.as_deref(),
2739 Some("Echoing test arguments")
2740 );
2741
2742 let tool_started = events
2743 .iter()
2744 .find(|event| event.event_type == "tool.started")
2745 .expect("tool.started event");
2746 let crate::events::EventData::ToolStarted(data) = &tool_started.data else {
2747 panic!("expected tool.started data");
2748 };
2749 let started_fingerprint = data
2750 .tool_call_fingerprint
2751 .as_ref()
2752 .expect("tool.started call fingerprint");
2753 assert_eq!(data.narration.as_deref(), Some("Echoing test arguments"));
2754
2755 let tool_completed = events
2756 .iter()
2757 .find(|event| event.event_type == "tool.completed")
2758 .expect("tool.completed event");
2759 let crate::events::EventData::ToolCompleted(data) = &tool_completed.data else {
2760 panic!("expected tool.completed data");
2761 };
2762 assert_eq!(
2763 data.tool_call_fingerprint.as_ref(),
2764 Some(started_fingerprint)
2765 );
2766 assert!(data.tool_result_fingerprint.is_some());
2767 assert_eq!(data.narration.as_deref(), Some("Echoing test arguments"));
2768 }
2769
2770 #[tokio::test]
2771 async fn test_act_atom_strips_human_intent_from_client_tool_calls() {
2772 let executor = ToolRegistry::new();
2773 let emitter = crate::test_fixtures::TestEventEmitter::new();
2774 let atom = ActAtom::new(executor, emitter)
2775 .with_tool_call_hooks(vec![std::sync::Arc::new(HumanIntentFixtureHook)]);
2776
2777 let context = ExecutionContext::new(SessionId::new(), TurnId::new(), MessageId::new());
2778 let input = ActInput {
2779 org_id: Some(1),
2780 context,
2781 harness_id: HarnessId::from_seed(1),
2782 agent_id: Some(AgentId::new()),
2783 tool_calls: vec![ToolCall {
2784 id: "call_client".to_string(),
2785 name: "browser_click".to_string(),
2786 arguments: json!({
2787 "selector": "#btn",
2788 "human_intent": "Clicking approve"
2789 }),
2790 }],
2791 tool_definitions: vec![ToolDefinition::ClientSide(ClientSideTool::new(
2792 "browser_click",
2793 "Click button",
2794 json!({
2795 "type": "object",
2796 "properties": {
2797 "selector": {"type": "string"}
2798 },
2799 "required": ["selector"]
2800 }),
2801 ))],
2802 locale: None,
2803 blueprint_id: None,
2804 network_access: None,
2805 parallel_tool_calls: None,
2806 };
2807
2808 let result = atom.execute(input).await.unwrap();
2809
2810 assert_eq!(result.client_tool_calls.len(), 1);
2811 assert_eq!(
2812 result.client_tool_calls[0].arguments,
2813 json!({ "selector": "#btn" })
2814 );
2815 }
2816
2817 #[test]
2818 fn test_act_result_connection_required_serialization() {
2819 let result = ActResult {
2820 results: vec![ToolCallResult {
2821 tool_call: ToolCall {
2822 id: "call_1".to_string(),
2823 name: "daytona_create_sandbox".to_string(),
2824 arguments: json!({}),
2825 },
2826 result: ToolResult {
2827 tool_call_id: "call_1".to_string(),
2828 result: Some(json!({"connection_required": "daytona"})),
2829 images: None,
2830 error: None,
2831 connection_required: Some(ConnectionRequired::provider_only("daytona")),
2832 raw_output: None,
2833 },
2834 success: false,
2835 status: "success".to_string(),
2836 connection_required: Some(ConnectionRequired::provider_only("daytona")),
2837 determinism_fatal: None,
2838 }],
2839 completed: true,
2840 success_count: 0,
2841 error_count: 0,
2842 waiting_for_tool_results: true,
2843 waiting_for_url_elicitation: false,
2844 blocked: false,
2845 client_tool_calls: vec![],
2846 client_tool_definitions: vec![],
2847 };
2848 let json_str = serde_json::to_string(&result).unwrap();
2849 let parsed: ActResult = serde_json::from_str(&json_str).unwrap();
2850
2851 assert!(parsed.waiting_for_tool_results);
2852 assert_eq!(
2853 parsed.results[0].connection_required,
2854 result.results[0].connection_required
2855 );
2856 }
2857
2858 #[test]
2859 fn test_act_result_backward_compat_deserialization() {
2860 let json_str = r#"{
2862 "results": [],
2863 "completed": true,
2864 "success_count": 0,
2865 "error_count": 0
2866 }"#;
2867 let parsed: ActResult = serde_json::from_str(json_str).unwrap();
2868
2869 assert!(!parsed.waiting_for_tool_results);
2870 assert!(parsed.client_tool_calls.is_empty());
2871 }
2872
2873 #[tokio::test]
2876 async fn test_outbound_tool_rate_limiter_blocks_execution() {
2877 use crate::typed_id::OrgId;
2878
2879 struct DenyAll;
2880 #[async_trait]
2881 impl crate::tool_execution::OutboundToolRateLimiter for DenyAll {
2882 async fn check_org(&self, _org_id: &OrgId) -> bool {
2883 false
2884 }
2885 }
2886
2887 let mut executor = ToolRegistry::with_defaults();
2888 executor.register(ArgumentEchoTool);
2889 let atom = ActAtom::new(executor, NoopEventEmitter)
2890 .with_org_id(OrgId::from_seed(1))
2891 .with_outbound_tool_rate_limiter(Arc::new(DenyAll));
2892
2893 let context = ExecutionContext::new(SessionId::new(), TurnId::new(), MessageId::new());
2894 let input = ActInput {
2895 org_id: Some(1),
2896 context,
2897 harness_id: HarnessId::from_seed(1),
2898 agent_id: Some(AgentId::new()),
2899 tool_calls: vec![ToolCall {
2900 id: "call_1".to_string(),
2901 name: "argument_echo".to_string(),
2902 arguments: json!({"value": "should_not_reach"}),
2903 }],
2904 tool_definitions: vec![ToolDefinition::Builtin(BuiltinTool {
2905 name: "argument_echo".to_string(),
2906 display_name: None,
2907 description: "echo".to_string(),
2908 parameters: json!({"type": "object"}),
2909 policy: Default::default(),
2910 category: None,
2911 deferrable: Default::default(),
2912 hints: crate::tool_types::ToolHints::default(),
2913 full_parameters: None,
2914 })],
2915 locale: None,
2916 blueprint_id: None,
2917 network_access: None,
2918 parallel_tool_calls: None,
2919 };
2920
2921 let result = atom.execute(input).await.unwrap();
2922
2923 assert_eq!(result.success_count, 0);
2924 assert_eq!(result.error_count, 1);
2925 let tool_result = &result.results[0];
2926 assert!(!tool_result.success);
2927 assert_eq!(tool_result.status, "error");
2928 assert!(
2929 tool_result
2930 .result
2931 .error
2932 .as_deref()
2933 .unwrap_or("")
2934 .contains("rate limit exceeded")
2935 );
2936 assert!(tool_result.result.result.is_none());
2937 }
2938
2939 #[tokio::test]
2941 async fn test_outbound_tool_rate_limiter_allows_execution() {
2942 use crate::typed_id::OrgId;
2943
2944 struct AllowAll;
2945 #[async_trait]
2946 impl crate::tool_execution::OutboundToolRateLimiter for AllowAll {
2947 async fn check_org(&self, _org_id: &OrgId) -> bool {
2948 true
2949 }
2950 }
2951
2952 let mut executor = ToolRegistry::with_defaults();
2953 executor.register(ArgumentEchoTool);
2954 let atom = ActAtom::new(executor, NoopEventEmitter)
2955 .with_org_id(OrgId::from_seed(1))
2956 .with_outbound_tool_rate_limiter(Arc::new(AllowAll));
2957
2958 let context = ExecutionContext::new(SessionId::new(), TurnId::new(), MessageId::new());
2959 let input = ActInput {
2960 org_id: Some(1),
2961 context,
2962 harness_id: HarnessId::from_seed(1),
2963 agent_id: Some(AgentId::new()),
2964 tool_calls: vec![ToolCall {
2965 id: "call_1".to_string(),
2966 name: "argument_echo".to_string(),
2967 arguments: json!({"value": "hello"}),
2968 }],
2969 tool_definitions: vec![ToolDefinition::Builtin(BuiltinTool {
2970 name: "argument_echo".to_string(),
2971 display_name: None,
2972 description: "echo".to_string(),
2973 parameters: json!({"type": "object"}),
2974 policy: Default::default(),
2975 category: None,
2976 deferrable: Default::default(),
2977 hints: crate::tool_types::ToolHints::default(),
2978 full_parameters: None,
2979 })],
2980 locale: None,
2981 blueprint_id: None,
2982 network_access: None,
2983 parallel_tool_calls: None,
2984 };
2985
2986 let result = atom.execute(input).await.unwrap();
2987
2988 assert_eq!(result.success_count, 1);
2989 assert_eq!(result.error_count, 0);
2990 }
2991
2992 use crate::tool_types::{SideEffectClass, ToolHints};
2997 use crate::{durability::DurableToolResultStore, durability::ToolCallClaimResult};
2998 use std::collections::HashMap;
2999 use std::sync::Mutex;
3000
3001 #[derive(Default)]
3002 struct InMemoryDurableStore {
3003 rows: Mutex<HashMap<(String, String), StoreRow>>,
3004 }
3005
3006 #[derive(Clone)]
3007 struct StoreRow {
3008 status: String,
3009 result_json: serde_json::Value,
3010 args_fingerprint: String,
3011 #[allow(dead_code)]
3012 claim_token: Uuid,
3013 }
3014
3015 #[async_trait]
3016 impl DurableToolResultStore for InMemoryDurableStore {
3017 async fn try_claim_tool_call(
3018 &self,
3019 turn_id: &str,
3020 tool_call_id: &str,
3021 _tool_name: &str,
3022 args_fingerprint: &str,
3023 ) -> crate::error::Result<ToolCallClaimResult> {
3024 let key = (turn_id.to_string(), tool_call_id.to_string());
3025 let mut rows = self.rows.lock().unwrap();
3026 if let Some(row) = rows.get(&key) {
3027 match row.status.as_str() {
3028 "settled" => {
3029 if row.args_fingerprint != args_fingerprint {
3030 return Ok(ToolCallClaimResult::DeterminismViolation {
3031 stored_fingerprint: row.args_fingerprint.clone(),
3032 current_fingerprint: args_fingerprint.to_string(),
3033 });
3034 }
3035 return Ok(ToolCallClaimResult::AlreadySettled {
3036 result_json: row.result_json.clone(),
3037 args_fingerprint: row.args_fingerprint.clone(),
3038 });
3039 }
3040 _ => {
3041 return Ok(ToolCallClaimResult::AlreadyRunning {
3042 args_fingerprint: row.args_fingerprint.clone(),
3043 });
3044 }
3045 }
3046 }
3047 let token = Uuid::new_v4();
3048 rows.insert(
3049 key,
3050 StoreRow {
3051 status: "running".to_string(),
3052 result_json: serde_json::Value::Null,
3053 args_fingerprint: args_fingerprint.to_string(),
3054 claim_token: token,
3055 },
3056 );
3057 Ok(ToolCallClaimResult::Claimed { claim_token: token })
3058 }
3059
3060 async fn settle_tool_call(
3061 &self,
3062 turn_id: &str,
3063 tool_call_id: &str,
3064 result_json: serde_json::Value,
3065 status: &str,
3066 _claim_token: Uuid,
3067 ) -> crate::error::Result<bool> {
3068 let key = (turn_id.to_string(), tool_call_id.to_string());
3069 let mut rows = self.rows.lock().unwrap();
3070 if let Some(row) = rows.get_mut(&key) {
3071 row.status = status.to_string();
3072 row.result_json = result_json;
3073 return Ok(true);
3074 }
3075 Ok(false)
3076 }
3077
3078 async fn get_tool_call_status(
3079 &self,
3080 turn_id: &str,
3081 tool_call_id: &str,
3082 ) -> crate::error::Result<Option<crate::durability::DurableToolCallStatus>> {
3083 let key = (turn_id.to_string(), tool_call_id.to_string());
3084 let rows = self.rows.lock().unwrap();
3085 Ok(rows.get(&key).map(|row| match row.status.as_str() {
3086 "settled" => crate::durability::DurableToolCallStatus::Settled {
3087 result_json: row.result_json.clone(),
3088 },
3089 "interrupted" => crate::durability::DurableToolCallStatus::Interrupted {
3090 result_json: Some(row.result_json.clone()),
3091 },
3092 _ => crate::durability::DurableToolCallStatus::Running,
3093 }))
3094 }
3095 }
3096
3097 fn make_act_input_with_store(
3098 tool_call: ToolCall,
3099 tool_defs: Vec<ToolDefinition>,
3100 context: ExecutionContext,
3101 ) -> ActInput {
3102 ActInput {
3103 org_id: None,
3104 context,
3105 harness_id: HarnessId::from_seed(1),
3106 agent_id: Some(AgentId::new()),
3107 tool_calls: vec![tool_call],
3108 tool_definitions: tool_defs,
3109 locale: None,
3110 blueprint_id: None,
3111 network_access: None,
3112 parallel_tool_calls: None,
3113 }
3114 }
3115
3116 fn arg_echo_tool_def(side_effect: SideEffectClass) -> ToolDefinition {
3117 ToolDefinition::Builtin(BuiltinTool {
3118 name: "argument_echo".to_string(),
3119 display_name: None,
3120 description: "echo".to_string(),
3121 parameters: json!({"type": "object"}),
3122 policy: Default::default(),
3123 category: None,
3124 deferrable: Default::default(),
3125 hints: ToolHints::default().with_side_effect_class(side_effect),
3126 full_parameters: None,
3127 })
3128 }
3129
3130 #[tokio::test]
3132 async fn test_idempotency_first_execution_claims_and_settles() {
3133 let store = Arc::new(InMemoryDurableStore::default());
3134 let mut executor = ToolRegistry::with_defaults();
3135 executor.register(ArgumentEchoTool);
3136 let atom =
3137 ActAtom::new(executor, NoopEventEmitter).with_durable_tool_result_store(store.clone());
3138
3139 let context = ExecutionContext::new(SessionId::new(), TurnId::new(), MessageId::new());
3140 let tc = ToolCall {
3141 id: "c1".to_string(),
3142 name: "argument_echo".to_string(),
3143 arguments: json!({"value": "hello"}),
3144 };
3145 let input = make_act_input_with_store(
3146 tc,
3147 vec![arg_echo_tool_def(SideEffectClass::AtMostOnce)],
3148 context,
3149 );
3150
3151 let result = atom.execute(input).await.unwrap();
3152 assert_eq!(result.success_count, 1);
3153 assert_eq!(result.error_count, 0);
3154
3155 let rows = store.rows.lock().unwrap();
3157 let row = rows.values().next().unwrap();
3158 assert_eq!(row.status, "settled");
3159 }
3160
3161 #[tokio::test]
3163 async fn test_idempotency_replay_already_settled() {
3164 use crate::tool_fingerprint::tool_call_fingerprint;
3165
3166 let store = Arc::new(InMemoryDurableStore::default());
3167 let tc = ToolCall {
3168 id: "c1".to_string(),
3169 name: "argument_echo".to_string(),
3170 arguments: json!({"value": "hello"}),
3171 };
3172 let fp = tool_call_fingerprint(&tc);
3173
3174 {
3176 let stored_result = serde_json::to_value(ToolResult {
3177 tool_call_id: "c1".to_string(),
3178 result: Some(json!({"value": "hello"})),
3179 images: None,
3180 error: None,
3181 connection_required: None,
3182 raw_output: None,
3183 })
3184 .unwrap();
3185 store.rows.lock().unwrap().insert(
3186 (
3187 "turn_00000000000000000000000000000000".to_string(),
3188 "c1".to_string(),
3189 ),
3190 StoreRow {
3191 status: "settled".to_string(),
3192 result_json: stored_result,
3193 args_fingerprint: fp,
3194 claim_token: Uuid::new_v4(),
3195 },
3196 );
3197 }
3198
3199 let mut executor = ToolRegistry::with_defaults();
3200 executor.register(ArgumentEchoTool);
3201 let atom =
3202 ActAtom::new(executor, NoopEventEmitter).with_durable_tool_result_store(store.clone());
3203
3204 let context = ExecutionContext::new(
3205 SessionId::new(),
3206 TurnId::from_uuid(Uuid::nil()),
3207 MessageId::new(),
3208 );
3209 let input = make_act_input_with_store(
3210 tc,
3211 vec![arg_echo_tool_def(SideEffectClass::AtMostOnce)],
3212 context,
3213 );
3214
3215 let result = atom.execute(input).await.unwrap();
3216 assert_eq!(result.success_count, 1, "replay should count as success");
3217 assert_eq!(result.error_count, 0);
3218 }
3219
3220 #[tokio::test]
3222 async fn test_idempotency_at_most_once_stale_running_returns_interrupted() {
3223 use crate::tool_fingerprint::tool_call_fingerprint;
3224
3225 let store = Arc::new(InMemoryDurableStore::default());
3226 let tc = ToolCall {
3227 id: "c1".to_string(),
3228 name: "argument_echo".to_string(),
3229 arguments: json!({"value": "x"}),
3230 };
3231 let fp = tool_call_fingerprint(&tc);
3232
3233 store.rows.lock().unwrap().insert(
3235 (
3236 "turn_00000000000000000000000000000000".to_string(),
3237 "c1".to_string(),
3238 ),
3239 StoreRow {
3240 status: "running".to_string(),
3241 result_json: serde_json::Value::Null,
3242 args_fingerprint: fp,
3243 claim_token: Uuid::new_v4(),
3244 },
3245 );
3246
3247 let mut executor = ToolRegistry::with_defaults();
3248 executor.register(ArgumentEchoTool);
3249 let atom =
3250 ActAtom::new(executor, NoopEventEmitter).with_durable_tool_result_store(store.clone());
3251
3252 let context = ExecutionContext::new(
3253 SessionId::new(),
3254 TurnId::from_uuid(Uuid::nil()),
3255 MessageId::new(),
3256 );
3257 let input = make_act_input_with_store(
3258 tc,
3259 vec![arg_echo_tool_def(SideEffectClass::AtMostOnce)],
3260 context,
3261 );
3262
3263 let result = atom.execute(input).await.unwrap();
3264 assert_eq!(
3265 result.error_count, 1,
3266 "AtMostOnce stale running should error"
3267 );
3268 let err = result.results[0].result.error.as_deref().unwrap_or("");
3269 assert!(
3270 err.contains("interrupted"),
3271 "error should mention interrupted: {err}"
3272 );
3273
3274 let rows = store.rows.lock().unwrap();
3276 let row = rows.values().next().unwrap();
3277 assert_eq!(row.status, "interrupted");
3278 }
3279
3280 #[tokio::test]
3282 async fn test_idempotency_idempotent_tool_stale_running_reexecutes() {
3283 use crate::tool_fingerprint::tool_call_fingerprint;
3284
3285 let store = Arc::new(InMemoryDurableStore::default());
3286 let tc = ToolCall {
3287 id: "c1".to_string(),
3288 name: "argument_echo".to_string(),
3289 arguments: json!({"value": "x"}),
3290 };
3291 let fp = tool_call_fingerprint(&tc);
3292
3293 store.rows.lock().unwrap().insert(
3294 (
3295 "turn_00000000000000000000000000000000".to_string(),
3296 "c1".to_string(),
3297 ),
3298 StoreRow {
3299 status: "running".to_string(),
3300 result_json: serde_json::Value::Null,
3301 args_fingerprint: fp,
3302 claim_token: Uuid::new_v4(),
3303 },
3304 );
3305
3306 let mut executor = ToolRegistry::with_defaults();
3307 executor.register(ArgumentEchoTool);
3308 let atom =
3309 ActAtom::new(executor, NoopEventEmitter).with_durable_tool_result_store(store.clone());
3310
3311 let context = ExecutionContext::new(
3312 SessionId::new(),
3313 TurnId::from_uuid(Uuid::nil()),
3314 MessageId::new(),
3315 );
3316 let input = make_act_input_with_store(
3317 tc,
3318 vec![arg_echo_tool_def(SideEffectClass::Idempotent)],
3319 context,
3320 );
3321
3322 let result = atom.execute(input).await.unwrap();
3323 assert_eq!(
3324 result.success_count, 1,
3325 "Idempotent should re-execute successfully"
3326 );
3327 assert_eq!(result.error_count, 0);
3328 }
3329}