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