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 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
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: 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// ============================================================================
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 appends synthetic `ask_user` calls, so it must
320    /// 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::ClientSideToolHook),
327        ]
328    }
329
330    /// Default final post-tool-exec hooks (infrastructure, always-on).
331    /// These run after all capability-contributed hooks and cannot be removed.
332    fn default_final_hooks() -> Vec<Arc<dyn act_hooks::PostToolExecHook>> {
333        vec![Arc::new(act_hooks::OutputHardLimitHook)]
334    }
335
336    /// Set the session storage store on this atom
337    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    /// Set the image artifact store on this atom
346    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    /// Set the provider credential store on this atom
355    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    /// Set the utility LLM service on this atom.
364    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    /// Set the scoped-MCP tool invoker on this atom (guardrails `mcp` check).
370    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    /// Set the outbound egress service on this atom.
376    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    /// Set the user connection resolver on this atom
382    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    /// Set session store for context-aware tools.
391    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    /// Set agent store for context-aware tools.
397    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    /// Set session schedule store for scheduling tools.
403    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    /// Set platform store for org-level management tools.
412    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    /// Set leased resource store for lifecycle-managed provider resources.
421    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    /// Set session resource registry.
430    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    /// Add a session task registry passed to tool contexts.
439    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    /// Set the active built-in tool registry for meta-tools like `spawn_background`.
456    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    /// Add capability-contributed post-tool-exec hooks.
462    /// Callers should pass hooks from the *active* capabilities for this session,
463    /// not from the full platform registry.
464    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    /// Add capability-contributed pre-tool-use hooks. Pre-hooks fire before
473    /// each tool call and can mutate or block it; see
474    /// `act_hooks::PreToolUseHook` and `knowledge/runtime-resources/user-hooks.md`.
475    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    /// Set org ID for org-scoped operations.
489    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    /// Set the merged network access list for URL filtering in tools.
495    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    /// Set the budget checker for the check_budget tool.
504    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    /// Set the internal payment authority for paid capability tools.
513    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    /// Set the authority used to authorize detached peer-session creation.
522    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    /// Set the per-org outbound tool-call rate limiter (TM-TOOL-009).
531    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    /// Set the durable per-tool-call idempotency store (EVE-530).
540    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    /// Set the durable subagent spawn handle store (EVE-535).
549    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    /// Set the resolved subagent nesting policy for tool contexts.
558    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    /// Set the live reasoning-effort handle (EVE-595). When set, each tool's
567    /// `ToolContext` receives a clone so a tool can change the reasoning effort
568    /// mid-turn for subsequent LLM steps in the same turn.
569    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    /// Stable phase name used by logs and durable activity adapters.
584    pub fn name(&self) -> &'static str {
585        "act"
586    }
587
588    /// Execute one scheduled tool-call batch through injected contracts.
589    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            .. // agent_id/org_id not needed here, just passed through workflow
598        } = input;
599
600        // Partition tool calls: server-side tools get executed, client-side tools
601        // are stored on ActResult for the ClientSideToolHook to emit.
602        let (server_tool_calls, client_tool_calls): (Vec<_>, Vec<_>) = tool_calls
603            .into_iter()
604            .partition(|tc| act_hooks::runs_on_server(tc, &tool_definitions));
605
606        let client_tool_calls: Vec<_> = client_tool_calls
607            .into_iter()
608            .map(|tool_call| self.transform_tool_call_for_execution(tool_call))
609            .collect();
610
611        let client_tool_definitions: Vec<_> = if client_tool_calls.is_empty() {
612            vec![]
613        } else {
614            tool_definitions
615                .iter()
616                .filter(|td| {
617                    if let ToolDefinition::ClientSide(ct) = td {
618                        client_tool_calls.iter().any(|tc| tc.name == ct.name)
619                    } else {
620                        false
621                    }
622                })
623                .cloned()
624                .collect()
625        };
626
627        if server_tool_calls.is_empty() && client_tool_calls.is_empty() {
628            return Ok(ActResult {
629                results: vec![],
630                completed: true,
631                success_count: 0,
632                error_count: 0,
633                waiting_for_tool_results: false,
634                waiting_for_url_elicitation: false,
635                blocked: false,
636                client_tool_calls: vec![],
637                client_tool_definitions: vec![],
638            });
639        }
640
641        // If only client-side tools (no server-side), skip tool execution entirely.
642        // Just run hooks to emit tool.call_requested.
643        if server_tool_calls.is_empty() {
644            let mut result = ActResult {
645                results: vec![],
646                completed: true,
647                success_count: 0,
648                error_count: 0,
649                waiting_for_tool_results: false,
650                waiting_for_url_elicitation: false,
651                blocked: false,
652                client_tool_calls,
653                client_tool_definitions,
654            };
655            act_hooks::run_post_act_hooks(
656                &self.hooks,
657                &context,
658                &mut result,
659                &tool_definitions,
660                &self.event_emitter,
661                locale.as_deref(),
662            )
663            .await;
664            return Ok(result);
665        }
666
667        // Replace tool_calls with only server-side tools for execution
668        let tool_calls = server_tool_calls;
669
670        tracing::info!(
671            session_id = %context.session_id,
672            turn_id = %context.turn_id,
673            exec_id = %context.exec_id,
674            tool_count = %tool_calls.len(),
675            "ActAtom: executing tools in parallel"
676        );
677
678        // Generate OTel-style span IDs for hierarchical tracing
679        // trace_id: groups all events in this turn
680        // span_id: unique identifier for this act span (shared by started/completed)
681        // parent_span_id: links to turn as parent
682        //
683        // NOTE: TurnId::to_string() returns prefixed format (e.g., "turn_abc123")
684        // matching the format used by turn.started/completed events in Braintrust.
685        let trace_id = context.turn_id.to_string();
686        let act_span_id = Uuid::now_v7().to_string();
687        let parent_span_id = trace_id.clone(); // Parent is the turn
688
689        // Create event context from atom context with span info
690        let event_context = EventContext::from_execution_context(&context).with_span(
691            trace_id.clone(),
692            act_span_id.clone(),
693            Some(parent_span_id.clone()),
694        );
695
696        // Track act phase timing for Braintrust observability
697        let act_start = Instant::now();
698
699        let visible_tool_names = Arc::new(
700            tool_definitions
701                .iter()
702                .map(|def| def.name().to_string())
703                .collect::<HashSet<_>>(),
704        );
705
706        // Build tool name to definition map
707        let tool_map: std::collections::HashMap<&str, &ToolDefinition> = tool_definitions
708            .iter()
709            .map(|def| {
710                let name = def.name();
711                (name, def)
712            })
713            .collect();
714
715        let mut started_data = ActStartedData::with_definitions_and_locale(
716            &tool_calls,
717            &tool_definitions,
718            locale.as_deref(),
719        );
720        for summary in &mut started_data.tool_calls {
721            if let Some(tool_call) = tool_calls.iter().find(|tc| tc.id == summary.id) {
722                let tool_def = tool_map.get(tool_call.name.as_str()).copied();
723                summary.narration = Some(self.render_tool_narration(
724                    &context,
725                    tool_def,
726                    tool_call,
727                    ToolNarrationPhase::Started,
728                    locale.as_deref(),
729                ));
730                summary.completed_narration = Some(self.render_tool_narration(
731                    &context,
732                    tool_def,
733                    tool_call,
734                    ToolNarrationPhase::Completed,
735                    locale.as_deref(),
736                ));
737            }
738        }
739        started_data.headline = self.render_group_headline(
740            &context,
741            &tool_calls,
742            &tool_map,
743            ToolNarrationPhase::Started,
744            locale.as_deref(),
745        );
746
747        // Emit act.started event (with display names from tool definitions)
748        if let Err(e) = self
749            .event_emitter
750            .emit(EventRequest::new(
751                context.session_id,
752                event_context.clone(),
753                started_data,
754            ))
755            .await
756        {
757            tracing::warn!(
758                session_id = %context.session_id,
759                error = %e,
760                "ActAtom: failed to emit act.started event"
761            );
762        }
763
764        // Decide the execution schedule from per-tool metadata. Calls that
765        // share a concurrency class (mutations to the same shared resource) run
766        // sequentially in arrival order; everything else runs concurrently,
767        // bounded by a global cap. `parallel_tool_calls == Some(false)` forces a
768        // fully sequential schedule. Each tool event references the act span as
769        // its parent regardless of scheduling.
770        let classes: Vec<Option<String>> = tool_calls
771            .iter()
772            .map(|tool_call| {
773                tool_map
774                    .get(tool_call.name.as_str())
775                    .and_then(|def| def.concurrency_class())
776                    .map(|class| class.to_string())
777            })
778            .collect();
779        let schedule_config = tool_scheduler::ScheduleConfig {
780            serialize_all: parallel_tool_calls == Some(false),
781            ..tool_scheduler::ScheduleConfig::default()
782        };
783        let results =
784            tool_scheduler::schedule(tool_calls.len(), &classes, schedule_config, |index| {
785                let tool_call = &tool_calls[index];
786                let tool_def = tool_map.get(tool_call.name.as_str()).cloned();
787                self.execute_single_tool(
788                    &context,
789                    tool_call.clone(),
790                    tool_def,
791                    &trace_id,
792                    &act_span_id,
793                    locale.as_deref(),
794                    network_access.as_ref(),
795                    visible_tool_names.clone(),
796                )
797            })
798            .await;
799
800        // Count successes and errors
801        let success_count = results.iter().filter(|r| r.success).count() as u32;
802        let error_count = results.iter().filter(|r| !r.success).count() as u32;
803
804        // Calculate act phase duration
805        let act_duration_ms = act_start.elapsed().as_millis() as u64;
806
807        // Emit act.completed event (same span as act.started, parent is turn)
808        let completed_context = EventContext::from_execution_context(&context).with_span(
809            trace_id.clone(),
810            act_span_id.clone(), // Same span_id as started
811            Some(parent_span_id.clone()),
812        );
813        let mut completed_headline = self.render_group_headline(
814            &context,
815            &tool_calls,
816            &tool_map,
817            ToolNarrationPhase::Completed,
818            locale.as_deref(),
819        );
820        if error_count > 0 {
821            let suffix = crate::localization::format_error_suffix(locale.as_deref(), error_count);
822            completed_headline = Some(match completed_headline {
823                Some(text) => format!("{text}{suffix}"),
824                None => {
825                    crate::localization::format_completed_tool_batch(locale.as_deref(), error_count)
826                }
827            });
828        }
829
830        if let Err(e) = self
831            .event_emitter
832            .emit(EventRequest::new(
833                context.session_id,
834                completed_context,
835                ActCompletedData {
836                    completed: true,
837                    success_count,
838                    error_count,
839                    duration_ms: Some(act_duration_ms),
840                    headline: completed_headline,
841                },
842            ))
843            .await
844        {
845            tracing::warn!(
846                session_id = %context.session_id,
847                error = %e,
848                "ActAtom: failed to emit act.completed event"
849            );
850        }
851
852        tracing::info!(
853            session_id = %context.session_id,
854            turn_id = %context.turn_id,
855            success_count = %success_count,
856            error_count = %error_count,
857            "ActAtom: all tools completed"
858        );
859
860        // Fail the durable workflow fast on any determinism violation (EVE-530).
861        // All tool.completed events have already been emitted above for affected calls.
862        if let Some(fatal_msg) = results.iter().find_map(|r| r.determinism_fatal.as_deref()) {
863            return Err(crate::error::AgentLoopError::tool(format!(
864                "act activity aborted due to determinism violation: {fatal_msg}"
865            )));
866        }
867
868        let mut act_result = ActResult {
869            results,
870            completed: true,
871            success_count,
872            error_count,
873            waiting_for_tool_results: false,
874            waiting_for_url_elicitation: false,
875            blocked: false,
876            client_tool_calls,
877            client_tool_definitions,
878        };
879
880        // Run post-act hooks (connection setup, client-side tool emission, etc.)
881        act_hooks::run_post_act_hooks(
882            &self.hooks,
883            &context,
884            &mut act_result,
885            &tool_definitions,
886            &self.event_emitter,
887            locale.as_deref(),
888        )
889        .await;
890
891        Ok(act_result)
892    }
893}
894
895impl<T, E> ActAtom<T, E>
896where
897    T: ToolExecutor + Send + Sync + 'static,
898    E: EventEmitter + Send + Sync + 'static,
899{
900    fn render_tool_narration(
901        &self,
902        execution_context: &ExecutionContext,
903        tool_def: Option<&ToolDefinition>,
904        tool_call: &ToolCall,
905        phase: ToolNarrationPhase,
906        locale: Option<&str>,
907    ) -> String {
908        let wrapped_store = self.wrap_file_store_for_narration(execution_context);
909        let ctx = ToolNarrationContext::new(wrapped_store.as_deref());
910        for hook in &self.tool_call_hooks {
911            if let Some(narration) = hook.narration(tool_def, tool_call, phase, locale, ctx) {
912                return narration;
913            }
914        }
915        // Capability hooks only reach tools a capability lists in `tools()`.
916        // Tools assembled outside any capability — the unified `spawn_agent`
917        // dispatcher built from delegation targets, registry-augmented and
918        // proxied tools — still own narration, so ask the executing tool
919        // directly before falling back to the generic display-name phrasing.
920        if let Some(narration) = self
921            .context_services
922            .tool_registry
923            .as_ref()
924            .and_then(|registry| registry.get(&tool_call.name))
925            .and_then(|tool| tool.narrate(tool_call, phase, locale, ctx))
926        {
927            return narration;
928        }
929        render_tool_narration_with_locale(tool_def, tool_call, phase, locale)
930    }
931
932    fn render_group_headline(
933        &self,
934        execution_context: &ExecutionContext,
935        tool_calls: &[ToolCall],
936        tool_map: &std::collections::HashMap<&str, &ToolDefinition>,
937        phase: ToolNarrationPhase,
938        locale: Option<&str>,
939    ) -> Option<String> {
940        if tool_calls.is_empty() {
941            return None;
942        }
943        if let [tool_call] = tool_calls {
944            return Some(self.render_tool_narration(
945                execution_context,
946                tool_map.get(tool_call.name.as_str()).copied(),
947                tool_call,
948                phase,
949                locale,
950            ));
951        }
952
953        let actions = tool_calls
954            .iter()
955            .map(|tool_call| {
956                let tool_def = tool_map.get(tool_call.name.as_str()).copied();
957                let narration = self.render_tool_narration(
958                    execution_context,
959                    tool_def,
960                    tool_call,
961                    phase,
962                    locale,
963                );
964                let repeated_narration = self.render_tool_narration(
965                    execution_context,
966                    tool_def,
967                    &tool_call_for_group_summary(tool_call),
968                    phase,
969                    locale,
970                );
971                GroupHeadlineAction::new(tool_call, narration, repeated_narration)
972            })
973            .collect::<Vec<_>>();
974
975        Some(summarize_group_actions(&actions, locale))
976    }
977
978    /// Mirror the file-store wrapping applied during tool execution so
979    /// path-bearing narration uses the same mount resolver and workspace key.
980    fn wrap_file_store_for_narration(
981        &self,
982        execution_context: &ExecutionContext,
983    ) -> Option<Arc<dyn SessionFileSystem>> {
984        let store = self.context_services.file_store.as_ref()?.clone();
985        let store = if let Some(workspace_id) = execution_context.workspace_id {
986            crate::session_files::WorkspaceScopedFileSystem::wrap(store, workspace_id)
987        } else {
988            store
989        };
990        Some(crate::mount_fs::MountFs::wrap_if_needed(store))
991    }
992
993    fn transform_tool_call_for_execution(&self, tool_call: ToolCall) -> ToolCall {
994        self.tool_call_hooks
995            .iter()
996            .fold(tool_call, |tool_call, hook| {
997                hook.transform_for_execution(tool_call)
998            })
999    }
1000
1001    /// Execute a single tool call
1002    ///
1003    /// Note: OTel instrumentation is handled via event listeners.
1004    /// tool.started/completed events are emitted, and OtelEventListener
1005    /// creates gen-ai spans from those events.
1006    #[allow(clippy::too_many_arguments)]
1007    async fn execute_single_tool(
1008        &self,
1009        context: &ExecutionContext,
1010        tool_call: ToolCall,
1011        tool_def: Option<&ToolDefinition>,
1012        trace_id: &str,
1013        act_span_id: &str,
1014        locale: Option<&str>,
1015        network_access: Option<&crate::network_access::NetworkAccessList>,
1016        visible_tool_names: Arc<HashSet<String>>,
1017    ) -> ToolCallResult {
1018        tracing::debug!(
1019            session_id = %context.session_id,
1020            turn_id = %context.turn_id,
1021            tool_name = %tool_call.name,
1022            tool_call_id = %tool_call.id,
1023            "ActAtom: executing tool"
1024        );
1025
1026        // Generate a unique span_id for this tool call (child of act span)
1027        let tool_span_id = Uuid::now_v7().to_string();
1028
1029        // Create event context from atom context (with act span as parent)
1030        let event_context = EventContext::from_execution_context(context).with_span(
1031            trace_id.to_string(),
1032            tool_span_id.clone(),
1033            Some(act_span_id.to_string()),
1034        );
1035
1036        // Track tool call timing for Braintrust observability
1037        let tool_start = Instant::now();
1038        let tool_call_fingerprint = tool_call_fingerprint(&tool_call);
1039
1040        // Resolve display name from tool definition
1041        let display_name = crate::localization::localized_tool_display_name(
1042            &tool_call.name,
1043            tool_def.and_then(|d| d.display_name()),
1044            locale,
1045        );
1046        let capability_attribution = tool_def.and_then(|def| {
1047            def.capability_attribution()
1048                .map(|(id, name)| (id.to_string(), name.map(str::to_string)))
1049        });
1050
1051        // THREAT[TM-TOOL-009]: enforce the injected per-org outbound tool-call limit.
1052        // Checked before tool.started so a denied call emits no events and leaves
1053        // no unmatched started/completed pair in UI or telemetry.
1054        if let (Some(limiter), Some(ref org_id)) = (
1055            &self.outbound_tool_rate_limiter,
1056            self.context_services.org_id,
1057        ) && !limiter.check_org(org_id).await
1058        {
1059            tracing::warn!(
1060                session_id = %context.session_id,
1061                tool_name = %tool_call.name,
1062                "ActAtom: outbound tool rate limit exceeded for org"
1063            );
1064            return ToolCallResult {
1065                tool_call: tool_call.clone(),
1066                result: ToolResult {
1067                    tool_call_id: tool_call.id.clone(),
1068                    result: None,
1069                    images: None,
1070                    error: Some(
1071                        "Outbound tool rate limit exceeded for this organization; back off and retry later.".to_string(),
1072                    ),
1073                    connection_required: None,
1074                    raw_output: None,
1075                },
1076                success: false,
1077                status: "error".to_string(),
1078                connection_required: None,
1079                determinism_fatal: None,
1080            };
1081        }
1082
1083        // Per-tool-call idempotency (EVE-530): claim before dispatch, replay if
1084        // already settled, refuse AtMostOnce re-execution on stale running claims.
1085        let claim_token = if let Some(ref store) = self.durable_tool_result_store {
1086            let turn_id = context.turn_id.to_string();
1087            match store
1088                .try_claim_tool_call(
1089                    &turn_id,
1090                    &tool_call.id,
1091                    &tool_call.name,
1092                    &tool_call_fingerprint,
1093                )
1094                .await
1095            {
1096                Ok(ToolCallClaimResult::Claimed { claim_token }) => Some(claim_token),
1097
1098                Ok(ToolCallClaimResult::AlreadySettled {
1099                    result_json,
1100                    args_fingerprint: stored_fp,
1101                }) => {
1102                    // Determinism guard: stored args fingerprint must match current call.
1103                    if stored_fp != tool_call_fingerprint {
1104                        let err_msg = format!(
1105                            "determinism violation: tool '{}' replay args fingerprint \
1106                             does not match prior execution (stored={stored_fp}, \
1107                             current={})",
1108                            tool_call.name, tool_call_fingerprint
1109                        );
1110                        tracing::error!(
1111                            session_id = %context.session_id,
1112                            turn_id = %context.turn_id,
1113                            tool_call_id = %tool_call.id,
1114                            stored_fp = %stored_fp,
1115                            current_fp = %tool_call_fingerprint,
1116                            "ActAtom: determinism violation — replay args fingerprint mismatch"
1117                        );
1118                        let result_fp =
1119                            tool_result_fingerprint(&tool_call.name, &ToolResult::error(&err_msg));
1120                        let _ = self
1121                            .event_emitter
1122                            .emit(EventRequest::new(
1123                                context.session_id,
1124                                event_context,
1125                                ToolCompletedData::failure(
1126                                    tool_call.id.clone(),
1127                                    tool_call.name.clone(),
1128                                    "error".to_string(),
1129                                    err_msg.clone(),
1130                                    None,
1131                                )
1132                                .with_fingerprints(tool_call_fingerprint.clone(), result_fp)
1133                                .with_display_name(display_name.clone()),
1134                            ))
1135                            .await;
1136                        return ToolCallResult {
1137                            tool_call: tool_call.clone(),
1138                            result: ToolResult {
1139                                tool_call_id: tool_call.id.clone(),
1140                                result: None,
1141                                images: None,
1142                                error: Some(err_msg.clone()),
1143                                connection_required: None,
1144                                raw_output: None,
1145                            },
1146                            success: false,
1147                            status: "error".to_string(),
1148                            connection_required: None,
1149                            determinism_fatal: Some(err_msg),
1150                        };
1151                    }
1152                    tracing::debug!(
1153                        session_id = %context.session_id,
1154                        turn_id = %context.turn_id,
1155                        tool_call_id = %tool_call.id,
1156                        "ActAtom: replaying already-settled tool call"
1157                    );
1158                    // Emit a replayed tool.completed without re-emitting tool.started.
1159                    let replayed_result: ToolResult = serde_json::from_value(result_json.clone())
1160                        .unwrap_or(ToolResult {
1161                            tool_call_id: tool_call.id.clone(),
1162                            result: Some(result_json),
1163                            images: None,
1164                            error: None,
1165                            connection_required: None,
1166                            raw_output: None,
1167                        });
1168                    let success = replayed_result.error.is_none();
1169                    let status = if success { "success" } else { "error" };
1170                    let result_fp = tool_result_fingerprint(&tool_call.name, &replayed_result);
1171                    let completed_data = if success {
1172                        // Reconstruct content: text + images (preserves image-producing tools on replay)
1173                        let mut content = replayed_result
1174                            .result
1175                            .as_ref()
1176                            .map(|r| vec![ContentPart::tool_result_text(r)])
1177                            .unwrap_or_default();
1178                        if let Some(ref images) = replayed_result.images {
1179                            for img in images {
1180                                content.push(ContentPart::Image(
1181                                    crate::message::ImageContentPart::from_base64(
1182                                        &img.base64,
1183                                        &img.media_type,
1184                                    ),
1185                                ));
1186                            }
1187                        }
1188                        ToolCompletedData::success(
1189                            tool_call.id.clone(),
1190                            tool_call.name.clone(),
1191                            content,
1192                            None,
1193                        )
1194                        .with_fingerprints(tool_call_fingerprint.clone(), result_fp)
1195                        .with_display_name(display_name.clone())
1196                    } else {
1197                        ToolCompletedData::failure(
1198                            tool_call.id.clone(),
1199                            tool_call.name.clone(),
1200                            status.to_string(),
1201                            replayed_result.error.clone().unwrap_or_default(),
1202                            None,
1203                        )
1204                        .with_fingerprints(tool_call_fingerprint.clone(), result_fp)
1205                        .with_display_name(display_name.clone())
1206                    };
1207                    let _ = self
1208                        .event_emitter
1209                        .emit(EventRequest::new(
1210                            context.session_id,
1211                            event_context,
1212                            completed_data,
1213                        ))
1214                        .await;
1215                    let conn_req = replayed_result.connection_required.clone();
1216                    return ToolCallResult {
1217                        tool_call,
1218                        result: replayed_result,
1219                        success,
1220                        status: status.to_string(),
1221                        connection_required: conn_req,
1222                        determinism_fatal: None,
1223                    };
1224                }
1225
1226                Ok(ToolCallClaimResult::AlreadyRunning {
1227                    args_fingerprint: stored_fp,
1228                }) => {
1229                    // Determinism guard: even in the running state, a fingerprint mismatch
1230                    // means the workflow is replaying with different args — fail loudly.
1231                    if stored_fp != tool_call_fingerprint {
1232                        let err_msg = format!(
1233                            "determinism violation: tool '{}' args fingerprint changed \
1234                             while prior claim is still running (stored={stored_fp}, \
1235                             current={tool_call_fingerprint})",
1236                            tool_call.name
1237                        );
1238                        tracing::error!(
1239                            session_id = %context.session_id,
1240                            turn_id = %context.turn_id,
1241                            tool_call_id = %tool_call.id,
1242                            stored = %stored_fp,
1243                            current = %tool_call_fingerprint,
1244                            "ActAtom: determinism violation — running claim fingerprint mismatch"
1245                        );
1246                        let result_fp =
1247                            tool_result_fingerprint(&tool_call.name, &ToolResult::error(&err_msg));
1248                        let _ = self
1249                            .event_emitter
1250                            .emit(EventRequest::new(
1251                                context.session_id,
1252                                event_context,
1253                                ToolCompletedData::failure(
1254                                    tool_call.id.clone(),
1255                                    tool_call.name.clone(),
1256                                    "error".to_string(),
1257                                    err_msg.clone(),
1258                                    None,
1259                                )
1260                                .with_fingerprints(tool_call_fingerprint.clone(), result_fp)
1261                                .with_display_name(display_name.clone()),
1262                            ))
1263                            .await;
1264                        return ToolCallResult {
1265                            tool_call: tool_call.clone(),
1266                            result: ToolResult {
1267                                tool_call_id: tool_call.id.clone(),
1268                                result: None,
1269                                images: None,
1270                                error: Some(err_msg.clone()),
1271                                connection_required: None,
1272                                raw_output: None,
1273                            },
1274                            success: false,
1275                            status: "error".to_string(),
1276                            connection_required: None,
1277                            determinism_fatal: Some(err_msg),
1278                        };
1279                    }
1280
1281                    let sec = tool_def
1282                        .map(|d| d.side_effect_class())
1283                        .unwrap_or(SideEffectClass::AtMostOnce);
1284                    match sec {
1285                        SideEffectClass::Pure | SideEffectClass::Idempotent => {
1286                            // Safe to re-execute; proceed as normal (no claim token).
1287                            tracing::debug!(
1288                                session_id = %context.session_id,
1289                                tool_call_id = %tool_call.id,
1290                                "ActAtom: stale running claim for idempotent tool, re-executing"
1291                            );
1292                            None
1293                        }
1294                        SideEffectClass::AtMostOnce => {
1295                            tracing::warn!(
1296                                session_id = %context.session_id,
1297                                turn_id = %context.turn_id,
1298                                tool_call_id = %tool_call.id,
1299                                "ActAtom: AtMostOnce tool has stale running claim; returning interrupted result"
1300                            );
1301                            // Settle the stale claim as interrupted, then return an error.
1302                            let _ = store
1303                                .settle_tool_call(
1304                                    &turn_id,
1305                                    &tool_call.id,
1306                                    serde_json::Value::Null,
1307                                    "interrupted",
1308                                    Uuid::nil(), // sentinel — bypass token check for interrupt
1309                                )
1310                                .await;
1311                            let err_msg = format!(
1312                                "tool '{}' was interrupted mid-execution during a prior \
1313                                 worker failure; result is uncertain and was not re-run \
1314                                 (AtMostOnce safety)",
1315                                tool_call.name
1316                            );
1317                            let result_fp = tool_result_fingerprint(
1318                                &tool_call.name,
1319                                &ToolResult::error(&err_msg),
1320                            );
1321                            let _ = self
1322                                .event_emitter
1323                                .emit(EventRequest::new(
1324                                    context.session_id,
1325                                    event_context,
1326                                    ToolCompletedData::failure(
1327                                        tool_call.id.clone(),
1328                                        tool_call.name.clone(),
1329                                        "interrupted".to_string(),
1330                                        err_msg.clone(),
1331                                        None,
1332                                    )
1333                                    .with_fingerprints(tool_call_fingerprint.clone(), result_fp)
1334                                    .with_display_name(display_name.clone()),
1335                                ))
1336                                .await;
1337                            return ToolCallResult {
1338                                tool_call: tool_call.clone(),
1339                                result: ToolResult {
1340                                    tool_call_id: tool_call.id.clone(),
1341                                    result: None,
1342                                    images: None,
1343                                    error: Some(err_msg),
1344                                    connection_required: None,
1345                                    raw_output: None,
1346                                },
1347                                success: false,
1348                                status: "error".to_string(),
1349                                connection_required: None,
1350                                determinism_fatal: None,
1351                            };
1352                        }
1353                    }
1354                }
1355
1356                Ok(ToolCallClaimResult::DeterminismViolation {
1357                    stored_fingerprint,
1358                    current_fingerprint,
1359                }) => {
1360                    let err_msg = format!(
1361                        "determinism violation: tool '{}' args fingerprint changed \
1362                         on replay (stored={stored_fingerprint}, \
1363                         current={current_fingerprint})",
1364                        tool_call.name
1365                    );
1366                    tracing::error!(
1367                        session_id = %context.session_id,
1368                        turn_id = %context.turn_id,
1369                        tool_call_id = %tool_call.id,
1370                        stored = %stored_fingerprint,
1371                        current = %current_fingerprint,
1372                        "ActAtom: determinism violation on claim"
1373                    );
1374                    let result_fp =
1375                        tool_result_fingerprint(&tool_call.name, &ToolResult::error(&err_msg));
1376                    let _ = self
1377                        .event_emitter
1378                        .emit(EventRequest::new(
1379                            context.session_id,
1380                            event_context,
1381                            ToolCompletedData::failure(
1382                                tool_call.id.clone(),
1383                                tool_call.name.clone(),
1384                                "error".to_string(),
1385                                err_msg.clone(),
1386                                None,
1387                            )
1388                            .with_fingerprints(tool_call_fingerprint.clone(), result_fp)
1389                            .with_display_name(display_name.clone()),
1390                        ))
1391                        .await;
1392                    return ToolCallResult {
1393                        tool_call: tool_call.clone(),
1394                        result: ToolResult {
1395                            tool_call_id: tool_call.id.clone(),
1396                            result: None,
1397                            images: None,
1398                            error: Some(err_msg.clone()),
1399                            connection_required: None,
1400                            raw_output: None,
1401                        },
1402                        success: false,
1403                        status: "error".to_string(),
1404                        connection_required: None,
1405                        determinism_fatal: Some(err_msg),
1406                    };
1407                }
1408
1409                Err(e) => {
1410                    tracing::warn!(
1411                        session_id = %context.session_id,
1412                        tool_call_id = %tool_call.id,
1413                        error = %e,
1414                        "ActAtom: durable claim failed; proceeding without idempotency"
1415                    );
1416                    None
1417                }
1418            }
1419        } else {
1420            None
1421        };
1422
1423        // Emit tool.started event (child of act.started)
1424        if let Err(e) = self
1425            .event_emitter
1426            .emit(EventRequest::new(
1427                context.session_id,
1428                event_context.clone(),
1429                ToolStartedData {
1430                    tool_call: tool_call.clone(),
1431                    tool_call_fingerprint: Some(tool_call_fingerprint.clone()),
1432                    display_name: display_name.clone(),
1433                    narration: Some(self.render_tool_narration(
1434                        context,
1435                        tool_def,
1436                        &tool_call,
1437                        ToolNarrationPhase::Started,
1438                        locale,
1439                    )),
1440                },
1441            ))
1442            .await
1443        {
1444            tracing::warn!(
1445                session_id = %context.session_id,
1446                tool_call_id = %tool_call.id,
1447                error = %e,
1448                "ActAtom: failed to emit tool.started event"
1449            );
1450        }
1451
1452        // If tool definition not found, return error result
1453        let Some(tool_def) = tool_def else {
1454            let error_msg = format!("Tool definition not found: {}", tool_call.name);
1455            let tool_duration_ms = tool_start.elapsed().as_millis() as u64;
1456
1457            // Emit tool.completed event for error (child of act.started)
1458            if let Err(e) = self
1459                .event_emitter
1460                .emit(EventRequest::new(
1461                    context.session_id,
1462                    event_context,
1463                    ToolCompletedData::failure(
1464                        tool_call.id.clone(),
1465                        tool_call.name.clone(),
1466                        "error".to_string(),
1467                        error_msg.clone(),
1468                        Some(tool_duration_ms),
1469                    )
1470                    .with_fingerprints(
1471                        tool_call_fingerprint.clone(),
1472                        tool_error_fingerprint(&tool_call.name, "error", &error_msg),
1473                    )
1474                    .with_narration(Some(self.render_tool_narration(
1475                        context,
1476                        None,
1477                        &tool_call,
1478                        ToolNarrationPhase::Failed,
1479                        locale,
1480                    ))),
1481                ))
1482                .await
1483            {
1484                tracing::warn!(
1485                    session_id = %context.session_id,
1486                    tool_call_id = %tool_call.id,
1487                    error = %e,
1488                    "ActAtom: failed to emit tool.completed event"
1489                );
1490            }
1491
1492            return ToolCallResult {
1493                tool_call: tool_call.clone(),
1494                result: ToolResult {
1495                    tool_call_id: tool_call.id.clone(),
1496                    result: None,
1497                    images: None,
1498                    error: Some(error_msg),
1499                    connection_required: None,
1500                    raw_output: None,
1501                },
1502                success: false,
1503                status: "error".to_string(),
1504                connection_required: None,
1505                determinism_fatal: None,
1506            };
1507        };
1508
1509        // Execute the tool (always with context so tools can emit progress events)
1510        let mut tool_context =
1511            ToolContext::from_services(context.session_id, &self.context_services);
1512        if let Some(resolver) = tool_context.connection_resolver.as_ref()
1513            && let Some(bound) = resolver.for_execution(context.input_message_id.uuid())
1514        {
1515            tool_context.connection_resolver = Some(bound);
1516        }
1517
1518        if let Some(scope) =
1519            tool_context.extension::<everruns_core::tool_context::ExecutionServicesExt>()
1520        {
1521            scope
1522                .0
1523                .bind(&mut tool_context, context.input_message_id.uuid());
1524        }
1525        if let Some(invoker) = tool_context.mcp_invoker.as_ref()
1526            && let Some(bound) = invoker.for_execution(context.input_message_id.uuid())
1527        {
1528            tool_context.mcp_invoker = Some(bound);
1529        }
1530        if let Some(authority) = tool_context.session_creation_authority.as_ref()
1531            && let Some(bound) = authority.for_execution(context.input_message_id.uuid())
1532        {
1533            tool_context.session_creation_authority = Some(bound);
1534        }
1535        // Key file I/O by the attached workspace when known: pin the file store
1536        // to the workspace so shared-workspace sessions address the workspace's
1537        // files, not the session's own keyspace. For the default 1:1 case this
1538        // is a transparent pass-through.
1539        if let Some(workspace_id) = context.workspace_id {
1540            tool_context.workspace_id = workspace_id;
1541            if let Some(store) = tool_context.file_store.take() {
1542                tool_context.file_store = Some(
1543                    crate::session_files::WorkspaceScopedFileSystem::wrap(store, workspace_id),
1544                );
1545            }
1546        }
1547        // Resolve model paths through the mount resolver (EVE-660): `/workspace`
1548        // is a mount + cwd, not a per-store prefix. Applied over the
1549        // workspace-keyed store so resolution sits above re-keying.
1550        if let Some(store) = tool_context.file_store.take() {
1551            tool_context.file_store = Some(crate::mount_fs::MountFs::wrap_if_needed(store));
1552        }
1553        tool_context.visible_tool_names = Some(visible_tool_names.clone());
1554        // Input network_access (per-session, merged from harness+agent+session) takes precedence
1555        tool_context.network_access = network_access
1556            .cloned()
1557            .or_else(|| self.context_services.network_access.clone());
1558        // Provide event emitter + context so tools can emit tool.progress events
1559        if tool_context.event_emitter.is_none() {
1560            tool_context.event_emitter =
1561                Some(Arc::new(self.event_emitter.clone()) as Arc<dyn EventEmitter>);
1562        }
1563        tool_context.bind_to_turn(event_context.clone());
1564        tool_context.tool_call_id = Some(tool_call.id.clone());
1565
1566        // Cooperative cancellation for this call. The guard fires when this
1567        // future is dropped — which is what a cancelled turn looks like from
1568        // here — and also on normal return, so the contract a tool sees is
1569        // simply "this call is over". Work the tool leaves running (a child
1570        // process, a detached watcher) can hold a clone and die with the call
1571        // instead of outliving it; dropping the future alone cannot tell it
1572        // anything, because a dropped future is never polled again.
1573        let call_cancellation = tokio_util::sync::CancellationToken::new();
1574        tool_context.cancellation = Some(call_cancellation.clone());
1575        let _cancel_on_call_end = call_cancellation.drop_guard();
1576
1577        let execution_tool_call = self.transform_tool_call_for_execution(tool_call.clone());
1578
1579        // Run pre-tool-use hooks (capability-contributed). They can mutate
1580        // the tool call or block execution entirely. First Block wins; the
1581        // tool is not invoked, and the synthetic error result flows through
1582        // the same completion/event path as a tool failure.
1583        let (execution_tool_call, pre_block_reason) = if self.pre_tool_hooks.is_empty() {
1584            (execution_tool_call, None)
1585        } else {
1586            match act_hooks::run_pre_tool_use_hooks(
1587                &self.pre_tool_hooks,
1588                execution_tool_call.clone(),
1589                tool_def,
1590                &tool_context,
1591            )
1592            .await
1593            {
1594                act_hooks::PreToolUseDecision::Continue(updated) => (updated, None),
1595                act_hooks::PreToolUseDecision::Block {
1596                    tool_call: blocked,
1597                    reason,
1598                    ..
1599                } => (blocked, Some(reason)),
1600            }
1601        };
1602
1603        let result = if let Some(reason) = pre_block_reason {
1604            tracing::warn!(
1605                session_id = %context.session_id,
1606                tool_call_id = %execution_tool_call.id,
1607                tool_name = %execution_tool_call.name,
1608                reason = %reason,
1609                "ActAtom: pre_tool_use hook blocked execution"
1610            );
1611            Ok(crate::tool_types::ToolResult {
1612                tool_call_id: execution_tool_call.id.clone(),
1613                result: None,
1614                images: None,
1615                error: Some(format!("blocked by pre_tool_use hook: {reason}")),
1616                connection_required: None,
1617                raw_output: None,
1618            })
1619        } else if tool_def.is_cpu_bound() {
1620            // CPU-bound / non-yielding in-process tools (e.g. the bash
1621            // interpreter) get their own task so a long synchronous burst
1622            // cannot starve the cooperative polling of I/O-bound tools running
1623            // alongside them in this act batch. On the multi-thread runtime the
1624            // spawned task can also progress on another worker thread.
1625            let executor = self.tool_executor.clone();
1626            let call = execution_tool_call.clone();
1627            let def = tool_def.clone();
1628            let ctx = tool_context.clone();
1629            match AbortOnDropJoinHandle::new(tokio::spawn(async move {
1630                executor.execute_with_context(&call, &def, &ctx).await
1631            }))
1632            .await
1633            {
1634                Ok(result) => result,
1635                Err(join_err) => Err(crate::error::AgentLoopError::tool(format!(
1636                    "tool task failed to complete: {join_err}"
1637                ))),
1638            }
1639        } else {
1640            self.tool_executor
1641                .execute_with_context(&execution_tool_call, tool_def, &tool_context)
1642                .await
1643        };
1644
1645        match result {
1646            Ok(mut tool_result) => {
1647                // Run post-tool-exec hooks (capability then final/infrastructure)
1648                act_hooks::run_post_tool_exec_hooks(
1649                    &self.post_tool_hooks,
1650                    &self.final_post_tool_hooks,
1651                    &execution_tool_call,
1652                    tool_def,
1653                    &mut tool_result,
1654                    &tool_context,
1655                )
1656                .await;
1657
1658                let tool_duration_ms = tool_start.elapsed().as_millis() as u64;
1659                let success = tool_result.error.is_none();
1660                let status = if success { "success" } else { "error" };
1661
1662                // Emit tool.completed event
1663                let completed_data = if success {
1664                    let result_fingerprint = tool_result_fingerprint(&tool_call.name, &tool_result);
1665                    // Convert result to ContentPart (text + optional images)
1666                    let mut result_content = tool_result
1667                        .result
1668                        .as_ref()
1669                        .map(|r| vec![ContentPart::tool_result_text(r)])
1670                        .unwrap_or_default();
1671                    // Append images as native Image content parts
1672                    if let Some(ref images) = tool_result.images {
1673                        for img in images {
1674                            result_content.push(ContentPart::Image(
1675                                crate::message::ImageContentPart::from_base64(
1676                                    &img.base64,
1677                                    &img.media_type,
1678                                ),
1679                            ));
1680                        }
1681                    }
1682                    ToolCompletedData::success(
1683                        tool_call.id.clone(),
1684                        tool_call.name.clone(),
1685                        result_content,
1686                        Some(tool_duration_ms),
1687                    )
1688                    .with_fingerprints(tool_call_fingerprint.clone(), result_fingerprint)
1689                    .with_display_name(display_name.clone())
1690                    .with_capability_attribution(
1691                        capability_attribution.as_ref().map(|(id, _)| id.clone()),
1692                        capability_attribution
1693                            .as_ref()
1694                            .and_then(|(_, name)| name.clone()),
1695                    )
1696                    .with_narration(Some(self.render_tool_narration(
1697                        context,
1698                        Some(tool_def),
1699                        &tool_call,
1700                        ToolNarrationPhase::Completed,
1701                        locale,
1702                    )))
1703                } else {
1704                    let result_fingerprint = tool_result_fingerprint(&tool_call.name, &tool_result);
1705                    ToolCompletedData::failure(
1706                        tool_call.id.clone(),
1707                        tool_call.name.clone(),
1708                        status.to_string(),
1709                        tool_result.error.clone().unwrap_or_default(),
1710                        Some(tool_duration_ms),
1711                    )
1712                    .with_fingerprints(tool_call_fingerprint.clone(), result_fingerprint)
1713                    .with_display_name(display_name.clone())
1714                    .with_capability_attribution(
1715                        capability_attribution.as_ref().map(|(id, _)| id.clone()),
1716                        capability_attribution
1717                            .as_ref()
1718                            .and_then(|(_, name)| name.clone()),
1719                    )
1720                    .with_narration(Some(self.render_tool_narration(
1721                        context,
1722                        Some(tool_def),
1723                        &tool_call,
1724                        ToolNarrationPhase::Failed,
1725                        locale,
1726                    )))
1727                };
1728
1729                if let Err(e) = self
1730                    .event_emitter
1731                    .emit(EventRequest::new(
1732                        context.session_id,
1733                        event_context.clone(),
1734                        completed_data,
1735                    ))
1736                    .await
1737                {
1738                    tracing::warn!(
1739                        session_id = %context.session_id,
1740                        tool_call_id = %tool_call.id,
1741                        error = %e,
1742                        "ActAtom: failed to emit tool.completed event"
1743                    );
1744                }
1745
1746                tracing::debug!(
1747                    session_id = %context.session_id,
1748                    tool_name = %tool_call.name,
1749                    tool_call_id = %tool_call.id,
1750                    success = %success,
1751                    "ActAtom: tool execution completed"
1752                );
1753
1754                // Settle the durable claim (EVE-530).
1755                if let (Some(store), Some(token)) = (&self.durable_tool_result_store, claim_token) {
1756                    let result_snapshot =
1757                        serde_json::to_value(&tool_result).unwrap_or(serde_json::Value::Null);
1758                    match store
1759                        .settle_tool_call(
1760                            &context.turn_id.to_string(),
1761                            &tool_call.id,
1762                            result_snapshot,
1763                            "settled",
1764                            token,
1765                        )
1766                        .await
1767                    {
1768                        Ok(false) => {
1769                            tracing::warn!(
1770                                session_id = %context.session_id,
1771                                tool_call_id = %tool_call.id,
1772                                "ActAtom: settle ownership check failed (task reclaimed)"
1773                            );
1774                        }
1775                        Err(e) => {
1776                            tracing::warn!(
1777                                session_id = %context.session_id,
1778                                tool_call_id = %tool_call.id,
1779                                error = %e,
1780                                "ActAtom: settle_tool_call failed"
1781                            );
1782                        }
1783                        Ok(true) => {}
1784                    }
1785                }
1786
1787                let conn_req = tool_result.connection_required.clone();
1788                ToolCallResult {
1789                    tool_call,
1790                    result: tool_result,
1791                    success,
1792                    status: status.to_string(),
1793                    connection_required: conn_req,
1794                    determinism_fatal: None,
1795                }
1796            }
1797            Err(e) => {
1798                let tool_duration_ms = tool_start.elapsed().as_millis() as u64;
1799                let error_msg = e.to_string();
1800
1801                // Emit tool.completed event for error
1802                if let Err(emit_err) = self
1803                    .event_emitter
1804                    .emit(EventRequest::new(
1805                        context.session_id,
1806                        event_context,
1807                        ToolCompletedData::failure(
1808                            tool_call.id.clone(),
1809                            tool_call.name.clone(),
1810                            "error".to_string(),
1811                            error_msg.clone(),
1812                            Some(tool_duration_ms),
1813                        )
1814                        .with_fingerprints(
1815                            tool_call_fingerprint.clone(),
1816                            tool_error_fingerprint(&tool_call.name, "error", &error_msg),
1817                        )
1818                        .with_display_name(display_name.clone())
1819                        .with_capability_attribution(
1820                            capability_attribution.as_ref().map(|(id, _)| id.clone()),
1821                            capability_attribution
1822                                .as_ref()
1823                                .and_then(|(_, name)| name.clone()),
1824                        )
1825                        .with_narration(Some(self.render_tool_narration(
1826                            context,
1827                            Some(tool_def),
1828                            &tool_call,
1829                            ToolNarrationPhase::Failed,
1830                            locale,
1831                        ))),
1832                    ))
1833                    .await
1834                {
1835                    tracing::warn!(
1836                        session_id = %context.session_id,
1837                        tool_call_id = %tool_call.id,
1838                        error = %emit_err,
1839                        "ActAtom: failed to emit tool.completed event"
1840                    );
1841                }
1842
1843                tracing::warn!(
1844                    session_id = %context.session_id,
1845                    tool_name = %tool_call.name,
1846                    tool_call_id = %tool_call.id,
1847                    error = %e,
1848                    "ActAtom: tool execution failed"
1849                );
1850
1851                ToolCallResult {
1852                    tool_call: tool_call.clone(),
1853                    result: ToolResult {
1854                        tool_call_id: tool_call.id.clone(),
1855                        result: None,
1856                        images: None,
1857                        error: Some(error_msg),
1858                        connection_required: None,
1859                        raw_output: None,
1860                    },
1861                    success: false,
1862                    status: "error".to_string(),
1863                    connection_required: None,
1864                    determinism_fatal: None,
1865                }
1866            }
1867        }
1868    }
1869}
1870
1871// ============================================================================
1872// Tests
1873// ============================================================================
1874
1875#[cfg(test)]
1876#[path = "act_tests.rs"]
1877mod tests;