Skip to main content

everruns_engine/execution/
act.rs

1//! ActAtom - Atom for scheduled tool execution
2//!
3//! This atom handles:
4//! 1. Emitting act.started event
5//! 2. Executing the batch of tool calls via the [`tool_scheduler`] (with
6//!    tool.started/completed events). Calls run concurrently by default, but
7//!    calls that share a [`crate::tool_types::ToolHints::concurrency_class`] are
8//!    serialized to avoid mutation races, total concurrency is capped, and
9//!    `cpu_bound` tools are offloaded to their own task.
10//! 3. Handling errors, timeouts, and cancellations as "normal" results
11//! 4. Emitting act.completed event
12//! 5. Returning all tool results (success, error, timeout, or cancelled)
13//!
14//! Tool results are emitted as `tool.completed` events and returned in ActResult.
15//! Messages are derived from events - no separate message storage is needed.
16//!
17//! Note: OTel instrumentation is handled via the event-listener pattern.
18//! tool.started/completed events are emitted by this atom, and OtelEventListener
19//! creates the appropriate gen-ai spans from those events.
20//!
21//! NOTES from Python spec:
22//! - Tool calls run concurrently by default; the scheduler serializes only
23//!   conflicting (same-concurrency-class) calls. See [`tool_scheduler`].
24//! - Error from tool call is not an error for the whole Act, error from tool is "normal" result
25//! - Tool invocations should be timeouted, timeout is also "normal" result from tool
26//! - Exit of act should have all tool calls finished (successfully or with error/timeout)
27//! - Act and each tool call should emit start/end events
28//! - Act and each tool call should be cancellable, and this is also "normal" result
29
30use serde::{Deserialize, Serialize};
31use std::collections::HashSet;
32use std::future::Future;
33use std::pin::Pin;
34use std::sync::Arc;
35use std::task::{Context, Poll};
36use web_time::Instant;
37
38use super::ExecutionContext;
39use super::act_hooks::{self, PostActHook};
40use super::tool_scheduler;
41use crate::error::Result;
42use crate::events::{
43    ActCompletedData, ActStartedData, EventContext, EventRequest, ToolCompletedData,
44    ToolStartedData,
45};
46use crate::message::ContentPart;
47use crate::phase_effects::{PhaseEffectEmitter, PhaseEffectSink};
48use crate::tool_fingerprint::{
49    tool_call_fingerprint, tool_error_fingerprint, tool_result_fingerprint,
50};
51use crate::tool_narration::{
52    GroupHeadlineAction, ToolNarrationContext, ToolNarrationPhase,
53    render_tool_narration_with_locale, summarize_group_actions, tool_call_for_group_summary,
54};
55use crate::tool_types::ConnectionRequired;
56use crate::tool_types::{SideEffectClass, ToolCall, ToolDefinition, ToolResult};
57use crate::typed_id::{AgentId, HarnessId};
58use crate::{
59    durability::DurableToolResultStore, durability::ToolCallClaimResult,
60    event_emitter::EventEmitter, execution_loading::AgentStore, execution_loading::SessionStore,
61    session_files::SessionFileSystem, tool_context::ToolContext, tool_execution::ToolExecutor,
62};
63use uuid::Uuid;
64
65/// A Tokio task handle that aborts its task if the parent future is dropped
66/// before the task completes. Tokio detaches a bare [`tokio::task::JoinHandle`]
67/// on drop, but tool execution must not outlive Act cancellation.
68struct AbortOnDropJoinHandle<T> {
69    handle: everruns_contracts::rt::JoinHandle<T>,
70}
71
72impl<T> AbortOnDropJoinHandle<T> {
73    fn new(handle: everruns_contracts::rt::JoinHandle<T>) -> Self {
74        Self { handle }
75    }
76}
77
78impl<T> Future for AbortOnDropJoinHandle<T> {
79    type Output = std::result::Result<T, everruns_contracts::rt::JoinError>;
80
81    fn poll(mut self: Pin<&mut Self>, cx: &mut Context<'_>) -> Poll<Self::Output> {
82        Pin::new(&mut self.handle).poll(cx)
83    }
84}
85
86impl<T> Drop for AbortOnDropJoinHandle<T> {
87    fn drop(&mut self) {
88        if !self.handle.is_finished() {
89            self.handle.abort();
90        }
91    }
92}
93
94// ============================================================================
95// Input and Output Types
96// ============================================================================
97
98/// Input for ActAtom
99#[derive(Debug, Clone, Serialize, Deserialize)]
100pub struct ActInput {
101    /// Organization ID for scoped data access.
102    #[serde(skip_serializing_if = "Option::is_none")]
103    pub org_id: Option<i64>,
104    /// Atom execution context
105    pub context: ExecutionContext,
106    /// Harness ID (needed for scheduling follow-up reason activity)
107    pub harness_id: HarnessId,
108    /// Agent ID (needed for scheduling follow-up reason activity, optional)
109    #[serde(skip_serializing_if = "Option::is_none")]
110    pub agent_id: Option<AgentId>,
111    /// Tool calls to execute
112    pub tool_calls: Vec<ToolCall>,
113    /// Available tool definitions for resolution
114    pub tool_definitions: Vec<ToolDefinition>,
115    /// Resolved locale for backend-authored tool narration and labels.
116    #[serde(skip_serializing_if = "Option::is_none")]
117    pub locale: Option<String>,
118    /// Blueprint ID for blueprint-backed sessions. When set, act_activity
119    /// loads tools from the blueprint instead of from agent/harness capabilities.
120    #[serde(skip_serializing_if = "Option::is_none")]
121    pub blueprint_id: Option<String>,
122    /// Merged network access list (harness ∩ agent ∩ session) for URL filtering.
123    #[serde(default, skip_serializing_if = "Option::is_none")]
124    pub network_access: Option<crate::network_access::NetworkAccessList>,
125    /// Mirrors the request's `parallel_tool_calls`. `Some(false)` forces the
126    /// act scheduler to execute this batch strictly sequentially; `None` or
127    /// `Some(true)` uses the default class-aware concurrent schedule.
128    #[serde(default, skip_serializing_if = "Option::is_none")]
129    pub parallel_tool_calls: Option<bool>,
130}
131
132/// Result of a single tool call execution
133#[derive(Debug, Clone, Serialize, Deserialize)]
134pub struct ToolCallResult {
135    /// The original tool call
136    pub tool_call: ToolCall,
137    /// The result of the tool call
138    pub result: ToolResult,
139    /// Whether the execution was successful
140    pub success: bool,
141    /// Status: "success", "error", "timeout", or "cancelled"
142    pub status: String,
143    /// If set, the tool requires a connection before it can execute.
144    #[serde(default, skip_serializing_if = "Option::is_none")]
145    pub connection_required: Option<ConnectionRequired>,
146    /// Determinism violation message. When Some, ActAtom::execute returns Err to fail the
147    /// durable workflow fast rather than continuing with a corrupted replay.
148    #[serde(default, skip_serializing_if = "Option::is_none")]
149    pub determinism_fatal: Option<String>,
150}
151
152/// Result of the ActAtom
153#[derive(Debug, Clone, Serialize, Deserialize)]
154pub struct ActResult {
155    /// Results for all tool calls
156    pub results: Vec<ToolCallResult>,
157    /// Whether all tool calls completed (regardless of success/failure)
158    pub completed: bool,
159    /// Number of successful tool calls
160    pub success_count: u32,
161    /// Number of failed tool calls
162    pub error_count: u32,
163    /// When true, the act emitted client-side tool calls (connection setup,
164    /// client-side tools, etc.) and the worker should pause until tool results
165    /// arrive. Workers check this single flag — they never need to know *why*
166    /// the act paused.
167    #[serde(default)]
168    pub waiting_for_tool_results: bool,
169    /// True when the pause is specifically a URL mode elicitation waiting on a
170    /// human to open a link. Kept apart from the generic flag because only a
171    /// client that renders the consent card can answer it — see
172    /// `plan_after_act`, which will not hold a turn for a card nobody can draw.
173    #[serde(default, skip_serializing_if = "is_false")]
174    pub waiting_for_url_elicitation: bool,
175    /// True when execution stopped before tool execution because a dependency was archived or deleted.
176    #[serde(default, skip_serializing_if = "is_false")]
177    pub blocked: bool,
178    /// Client-side tool calls that were NOT executed by ActAtom but need to be
179    /// sent to the client. Populated by ActAtom's partitioning logic, consumed
180    /// by ClientSideToolHook.
181    #[serde(default, skip_serializing_if = "Vec::is_empty")]
182    pub client_tool_calls: Vec<ToolCall>,
183    /// Tool definitions for the client-side tool calls (for narration/display).
184    #[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
192// ============================================================================
193// ActAtom
194// ============================================================================
195
196/// Atom that executes a batch of tool calls via the [`tool_scheduler`]
197///
198/// This atom:
199/// 1. Emits act.started event
200/// 2. Schedules all tool calls (emitting tool.started/completed for each):
201///    concurrent by default, serialized within a concurrency class, capped, and
202///    with `cpu_bound` tools offloaded to their own task
203/// 3. Handles errors, timeouts, and cancellations gracefully
204/// 4. Emits act.completed event
205/// 5. Returns comprehensive results for all tools
206///
207/// Tool results are emitted as events and returned in ActResult.
208/// Messages are derived from events by the message store.
209pub struct ActAtom<T, E>
210where
211    T: ToolExecutor,
212    E: PhaseEffectSink,
213{
214    // Held as `Arc` so individual `cpu_bound` tool calls can be offloaded to
215    // their own task (`tokio::spawn`) without borrowing `self` for `'static`.
216    tool_executor: Arc<T>,
217    event_emitter: PhaseEffectEmitter<E>,
218    /// Runtime-owned service snapshot cloned into every per-call ToolContext.
219    context_services: crate::tool_context::ToolContextServices,
220    /// Optional per-org outbound tool-call rate limiter (TM-TOOL-009).
221    /// When present, each tool call increments the org counter; calls that
222    /// exceed the per-org window return a tool error rather than a hard failure.
223    outbound_tool_rate_limiter: Option<Arc<dyn crate::tool_execution::OutboundToolRateLimiter>>,
224    /// Per-tool-call idempotency store (EVE-530). When present, each tool call
225    /// is claimed before dispatch and settled after completion so that reclaiming
226    /// workers can skip already-settled calls and avoid double side-effects for
227    /// `AtMostOnce` tools.
228    durable_tool_result_store: Option<Arc<dyn DurableToolResultStore>>,
229    /// Post-act hooks that run after tool execution completes.
230    /// Hooks inspect the result and may emit events (e.g. tool.call_requested).
231    hooks: Vec<Box<dyn PostActHook>>,
232    /// Post-tool-exec hooks (capability-contributed): run after each individual
233    /// tool execution. Capabilities register these via `post_tool_exec_hooks()`.
234    post_tool_hooks: Vec<Arc<dyn act_hooks::PostToolExecHook>>,
235    /// Pre-tool-use hooks (capability-contributed): run before each individual
236    /// tool execution. Capabilities wire these in via the user-hooks
237    /// adapter chain (see `crate::hook_adapter`). Hooks can mutate the
238    /// `ToolCall` (returning `Continue`) or refuse execution
239    /// (returning `Block`).
240    pre_tool_hooks: Vec<Arc<dyn act_hooks::PreToolUseHook>>,
241    /// Tool-call hooks (capability-contributed): inspect model-authored tool
242    /// calls for UI narration and transform calls before actual execution.
243    tool_call_hooks: Vec<Arc<dyn crate::capabilities::ToolCallHook>>,
244    /// Final post-tool-exec hooks (infrastructure): run after capability hooks.
245    /// Always registered, cannot be removed by capabilities (EVE-225).
246    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    /// Create a new ActAtom with default hooks (ConnectionSetup + ClientSideTool).
255    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    /// Create a new ActAtom with a file store for context-aware tools
271    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    /// Replace the complete runtime-owned service snapshot used for every
294    /// per-call [`ToolContext`]. Production hosts should prefer this over
295    /// assembling individual services on the atom.
296    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    /// Add a custom post-act hook.
305    pub fn with_hook(mut self, hook: Box<dyn PostActHook>) -> Self {
306        self.hooks.push(hook);
307        self
308    }
309
310    /// Add a runtime-owned final post-tool hook. Hosts use this for portable
311    /// policies that must run after capability hooks but before the hard output
312    /// limit.
313    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    /// Default hooks, in order. FormElicitation and ToolApprovalPause append synthetic calls, so
320    /// they must precede ClientSideTool, which emits `tool.call_requested` for client-side calls.
321    fn default_hooks() -> Vec<Box<dyn PostActHook>> {
322        vec![
323            Box::new(act_hooks::ConnectionSetupHook),
324            Box::new(act_hooks::UrlElicitationHook),
325            Box::new(act_hooks::FormElicitationHook),
326            Box::new(act_hooks::ToolApprovalPauseHook),
327            Box::new(act_hooks::ClientSideToolHook),
328        ]
329    }
330
331    /// Default final post-tool-exec hooks (infrastructure, always-on).
332    /// These run after all capability-contributed hooks and cannot be removed.
333    fn default_final_hooks() -> Vec<Arc<dyn act_hooks::PostToolExecHook>> {
334        vec![Arc::new(act_hooks::OutputHardLimitHook)]
335    }
336
337    /// Set the session storage store on this atom
338    pub fn with_storage_store(
339        mut self,
340        store: Arc<dyn crate::session_services::SessionStorageStore>,
341    ) -> Self {
342        self.context_services.storage_store = Some(store);
343        self
344    }
345
346    /// Set the image artifact store on this atom
347    pub fn with_image_store(
348        mut self,
349        store: Arc<dyn crate::image_services::ImageArtifactStore>,
350    ) -> Self {
351        self.context_services.image_store = Some(store);
352        self
353    }
354
355    /// Set the provider credential store on this atom
356    pub fn with_provider_credential_store(
357        mut self,
358        store: Arc<dyn crate::connection_services::ProviderCredentialStore>,
359    ) -> Self {
360        self.context_services.provider_credential_store = Some(store);
361        self
362    }
363
364    /// Set the utility LLM service on this atom.
365    pub fn with_utility_llm_service(mut self, service: Arc<dyn crate::UtilityLlmService>) -> Self {
366        self.context_services.utility_llm_service = Some(service);
367        self
368    }
369
370    /// Set the scoped-MCP tool invoker on this atom (guardrails `mcp` check).
371    pub fn with_mcp_invoker(mut self, invoker: Arc<dyn crate::McpToolInvoker>) -> Self {
372        self.context_services.mcp_invoker = Some(invoker);
373        self
374    }
375
376    /// Set the outbound egress service on this atom.
377    pub fn with_egress_service(mut self, service: Arc<dyn crate::EgressService>) -> Self {
378        self.context_services.egress_service = Some(service);
379        self
380    }
381
382    /// Set the user connection resolver on this atom
383    pub fn with_connection_resolver(
384        mut self,
385        resolver: Arc<dyn crate::connection_services::UserConnectionResolver>,
386    ) -> Self {
387        self.context_services.connection_resolver = Some(resolver);
388        self
389    }
390
391    /// Set session store for context-aware tools.
392    pub fn with_session_store(mut self, store: Arc<dyn SessionStore>) -> Self {
393        self.context_services.session_store = Some(store);
394        self
395    }
396
397    /// Set agent store for context-aware tools.
398    pub fn with_agent_store(mut self, store: Arc<dyn AgentStore>) -> Self {
399        self.context_services.agent_store = Some(store);
400        self
401    }
402
403    /// Set session schedule store for scheduling tools.
404    pub fn with_schedule_store(
405        mut self,
406        store: Arc<dyn crate::session_services::SessionScheduleStore>,
407    ) -> Self {
408        self.context_services.schedule_store = Some(store);
409        self
410    }
411
412    /// Set platform store for org-level management tools.
413    pub fn with_subagent_delegate(
414        mut self,
415        store: Arc<dyn crate::subagent_delegation::SubagentSessionDelegate>,
416    ) -> Self {
417        self.context_services.subagent_delegate = Some(store);
418        self
419    }
420
421    /// Set leased resource store for lifecycle-managed provider resources.
422    pub fn with_leased_resource_store(
423        mut self,
424        store: Arc<dyn crate::session_services::LeasedResourceStore>,
425    ) -> Self {
426        self.context_services.leased_resource_store = Some(store);
427        self
428    }
429
430    /// Set session resource registry.
431    pub fn with_session_resource_registry(
432        mut self,
433        registry: Arc<dyn crate::session_services::SessionResourceRegistry>,
434    ) -> Self {
435        self.context_services.session_resource_registry = Some(registry);
436        self
437    }
438
439    /// Add a session task registry passed to tool contexts.
440    pub fn with_session_task_registry(
441        mut self,
442        registry: Arc<dyn crate::session_task::SessionTaskRegistry>,
443    ) -> Self {
444        self.context_services.session_task_registry = Some(registry);
445        self
446    }
447
448    pub fn with_capability_registry(
449        mut self,
450        registry: crate::capabilities::CapabilityRegistry,
451    ) -> Self {
452        self.context_services.capability_registry = Some(registry);
453        self
454    }
455
456    /// Set the active built-in tool registry for meta-tools like `spawn_background`.
457    pub fn with_tool_registry(mut self, registry: Arc<crate::tools::ToolRegistry>) -> Self {
458        self.context_services.tool_registry = Some(registry);
459        self
460    }
461
462    /// Add capability-contributed post-tool-exec hooks.
463    /// Callers should pass hooks from the *active* capabilities for this session,
464    /// not from the full platform registry.
465    pub fn with_post_tool_hooks(
466        mut self,
467        hooks: Vec<Arc<dyn act_hooks::PostToolExecHook>>,
468    ) -> Self {
469        self.post_tool_hooks.extend(hooks);
470        self
471    }
472
473    /// Add capability-contributed pre-tool-use hooks. Pre-hooks fire before
474    /// each tool call and can mutate or block it; see
475    /// `act_hooks::PreToolUseHook` and `knowledge/runtime-resources/user-hooks.md`.
476    pub fn with_pre_tool_hooks(mut self, hooks: Vec<Arc<dyn act_hooks::PreToolUseHook>>) -> Self {
477        self.pre_tool_hooks.extend(hooks);
478        self
479    }
480
481    pub fn with_tool_call_hooks(
482        mut self,
483        hooks: Vec<Arc<dyn crate::capabilities::ToolCallHook>>,
484    ) -> Self {
485        self.tool_call_hooks.extend(hooks);
486        self
487    }
488
489    /// Set org ID for org-scoped operations.
490    pub fn with_org_id(mut self, org_id: crate::typed_id::OrgId) -> Self {
491        self.context_services.org_id = Some(org_id);
492        self
493    }
494
495    /// Set the merged network access list for URL filtering in tools.
496    pub fn with_network_access(
497        mut self,
498        network_access: Option<crate::network_access::NetworkAccessList>,
499    ) -> Self {
500        self.context_services.network_access = network_access;
501        self
502    }
503
504    /// Set the budget checker for the check_budget tool.
505    pub fn with_budget_checker(
506        mut self,
507        checker: Arc<dyn crate::tool_execution::BudgetChecker>,
508    ) -> Self {
509        self.context_services.budget_checker = Some(checker);
510        self
511    }
512
513    /// Set the internal payment authority for paid capability tools.
514    pub fn with_payment_authority(
515        mut self,
516        authority: Arc<dyn crate::tool_execution::PaymentAuthority>,
517    ) -> Self {
518        self.context_services.payment_authority = Some(authority);
519        self
520    }
521
522    /// Set the authority used to authorize detached peer-session creation.
523    pub fn with_session_creation_authority(
524        mut self,
525        authority: Arc<dyn crate::delegation_services::SessionCreationAuthority>,
526    ) -> Self {
527        self.context_services.session_creation_authority = Some(authority);
528        self
529    }
530
531    /// Set the per-org outbound tool-call rate limiter (TM-TOOL-009).
532    pub fn with_outbound_tool_rate_limiter(
533        mut self,
534        limiter: Arc<dyn crate::tool_execution::OutboundToolRateLimiter>,
535    ) -> Self {
536        self.outbound_tool_rate_limiter = Some(limiter);
537        self
538    }
539
540    /// Set the durable per-tool-call idempotency store (EVE-530).
541    pub fn with_durable_tool_result_store(
542        mut self,
543        store: Arc<dyn DurableToolResultStore>,
544    ) -> Self {
545        self.durable_tool_result_store = Some(store);
546        self
547    }
548
549    /// Set the durable subagent spawn handle store (EVE-535).
550    pub fn with_subagent_spawn_store(
551        mut self,
552        store: Arc<dyn crate::delegation_services::SubagentSpawnStore>,
553    ) -> Self {
554        self.context_services.subagent_spawn_store = Some(store);
555        self
556    }
557
558    /// Set the resolved subagent nesting policy for tool contexts.
559    pub fn with_subagent_nesting_policy(
560        mut self,
561        policy: crate::delegation_services::SubagentNestingPolicy,
562    ) -> Self {
563        self.context_services.subagent_nesting_policy = policy;
564        self
565    }
566
567    /// Set the live reasoning-effort handle (EVE-595). When set, each tool's
568    /// `ToolContext` receives a clone so a tool can change the reasoning effort
569    /// mid-turn for subsequent LLM steps in the same turn.
570    pub fn with_reasoning_effort_handle(
571        mut self,
572        handle: crate::tool_context::ReasoningEffortHandle,
573    ) -> Self {
574        self.context_services.reasoning_effort_handle = Some(handle);
575        self
576    }
577}
578
579impl<T, E> ActAtom<T, E>
580where
581    T: ToolExecutor + Send + Sync + 'static,
582    E: EventEmitter + Send + Sync + 'static,
583{
584    /// Stable phase name used by logs and durable activity adapters.
585    pub fn name(&self) -> &'static str {
586        "act"
587    }
588
589    /// Execute one scheduled tool-call batch through injected contracts.
590    pub async fn execute(&self, input: ActInput) -> Result<ActResult> {
591        let ActInput {
592            context,
593            tool_calls,
594            tool_definitions,
595            locale,
596            network_access,
597            parallel_tool_calls,
598            .. // agent_id/org_id not needed here, just passed through workflow
599        } = input;
600
601        // Partition tool calls: server-side tools get executed, client-side tools
602        // are stored on ActResult for the ClientSideToolHook to emit.
603        let (server_tool_calls, client_tool_calls): (Vec<_>, Vec<_>) = tool_calls
604            .into_iter()
605            .partition(|tc| act_hooks::runs_on_server(tc, &tool_definitions));
606
607        let client_tool_calls: Vec<_> = client_tool_calls
608            .into_iter()
609            .map(|tool_call| self.transform_tool_call_for_execution(tool_call))
610            .collect();
611
612        let client_tool_definitions: Vec<_> = if client_tool_calls.is_empty() {
613            vec![]
614        } else {
615            tool_definitions
616                .iter()
617                .filter(|td| {
618                    if let ToolDefinition::ClientSide(ct) = td {
619                        client_tool_calls.iter().any(|tc| tc.name == ct.name)
620                    } else {
621                        false
622                    }
623                })
624                .cloned()
625                .collect()
626        };
627
628        if server_tool_calls.is_empty() && client_tool_calls.is_empty() {
629            return Ok(ActResult {
630                results: vec![],
631                completed: true,
632                success_count: 0,
633                error_count: 0,
634                waiting_for_tool_results: false,
635                waiting_for_url_elicitation: false,
636                blocked: false,
637                client_tool_calls: vec![],
638                client_tool_definitions: vec![],
639            });
640        }
641
642        // If only client-side tools (no server-side), skip tool execution entirely.
643        // Just run hooks to emit tool.call_requested.
644        if server_tool_calls.is_empty() {
645            let mut result = ActResult {
646                results: vec![],
647                completed: true,
648                success_count: 0,
649                error_count: 0,
650                waiting_for_tool_results: false,
651                waiting_for_url_elicitation: false,
652                blocked: false,
653                client_tool_calls,
654                client_tool_definitions,
655            };
656            act_hooks::run_post_act_hooks(
657                &self.hooks,
658                &context,
659                &mut result,
660                &tool_definitions,
661                &self.event_emitter,
662                locale.as_deref(),
663            )
664            .await;
665            return Ok(result);
666        }
667
668        // Replace tool_calls with only server-side tools for execution
669        let tool_calls = server_tool_calls;
670
671        tracing::info!(
672            session_id = %context.session_id,
673            turn_id = %context.turn_id,
674            exec_id = %context.exec_id,
675            tool_count = %tool_calls.len(),
676            "ActAtom: executing tools in parallel"
677        );
678
679        // Generate OTel-style span IDs for hierarchical tracing
680        // trace_id: groups all events in this turn
681        // span_id: unique identifier for this act span (shared by started/completed)
682        // parent_span_id: links to turn as parent
683        //
684        // NOTE: TurnId::to_string() returns prefixed format (e.g., "turn_abc123")
685        // matching the format used by turn.started/completed events in Braintrust.
686        let trace_id = context.turn_id.to_string();
687        let act_span_id = Uuid::now_v7().to_string();
688        let parent_span_id = trace_id.clone(); // Parent is the turn
689
690        // Create event context from atom context with span info
691        let event_context = EventContext::from_execution_context(&context).with_span(
692            trace_id.clone(),
693            act_span_id.clone(),
694            Some(parent_span_id.clone()),
695        );
696
697        // Track act phase timing for Braintrust observability
698        let act_start = Instant::now();
699
700        let visible_tool_names = Arc::new(
701            tool_definitions
702                .iter()
703                .map(|def| def.name().to_string())
704                .collect::<HashSet<_>>(),
705        );
706
707        // Build tool name to definition map
708        let tool_map: std::collections::HashMap<&str, &ToolDefinition> = tool_definitions
709            .iter()
710            .map(|def| {
711                let name = def.name();
712                (name, def)
713            })
714            .collect();
715
716        let mut started_data = ActStartedData::with_definitions_and_locale(
717            &tool_calls,
718            &tool_definitions,
719            locale.as_deref(),
720        );
721        for summary in &mut started_data.tool_calls {
722            if let Some(tool_call) = tool_calls.iter().find(|tc| tc.id == summary.id) {
723                let tool_def = tool_map.get(tool_call.name.as_str()).copied();
724                summary.narration = Some(self.render_tool_narration(
725                    &context,
726                    tool_def,
727                    tool_call,
728                    ToolNarrationPhase::Started,
729                    locale.as_deref(),
730                ));
731                summary.completed_narration = Some(self.render_tool_narration(
732                    &context,
733                    tool_def,
734                    tool_call,
735                    ToolNarrationPhase::Completed,
736                    locale.as_deref(),
737                ));
738            }
739        }
740        started_data.headline = self.render_group_headline(
741            &context,
742            &tool_calls,
743            &tool_map,
744            ToolNarrationPhase::Started,
745            locale.as_deref(),
746        );
747
748        // Emit act.started event (with display names from tool definitions)
749        if let Err(e) = self
750            .event_emitter
751            .emit(EventRequest::new(
752                context.session_id,
753                event_context.clone(),
754                started_data,
755            ))
756            .await
757        {
758            tracing::warn!(
759                session_id = %context.session_id,
760                error = %e,
761                "ActAtom: failed to emit act.started event"
762            );
763        }
764
765        // Decide the execution schedule from per-tool metadata. Calls that
766        // share a concurrency class (mutations to the same shared resource) run
767        // sequentially in arrival order; everything else runs concurrently,
768        // bounded by a global cap. `parallel_tool_calls == Some(false)` forces a
769        // fully sequential schedule. Each tool event references the act span as
770        // its parent regardless of scheduling.
771        let classes: Vec<Option<String>> = tool_calls
772            .iter()
773            .map(|tool_call| {
774                tool_map
775                    .get(tool_call.name.as_str())
776                    .and_then(|def| def.concurrency_class())
777                    .map(|class| class.to_string())
778            })
779            .collect();
780        let schedule_config = tool_scheduler::ScheduleConfig {
781            serialize_all: parallel_tool_calls == Some(false),
782            ..tool_scheduler::ScheduleConfig::default()
783        };
784        let results =
785            tool_scheduler::schedule(tool_calls.len(), &classes, schedule_config, |index| {
786                let tool_call = &tool_calls[index];
787                let tool_def = tool_map.get(tool_call.name.as_str()).cloned();
788                self.execute_single_tool(
789                    &context,
790                    tool_call.clone(),
791                    tool_def,
792                    &trace_id,
793                    &act_span_id,
794                    locale.as_deref(),
795                    network_access.as_ref(),
796                    visible_tool_names.clone(),
797                )
798            })
799            .await;
800
801        // Count successes and errors
802        let success_count = results.iter().filter(|r| r.success).count() as u32;
803        let error_count = results.iter().filter(|r| !r.success).count() as u32;
804
805        // Calculate act phase duration
806        let act_duration_ms = act_start.elapsed().as_millis() as u64;
807
808        // Emit act.completed event (same span as act.started, parent is turn)
809        let completed_context = EventContext::from_execution_context(&context).with_span(
810            trace_id.clone(),
811            act_span_id.clone(), // Same span_id as started
812            Some(parent_span_id.clone()),
813        );
814        let mut completed_headline = self.render_group_headline(
815            &context,
816            &tool_calls,
817            &tool_map,
818            ToolNarrationPhase::Completed,
819            locale.as_deref(),
820        );
821        if error_count > 0 {
822            let suffix = crate::localization::format_error_suffix(locale.as_deref(), error_count);
823            completed_headline = Some(match completed_headline {
824                Some(text) => format!("{text}{suffix}"),
825                None => {
826                    crate::localization::format_completed_tool_batch(locale.as_deref(), error_count)
827                }
828            });
829        }
830
831        if let Err(e) = self
832            .event_emitter
833            .emit(EventRequest::new(
834                context.session_id,
835                completed_context,
836                ActCompletedData {
837                    completed: true,
838                    success_count,
839                    error_count,
840                    duration_ms: Some(act_duration_ms),
841                    headline: completed_headline,
842                },
843            ))
844            .await
845        {
846            tracing::warn!(
847                session_id = %context.session_id,
848                error = %e,
849                "ActAtom: failed to emit act.completed event"
850            );
851        }
852
853        tracing::info!(
854            session_id = %context.session_id,
855            turn_id = %context.turn_id,
856            success_count = %success_count,
857            error_count = %error_count,
858            "ActAtom: all tools completed"
859        );
860
861        // Fail the durable workflow fast on any determinism violation (EVE-530).
862        // All tool.completed events have already been emitted above for affected calls.
863        if let Some(fatal_msg) = results.iter().find_map(|r| r.determinism_fatal.as_deref()) {
864            return Err(crate::error::AgentLoopError::tool(format!(
865                "act activity aborted due to determinism violation: {fatal_msg}"
866            )));
867        }
868
869        let mut act_result = ActResult {
870            results,
871            completed: true,
872            success_count,
873            error_count,
874            waiting_for_tool_results: false,
875            waiting_for_url_elicitation: false,
876            blocked: false,
877            client_tool_calls,
878            client_tool_definitions,
879        };
880
881        // Run post-act hooks (connection setup, client-side tool emission, etc.)
882        act_hooks::run_post_act_hooks(
883            &self.hooks,
884            &context,
885            &mut act_result,
886            &tool_definitions,
887            &self.event_emitter,
888            locale.as_deref(),
889        )
890        .await;
891
892        Ok(act_result)
893    }
894}
895
896impl<T, E> ActAtom<T, E>
897where
898    T: ToolExecutor + Send + Sync + 'static,
899    E: EventEmitter + Send + Sync + 'static,
900{
901    fn render_tool_narration(
902        &self,
903        execution_context: &ExecutionContext,
904        tool_def: Option<&ToolDefinition>,
905        tool_call: &ToolCall,
906        phase: ToolNarrationPhase,
907        locale: Option<&str>,
908    ) -> String {
909        let wrapped_store = self.wrap_file_store_for_narration(execution_context);
910        let ctx = ToolNarrationContext::new(wrapped_store.as_deref());
911        for hook in &self.tool_call_hooks {
912            if let Some(narration) = hook.narration(tool_def, tool_call, phase, locale, ctx) {
913                return narration;
914            }
915        }
916        // Capability hooks only reach tools a capability lists in `tools()`.
917        // Tools assembled outside any capability — the unified `spawn_agent`
918        // dispatcher built from delegation targets, registry-augmented and
919        // proxied tools — still own narration, so ask the executing tool
920        // directly before falling back to the generic display-name phrasing.
921        if let Some(narration) = self
922            .context_services
923            .tool_registry
924            .as_ref()
925            .and_then(|registry| registry.get(&tool_call.name))
926            .and_then(|tool| tool.narrate(tool_call, phase, locale, ctx))
927        {
928            return narration;
929        }
930        render_tool_narration_with_locale(tool_def, tool_call, phase, locale)
931    }
932
933    fn render_group_headline(
934        &self,
935        execution_context: &ExecutionContext,
936        tool_calls: &[ToolCall],
937        tool_map: &std::collections::HashMap<&str, &ToolDefinition>,
938        phase: ToolNarrationPhase,
939        locale: Option<&str>,
940    ) -> Option<String> {
941        if tool_calls.is_empty() {
942            return None;
943        }
944        if let [tool_call] = tool_calls {
945            return Some(self.render_tool_narration(
946                execution_context,
947                tool_map.get(tool_call.name.as_str()).copied(),
948                tool_call,
949                phase,
950                locale,
951            ));
952        }
953
954        let actions = tool_calls
955            .iter()
956            .map(|tool_call| {
957                let tool_def = tool_map.get(tool_call.name.as_str()).copied();
958                let narration = self.render_tool_narration(
959                    execution_context,
960                    tool_def,
961                    tool_call,
962                    phase,
963                    locale,
964                );
965                let repeated_narration = self.render_tool_narration(
966                    execution_context,
967                    tool_def,
968                    &tool_call_for_group_summary(tool_call),
969                    phase,
970                    locale,
971                );
972                GroupHeadlineAction::new(tool_call, narration, repeated_narration)
973            })
974            .collect::<Vec<_>>();
975
976        Some(summarize_group_actions(&actions, locale))
977    }
978
979    /// Mirror the file-store wrapping applied during tool execution so
980    /// path-bearing narration uses the same mount resolver and workspace key.
981    fn wrap_file_store_for_narration(
982        &self,
983        execution_context: &ExecutionContext,
984    ) -> Option<Arc<dyn SessionFileSystem>> {
985        let store = self.context_services.file_store.as_ref()?.clone();
986        let store = if let Some(workspace_id) = execution_context.workspace_id {
987            crate::session_files::WorkspaceScopedFileSystem::wrap(store, workspace_id)
988        } else {
989            store
990        };
991        Some(crate::mount_fs::MountFs::wrap_if_needed(store))
992    }
993
994    fn transform_tool_call_for_execution(&self, tool_call: ToolCall) -> ToolCall {
995        self.tool_call_hooks
996            .iter()
997            .fold(tool_call, |tool_call, hook| {
998                hook.transform_for_execution(tool_call)
999            })
1000    }
1001
1002    /// Execute a single tool call
1003    ///
1004    /// Note: OTel instrumentation is handled via event listeners.
1005    /// tool.started/completed events are emitted, and OtelEventListener
1006    /// creates gen-ai spans from those events.
1007    #[allow(clippy::too_many_arguments)]
1008    async fn execute_single_tool(
1009        &self,
1010        context: &ExecutionContext,
1011        tool_call: ToolCall,
1012        tool_def: Option<&ToolDefinition>,
1013        trace_id: &str,
1014        act_span_id: &str,
1015        locale: Option<&str>,
1016        network_access: Option<&crate::network_access::NetworkAccessList>,
1017        visible_tool_names: Arc<HashSet<String>>,
1018    ) -> ToolCallResult {
1019        tracing::debug!(
1020            session_id = %context.session_id,
1021            turn_id = %context.turn_id,
1022            tool_name = %tool_call.name,
1023            tool_call_id = %tool_call.id,
1024            "ActAtom: executing tool"
1025        );
1026
1027        // Generate a unique span_id for this tool call (child of act span)
1028        let tool_span_id = Uuid::now_v7().to_string();
1029
1030        // Create event context from atom context (with act span as parent)
1031        let event_context = EventContext::from_execution_context(context).with_span(
1032            trace_id.to_string(),
1033            tool_span_id.clone(),
1034            Some(act_span_id.to_string()),
1035        );
1036
1037        // Track tool call timing for Braintrust observability
1038        let tool_start = Instant::now();
1039        let tool_call_fingerprint = tool_call_fingerprint(&tool_call);
1040
1041        // Resolve display name from tool definition
1042        let display_name = crate::localization::localized_tool_display_name(
1043            &tool_call.name,
1044            tool_def.and_then(|d| d.display_name()),
1045            locale,
1046        );
1047        let capability_attribution = tool_def.and_then(|def| {
1048            def.capability_attribution()
1049                .map(|(id, name)| (id.to_string(), name.map(str::to_string)))
1050        });
1051
1052        // THREAT[TM-TOOL-009]: enforce the injected per-org outbound tool-call limit.
1053        // Checked before tool.started so a denied call emits no events and leaves
1054        // no unmatched started/completed pair in UI or telemetry.
1055        if let (Some(limiter), Some(ref org_id)) = (
1056            &self.outbound_tool_rate_limiter,
1057            self.context_services.org_id,
1058        ) && !limiter.check_org(org_id).await
1059        {
1060            tracing::warn!(
1061                session_id = %context.session_id,
1062                tool_name = %tool_call.name,
1063                "ActAtom: outbound tool rate limit exceeded for org"
1064            );
1065            return ToolCallResult {
1066                tool_call: tool_call.clone(),
1067                result: ToolResult {
1068                    tool_call_id: tool_call.id.clone(),
1069                    result: None,
1070                    images: None,
1071                    error: Some(
1072                        "Outbound tool rate limit exceeded for this organization; back off and retry later.".to_string(),
1073                    ),
1074                    connection_required: None,
1075                    raw_output: None,
1076                },
1077                success: false,
1078                status: "error".to_string(),
1079                connection_required: None,
1080                determinism_fatal: None,
1081            };
1082        }
1083
1084        // Per-tool-call idempotency (EVE-530): claim before dispatch, replay if
1085        // already settled, refuse AtMostOnce re-execution on stale running claims.
1086        let claim_token = if let Some(ref store) = self.durable_tool_result_store {
1087            let turn_id = context.turn_id.to_string();
1088            match store
1089                .try_claim_tool_call(
1090                    &turn_id,
1091                    &tool_call.id,
1092                    &tool_call.name,
1093                    &tool_call_fingerprint,
1094                )
1095                .await
1096            {
1097                Ok(ToolCallClaimResult::Claimed { claim_token }) => Some(claim_token),
1098
1099                Ok(ToolCallClaimResult::AlreadySettled {
1100                    result_json,
1101                    args_fingerprint: stored_fp,
1102                }) => {
1103                    // Determinism guard: stored args fingerprint must match current call.
1104                    if stored_fp != tool_call_fingerprint {
1105                        let err_msg = format!(
1106                            "determinism violation: tool '{}' replay args fingerprint \
1107                             does not match prior execution (stored={stored_fp}, \
1108                             current={})",
1109                            tool_call.name, tool_call_fingerprint
1110                        );
1111                        tracing::error!(
1112                            session_id = %context.session_id,
1113                            turn_id = %context.turn_id,
1114                            tool_call_id = %tool_call.id,
1115                            stored_fp = %stored_fp,
1116                            current_fp = %tool_call_fingerprint,
1117                            "ActAtom: determinism violation — replay args fingerprint mismatch"
1118                        );
1119                        let result_fp =
1120                            tool_result_fingerprint(&tool_call.name, &ToolResult::error(&err_msg));
1121                        let _ = self
1122                            .event_emitter
1123                            .emit(EventRequest::new(
1124                                context.session_id,
1125                                event_context,
1126                                ToolCompletedData::failure(
1127                                    tool_call.id.clone(),
1128                                    tool_call.name.clone(),
1129                                    "error".to_string(),
1130                                    err_msg.clone(),
1131                                    None,
1132                                )
1133                                .with_fingerprints(tool_call_fingerprint.clone(), result_fp)
1134                                .with_display_name(display_name.clone()),
1135                            ))
1136                            .await;
1137                        return ToolCallResult {
1138                            tool_call: tool_call.clone(),
1139                            result: ToolResult {
1140                                tool_call_id: tool_call.id.clone(),
1141                                result: None,
1142                                images: None,
1143                                error: Some(err_msg.clone()),
1144                                connection_required: None,
1145                                raw_output: None,
1146                            },
1147                            success: false,
1148                            status: "error".to_string(),
1149                            connection_required: None,
1150                            determinism_fatal: Some(err_msg),
1151                        };
1152                    }
1153                    tracing::debug!(
1154                        session_id = %context.session_id,
1155                        turn_id = %context.turn_id,
1156                        tool_call_id = %tool_call.id,
1157                        "ActAtom: replaying already-settled tool call"
1158                    );
1159                    // Emit a replayed tool.completed without re-emitting tool.started.
1160                    let replayed_result: ToolResult = serde_json::from_value(result_json.clone())
1161                        .unwrap_or(ToolResult {
1162                            tool_call_id: tool_call.id.clone(),
1163                            result: Some(result_json),
1164                            images: None,
1165                            error: None,
1166                            connection_required: None,
1167                            raw_output: None,
1168                        });
1169                    let success = replayed_result.error.is_none();
1170                    let status = if success { "success" } else { "error" };
1171                    let result_fp = tool_result_fingerprint(&tool_call.name, &replayed_result);
1172                    let completed_data = if success {
1173                        // Reconstruct content: text + images (preserves image-producing tools on replay)
1174                        let mut content = replayed_result
1175                            .result
1176                            .as_ref()
1177                            .map(|r| vec![ContentPart::tool_result_text(r)])
1178                            .unwrap_or_default();
1179                        if let Some(ref images) = replayed_result.images {
1180                            for img in images {
1181                                content.push(ContentPart::Image(
1182                                    crate::message::ImageContentPart::from_base64(
1183                                        &img.base64,
1184                                        &img.media_type,
1185                                    ),
1186                                ));
1187                            }
1188                        }
1189                        ToolCompletedData::success(
1190                            tool_call.id.clone(),
1191                            tool_call.name.clone(),
1192                            content,
1193                            None,
1194                        )
1195                        .with_fingerprints(tool_call_fingerprint.clone(), result_fp)
1196                        .with_display_name(display_name.clone())
1197                    } else {
1198                        ToolCompletedData::failure(
1199                            tool_call.id.clone(),
1200                            tool_call.name.clone(),
1201                            status.to_string(),
1202                            replayed_result.error.clone().unwrap_or_default(),
1203                            None,
1204                        )
1205                        .with_fingerprints(tool_call_fingerprint.clone(), result_fp)
1206                        .with_display_name(display_name.clone())
1207                    };
1208                    let _ = self
1209                        .event_emitter
1210                        .emit(EventRequest::new(
1211                            context.session_id,
1212                            event_context,
1213                            completed_data,
1214                        ))
1215                        .await;
1216                    let conn_req = replayed_result.connection_required.clone();
1217                    return ToolCallResult {
1218                        tool_call,
1219                        result: replayed_result,
1220                        success,
1221                        status: status.to_string(),
1222                        connection_required: conn_req,
1223                        determinism_fatal: None,
1224                    };
1225                }
1226
1227                Ok(ToolCallClaimResult::AlreadyRunning {
1228                    args_fingerprint: stored_fp,
1229                }) => {
1230                    // Determinism guard: even in the running state, a fingerprint mismatch
1231                    // means the workflow is replaying with different args — fail loudly.
1232                    if stored_fp != tool_call_fingerprint {
1233                        let err_msg = format!(
1234                            "determinism violation: tool '{}' args fingerprint changed \
1235                             while prior claim is still running (stored={stored_fp}, \
1236                             current={tool_call_fingerprint})",
1237                            tool_call.name
1238                        );
1239                        tracing::error!(
1240                            session_id = %context.session_id,
1241                            turn_id = %context.turn_id,
1242                            tool_call_id = %tool_call.id,
1243                            stored = %stored_fp,
1244                            current = %tool_call_fingerprint,
1245                            "ActAtom: determinism violation — running claim fingerprint mismatch"
1246                        );
1247                        let result_fp =
1248                            tool_result_fingerprint(&tool_call.name, &ToolResult::error(&err_msg));
1249                        let _ = self
1250                            .event_emitter
1251                            .emit(EventRequest::new(
1252                                context.session_id,
1253                                event_context,
1254                                ToolCompletedData::failure(
1255                                    tool_call.id.clone(),
1256                                    tool_call.name.clone(),
1257                                    "error".to_string(),
1258                                    err_msg.clone(),
1259                                    None,
1260                                )
1261                                .with_fingerprints(tool_call_fingerprint.clone(), result_fp)
1262                                .with_display_name(display_name.clone()),
1263                            ))
1264                            .await;
1265                        return ToolCallResult {
1266                            tool_call: tool_call.clone(),
1267                            result: ToolResult {
1268                                tool_call_id: tool_call.id.clone(),
1269                                result: None,
1270                                images: None,
1271                                error: Some(err_msg.clone()),
1272                                connection_required: None,
1273                                raw_output: None,
1274                            },
1275                            success: false,
1276                            status: "error".to_string(),
1277                            connection_required: None,
1278                            determinism_fatal: Some(err_msg),
1279                        };
1280                    }
1281
1282                    let sec = tool_def
1283                        .map(|d| d.side_effect_class())
1284                        .unwrap_or(SideEffectClass::AtMostOnce);
1285                    match sec {
1286                        SideEffectClass::Pure | SideEffectClass::Idempotent => {
1287                            // Safe to re-execute; proceed as normal (no claim token).
1288                            tracing::debug!(
1289                                session_id = %context.session_id,
1290                                tool_call_id = %tool_call.id,
1291                                "ActAtom: stale running claim for idempotent tool, re-executing"
1292                            );
1293                            None
1294                        }
1295                        SideEffectClass::AtMostOnce => {
1296                            tracing::warn!(
1297                                session_id = %context.session_id,
1298                                turn_id = %context.turn_id,
1299                                tool_call_id = %tool_call.id,
1300                                "ActAtom: AtMostOnce tool has stale running claim; returning interrupted result"
1301                            );
1302                            // Settle the stale claim as interrupted, then return an error.
1303                            let _ = store
1304                                .settle_tool_call(
1305                                    &turn_id,
1306                                    &tool_call.id,
1307                                    serde_json::Value::Null,
1308                                    "interrupted",
1309                                    Uuid::nil(), // sentinel — bypass token check for interrupt
1310                                )
1311                                .await;
1312                            let err_msg = format!(
1313                                "tool '{}' was interrupted mid-execution during a prior \
1314                                 worker failure; result is uncertain and was not re-run \
1315                                 (AtMostOnce safety)",
1316                                tool_call.name
1317                            );
1318                            let result_fp = tool_result_fingerprint(
1319                                &tool_call.name,
1320                                &ToolResult::error(&err_msg),
1321                            );
1322                            let _ = self
1323                                .event_emitter
1324                                .emit(EventRequest::new(
1325                                    context.session_id,
1326                                    event_context,
1327                                    ToolCompletedData::failure(
1328                                        tool_call.id.clone(),
1329                                        tool_call.name.clone(),
1330                                        "interrupted".to_string(),
1331                                        err_msg.clone(),
1332                                        None,
1333                                    )
1334                                    .with_fingerprints(tool_call_fingerprint.clone(), result_fp)
1335                                    .with_display_name(display_name.clone()),
1336                                ))
1337                                .await;
1338                            return ToolCallResult {
1339                                tool_call: tool_call.clone(),
1340                                result: ToolResult {
1341                                    tool_call_id: tool_call.id.clone(),
1342                                    result: None,
1343                                    images: None,
1344                                    error: Some(err_msg),
1345                                    connection_required: None,
1346                                    raw_output: None,
1347                                },
1348                                success: false,
1349                                status: "error".to_string(),
1350                                connection_required: None,
1351                                determinism_fatal: None,
1352                            };
1353                        }
1354                    }
1355                }
1356
1357                Ok(ToolCallClaimResult::DeterminismViolation {
1358                    stored_fingerprint,
1359                    current_fingerprint,
1360                }) => {
1361                    let err_msg = format!(
1362                        "determinism violation: tool '{}' args fingerprint changed \
1363                         on replay (stored={stored_fingerprint}, \
1364                         current={current_fingerprint})",
1365                        tool_call.name
1366                    );
1367                    tracing::error!(
1368                        session_id = %context.session_id,
1369                        turn_id = %context.turn_id,
1370                        tool_call_id = %tool_call.id,
1371                        stored = %stored_fingerprint,
1372                        current = %current_fingerprint,
1373                        "ActAtom: determinism violation on claim"
1374                    );
1375                    let result_fp =
1376                        tool_result_fingerprint(&tool_call.name, &ToolResult::error(&err_msg));
1377                    let _ = self
1378                        .event_emitter
1379                        .emit(EventRequest::new(
1380                            context.session_id,
1381                            event_context,
1382                            ToolCompletedData::failure(
1383                                tool_call.id.clone(),
1384                                tool_call.name.clone(),
1385                                "error".to_string(),
1386                                err_msg.clone(),
1387                                None,
1388                            )
1389                            .with_fingerprints(tool_call_fingerprint.clone(), result_fp)
1390                            .with_display_name(display_name.clone()),
1391                        ))
1392                        .await;
1393                    return ToolCallResult {
1394                        tool_call: tool_call.clone(),
1395                        result: ToolResult {
1396                            tool_call_id: tool_call.id.clone(),
1397                            result: None,
1398                            images: None,
1399                            error: Some(err_msg.clone()),
1400                            connection_required: None,
1401                            raw_output: None,
1402                        },
1403                        success: false,
1404                        status: "error".to_string(),
1405                        connection_required: None,
1406                        determinism_fatal: Some(err_msg),
1407                    };
1408                }
1409
1410                Err(e) => {
1411                    tracing::warn!(
1412                        session_id = %context.session_id,
1413                        tool_call_id = %tool_call.id,
1414                        error = %e,
1415                        "ActAtom: durable claim failed; proceeding without idempotency"
1416                    );
1417                    None
1418                }
1419            }
1420        } else {
1421            None
1422        };
1423
1424        // Emit tool.started event (child of act.started)
1425        if let Err(e) = self
1426            .event_emitter
1427            .emit(EventRequest::new(
1428                context.session_id,
1429                event_context.clone(),
1430                ToolStartedData {
1431                    tool_call: tool_call.clone(),
1432                    tool_call_fingerprint: Some(tool_call_fingerprint.clone()),
1433                    display_name: display_name.clone(),
1434                    narration: Some(self.render_tool_narration(
1435                        context,
1436                        tool_def,
1437                        &tool_call,
1438                        ToolNarrationPhase::Started,
1439                        locale,
1440                    )),
1441                },
1442            ))
1443            .await
1444        {
1445            tracing::warn!(
1446                session_id = %context.session_id,
1447                tool_call_id = %tool_call.id,
1448                error = %e,
1449                "ActAtom: failed to emit tool.started event"
1450            );
1451        }
1452
1453        // If tool definition not found, return error result
1454        let Some(tool_def) = tool_def else {
1455            let error_msg = format!("Tool definition not found: {}", tool_call.name);
1456            let tool_duration_ms = tool_start.elapsed().as_millis() as u64;
1457
1458            // Emit tool.completed event for error (child of act.started)
1459            if let Err(e) = self
1460                .event_emitter
1461                .emit(EventRequest::new(
1462                    context.session_id,
1463                    event_context,
1464                    ToolCompletedData::failure(
1465                        tool_call.id.clone(),
1466                        tool_call.name.clone(),
1467                        "error".to_string(),
1468                        error_msg.clone(),
1469                        Some(tool_duration_ms),
1470                    )
1471                    .with_fingerprints(
1472                        tool_call_fingerprint.clone(),
1473                        tool_error_fingerprint(&tool_call.name, "error", &error_msg),
1474                    )
1475                    .with_narration(Some(self.render_tool_narration(
1476                        context,
1477                        None,
1478                        &tool_call,
1479                        ToolNarrationPhase::Failed,
1480                        locale,
1481                    ))),
1482                ))
1483                .await
1484            {
1485                tracing::warn!(
1486                    session_id = %context.session_id,
1487                    tool_call_id = %tool_call.id,
1488                    error = %e,
1489                    "ActAtom: failed to emit tool.completed event"
1490                );
1491            }
1492
1493            return ToolCallResult {
1494                tool_call: tool_call.clone(),
1495                result: ToolResult {
1496                    tool_call_id: tool_call.id.clone(),
1497                    result: None,
1498                    images: None,
1499                    error: Some(error_msg),
1500                    connection_required: None,
1501                    raw_output: None,
1502                },
1503                success: false,
1504                status: "error".to_string(),
1505                connection_required: None,
1506                determinism_fatal: None,
1507            };
1508        };
1509
1510        // Execute the tool (always with context so tools can emit progress events)
1511        let mut tool_context =
1512            ToolContext::from_services(context.session_id, &self.context_services);
1513        if let Some(resolver) = tool_context.connection_resolver.as_ref()
1514            && let Some(bound) = resolver.for_execution(context.input_message_id.uuid())
1515        {
1516            tool_context.connection_resolver = Some(bound);
1517        }
1518
1519        if let Some(scope) =
1520            tool_context.extension::<everruns_core::tool_context::ExecutionServicesExt>()
1521        {
1522            scope
1523                .0
1524                .bind(&mut tool_context, context.input_message_id.uuid());
1525        }
1526        if let Some(invoker) = tool_context.mcp_invoker.as_ref()
1527            && let Some(bound) = invoker.for_execution(context.input_message_id.uuid())
1528        {
1529            tool_context.mcp_invoker = Some(bound);
1530        }
1531        if let Some(authority) = tool_context.session_creation_authority.as_ref()
1532            && let Some(bound) = authority.for_execution(context.input_message_id.uuid())
1533        {
1534            tool_context.session_creation_authority = Some(bound);
1535        }
1536        // Key file I/O by the attached workspace when known: pin the file store
1537        // to the workspace so shared-workspace sessions address the workspace's
1538        // files, not the session's own keyspace. For the default 1:1 case this
1539        // is a transparent pass-through.
1540        if let Some(workspace_id) = context.workspace_id {
1541            tool_context.workspace_id = workspace_id;
1542            if let Some(store) = tool_context.file_store.take() {
1543                tool_context.file_store = Some(
1544                    crate::session_files::WorkspaceScopedFileSystem::wrap(store, workspace_id),
1545                );
1546            }
1547        }
1548        // Resolve model paths through the mount resolver (EVE-660): `/workspace`
1549        // is a mount + cwd, not a per-store prefix. Applied over the
1550        // workspace-keyed store so resolution sits above re-keying.
1551        if let Some(store) = tool_context.file_store.take() {
1552            tool_context.file_store = Some(crate::mount_fs::MountFs::wrap_if_needed(store));
1553        }
1554        tool_context.visible_tool_names = Some(visible_tool_names.clone());
1555        // Input network_access (per-session, merged from harness+agent+session) takes precedence
1556        tool_context.network_access = network_access
1557            .cloned()
1558            .or_else(|| self.context_services.network_access.clone());
1559        // Provide event emitter + context so tools can emit tool.progress events
1560        if tool_context.event_emitter.is_none() {
1561            tool_context.event_emitter =
1562                Some(Arc::new(self.event_emitter.clone()) as Arc<dyn EventEmitter>);
1563        }
1564        tool_context.bind_to_turn(event_context.clone());
1565        tool_context.tool_call_id = Some(tool_call.id.clone());
1566
1567        // Cooperative cancellation for this call. The guard fires when this
1568        // future is dropped — which is what a cancelled turn looks like from
1569        // here — and also on normal return, so the contract a tool sees is
1570        // simply "this call is over". Work the tool leaves running (a child
1571        // process, a detached watcher) can hold a clone and die with the call
1572        // instead of outliving it; dropping the future alone cannot tell it
1573        // anything, because a dropped future is never polled again.
1574        let call_cancellation = tokio_util::sync::CancellationToken::new();
1575        tool_context.cancellation = Some(call_cancellation.clone());
1576        let _cancel_on_call_end = call_cancellation.drop_guard();
1577
1578        let execution_tool_call = self.transform_tool_call_for_execution(tool_call.clone());
1579
1580        // Run pre-tool-use hooks (capability-contributed). They can mutate the
1581        // call, block it, or defer it; a blocked or deferred call is not
1582        // invoked and its result flows through the ordinary completion path.
1583        let (execution_tool_call, pre_hook_result) = act_hooks::pre_tool_use_outcome(
1584            &self.pre_tool_hooks,
1585            execution_tool_call,
1586            tool_def,
1587            &tool_context,
1588        )
1589        .await;
1590
1591        let result = if let Some(pre_hook_result) = pre_hook_result {
1592            Ok(pre_hook_result)
1593        } else if tool_def.is_cpu_bound() {
1594            // CPU-bound / non-yielding in-process tools (e.g. the bash
1595            // interpreter) get their own task so a long synchronous burst
1596            // cannot starve the cooperative polling of I/O-bound tools running
1597            // alongside them in this act batch. On the multi-thread runtime the
1598            // spawned task can also progress on another worker thread.
1599            let executor = self.tool_executor.clone();
1600            let call = execution_tool_call.clone();
1601            let def = tool_def.clone();
1602            let ctx = tool_context.clone();
1603            match AbortOnDropJoinHandle::new(everruns_contracts::rt::spawn(async move {
1604                executor.execute_with_context(&call, &def, &ctx).await
1605            }))
1606            .await
1607            {
1608                Ok(result) => result,
1609                Err(join_err) => Err(crate::error::AgentLoopError::tool(format!(
1610                    "tool task failed to complete: {join_err}"
1611                ))),
1612            }
1613        } else {
1614            self.tool_executor
1615                .execute_with_context(&execution_tool_call, tool_def, &tool_context)
1616                .await
1617        };
1618
1619        match result {
1620            Ok(mut tool_result) => {
1621                // Run post-tool-exec hooks (capability then final/infrastructure)
1622                act_hooks::run_post_tool_exec_hooks(
1623                    &self.post_tool_hooks,
1624                    &self.final_post_tool_hooks,
1625                    &execution_tool_call,
1626                    tool_def,
1627                    &mut tool_result,
1628                    &tool_context,
1629                )
1630                .await;
1631
1632                let tool_duration_ms = tool_start.elapsed().as_millis() as u64;
1633                let success = tool_result.error.is_none();
1634                let status = if success { "success" } else { "error" };
1635
1636                // Emit tool.completed event
1637                let completed_data = if success {
1638                    let result_fingerprint = tool_result_fingerprint(&tool_call.name, &tool_result);
1639                    // Convert result to ContentPart (text + optional images)
1640                    let mut result_content = tool_result
1641                        .result
1642                        .as_ref()
1643                        .map(|r| vec![ContentPart::tool_result_text(r)])
1644                        .unwrap_or_default();
1645                    // Append images as native Image content parts
1646                    if let Some(ref images) = tool_result.images {
1647                        for img in images {
1648                            result_content.push(ContentPart::Image(
1649                                crate::message::ImageContentPart::from_base64(
1650                                    &img.base64,
1651                                    &img.media_type,
1652                                ),
1653                            ));
1654                        }
1655                    }
1656                    ToolCompletedData::success(
1657                        tool_call.id.clone(),
1658                        tool_call.name.clone(),
1659                        result_content,
1660                        Some(tool_duration_ms),
1661                    )
1662                    .with_fingerprints(tool_call_fingerprint.clone(), result_fingerprint)
1663                    .with_display_name(display_name.clone())
1664                    .with_capability_attribution(
1665                        capability_attribution.as_ref().map(|(id, _)| id.clone()),
1666                        capability_attribution
1667                            .as_ref()
1668                            .and_then(|(_, name)| name.clone()),
1669                    )
1670                    .with_narration(Some(self.render_tool_narration(
1671                        context,
1672                        Some(tool_def),
1673                        &tool_call,
1674                        ToolNarrationPhase::Completed,
1675                        locale,
1676                    )))
1677                } else {
1678                    let result_fingerprint = tool_result_fingerprint(&tool_call.name, &tool_result);
1679                    ToolCompletedData::failure(
1680                        tool_call.id.clone(),
1681                        tool_call.name.clone(),
1682                        status.to_string(),
1683                        tool_result.error.clone().unwrap_or_default(),
1684                        Some(tool_duration_ms),
1685                    )
1686                    .with_fingerprints(tool_call_fingerprint.clone(), result_fingerprint)
1687                    .with_display_name(display_name.clone())
1688                    .with_capability_attribution(
1689                        capability_attribution.as_ref().map(|(id, _)| id.clone()),
1690                        capability_attribution
1691                            .as_ref()
1692                            .and_then(|(_, name)| name.clone()),
1693                    )
1694                    .with_narration(Some(self.render_tool_narration(
1695                        context,
1696                        Some(tool_def),
1697                        &tool_call,
1698                        ToolNarrationPhase::Failed,
1699                        locale,
1700                    )))
1701                };
1702
1703                if let Err(e) = self
1704                    .event_emitter
1705                    .emit(EventRequest::new(
1706                        context.session_id,
1707                        event_context.clone(),
1708                        completed_data,
1709                    ))
1710                    .await
1711                {
1712                    tracing::warn!(
1713                        session_id = %context.session_id,
1714                        tool_call_id = %tool_call.id,
1715                        error = %e,
1716                        "ActAtom: failed to emit tool.completed event"
1717                    );
1718                }
1719
1720                tracing::debug!(
1721                    session_id = %context.session_id,
1722                    tool_name = %tool_call.name,
1723                    tool_call_id = %tool_call.id,
1724                    success = %success,
1725                    "ActAtom: tool execution completed"
1726                );
1727
1728                // Settle the durable claim (EVE-530).
1729                if let (Some(store), Some(token)) = (&self.durable_tool_result_store, claim_token) {
1730                    let result_snapshot =
1731                        serde_json::to_value(&tool_result).unwrap_or(serde_json::Value::Null);
1732                    match store
1733                        .settle_tool_call(
1734                            &context.turn_id.to_string(),
1735                            &tool_call.id,
1736                            result_snapshot,
1737                            "settled",
1738                            token,
1739                        )
1740                        .await
1741                    {
1742                        Ok(false) => {
1743                            tracing::warn!(
1744                                session_id = %context.session_id,
1745                                tool_call_id = %tool_call.id,
1746                                "ActAtom: settle ownership check failed (task reclaimed)"
1747                            );
1748                        }
1749                        Err(e) => {
1750                            tracing::warn!(
1751                                session_id = %context.session_id,
1752                                tool_call_id = %tool_call.id,
1753                                error = %e,
1754                                "ActAtom: settle_tool_call failed"
1755                            );
1756                        }
1757                        Ok(true) => {}
1758                    }
1759                }
1760
1761                let conn_req = tool_result.connection_required.clone();
1762                ToolCallResult {
1763                    tool_call,
1764                    result: tool_result,
1765                    success,
1766                    status: status.to_string(),
1767                    connection_required: conn_req,
1768                    determinism_fatal: None,
1769                }
1770            }
1771            Err(e) => {
1772                let tool_duration_ms = tool_start.elapsed().as_millis() as u64;
1773                let error_msg = e.to_string();
1774
1775                // Emit tool.completed event for error
1776                if let Err(emit_err) = self
1777                    .event_emitter
1778                    .emit(EventRequest::new(
1779                        context.session_id,
1780                        event_context,
1781                        ToolCompletedData::failure(
1782                            tool_call.id.clone(),
1783                            tool_call.name.clone(),
1784                            "error".to_string(),
1785                            error_msg.clone(),
1786                            Some(tool_duration_ms),
1787                        )
1788                        .with_fingerprints(
1789                            tool_call_fingerprint.clone(),
1790                            tool_error_fingerprint(&tool_call.name, "error", &error_msg),
1791                        )
1792                        .with_display_name(display_name.clone())
1793                        .with_capability_attribution(
1794                            capability_attribution.as_ref().map(|(id, _)| id.clone()),
1795                            capability_attribution
1796                                .as_ref()
1797                                .and_then(|(_, name)| name.clone()),
1798                        )
1799                        .with_narration(Some(self.render_tool_narration(
1800                            context,
1801                            Some(tool_def),
1802                            &tool_call,
1803                            ToolNarrationPhase::Failed,
1804                            locale,
1805                        ))),
1806                    ))
1807                    .await
1808                {
1809                    tracing::warn!(
1810                        session_id = %context.session_id,
1811                        tool_call_id = %tool_call.id,
1812                        error = %emit_err,
1813                        "ActAtom: failed to emit tool.completed event"
1814                    );
1815                }
1816
1817                tracing::warn!(
1818                    session_id = %context.session_id,
1819                    tool_name = %tool_call.name,
1820                    tool_call_id = %tool_call.id,
1821                    error = %e,
1822                    "ActAtom: tool execution failed"
1823                );
1824
1825                ToolCallResult {
1826                    tool_call: tool_call.clone(),
1827                    result: ToolResult {
1828                        tool_call_id: tool_call.id.clone(),
1829                        result: None,
1830                        images: None,
1831                        error: Some(error_msg),
1832                        connection_required: None,
1833                        raw_output: None,
1834                    },
1835                    success: false,
1836                    status: "error".to_string(),
1837                    connection_required: None,
1838                    determinism_fatal: None,
1839                }
1840            }
1841        }
1842    }
1843}
1844
1845// ============================================================================
1846// Tests
1847// ============================================================================
1848
1849#[cfg(test)]
1850#[path = "act_tests.rs"]
1851mod tests;
1852
1853#[cfg(test)]
1854#[path = "act_approval_tests.rs"]
1855mod approval_tests;