use std::collections::HashSet;
use log::info;
use std::path::Path;
use std::sync::Arc;
use std::sync::atomic::AtomicU64;
use tokio::sync::mpsc;
static PARALLEL_READONLY_TOOL_BATCH_SEQ: AtomicU64 = AtomicU64::new(1);
use crate::agent::per_coord::PerCoordinator;
use crate::agent::plan_artifact::PlanStepExecutorKind;
use crate::config::AgentConfig;
use crate::memory::long_term_memory::LongTermMemoryRuntime;
use crate::sse::{SseEncoder, SsePayload, encode_message};
use crate::tool_registry;
use crate::tool_result::ToolEnvelopeContext;
use crate::types::{Message, Tool, ToolCall};
use crate::workspace::changelist::WorkspaceChangelist;
use crate::cm_agent::agent_turn::{
ToolBatchExecutionMode, ToolBatchModeParams, replay_force_serial_from_env,
resolve_tool_batch_execution_mode,
};
mod emit;
mod parallel_readonly;
mod run_command_guard;
mod serial;
use emit::{
emit_sse_tool_running, emit_timeline_log_sse, emit_tool_call_summary_sse,
emit_tool_result_sse_and_append, emit_turn_tool_phase_end_sse,
};
use parallel_readonly::execute_tools_parallel;
use serial::execute_tools_serial;
const LOG_TARGET: &str = "crabmate::execute_tools";
fn trace_parallel_tool_child_span(
tracing_turn: Option<&Arc<crate::observability::TracingChatTurn>>,
tool_call_id: &str,
) -> tracing::Span {
match tracing_turn {
Some(t) => {
let tool_call_id_label = t.record_tool_call_id_for_log(tool_call_id);
tracing::span!(
parent: t.span.id(),
tracing::Level::INFO,
"parallel_tool",
tool_call_id = %tool_call_id_label,
)
}
None => tracing::Span::none(),
}
}
pub(crate) struct WebExecuteCtx<'a> {
pub cfg: &'a Arc<AgentConfig>,
pub effective_working_dir: &'a Path,
pub workspace_is_set: bool,
pub read_file_turn_cache: Option<Arc<crate::read_file_turn_cache::ReadFileTurnCache>>,
pub control: crate::agent::agent_turn::TurnControlSink<'a>,
pub web_tool_ctx: Option<&'a tool_registry::WebToolRuntime>,
pub mcp_turn: Option<&'a crate::mcp::McpTurnHandle>,
pub workspace_changelist: Option<&'a Arc<WorkspaceChangelist>>,
pub request_chrome_trace: Option<Arc<crate::request_chrome_trace::RequestTurnTrace>>,
pub step_executor_constraint: Option<PlanStepExecutorKind>,
pub tools_defs_full: &'a [Tool],
pub turn_allow: Option<&'a HashSet<String>>,
pub long_term_memory: Option<Arc<LongTermMemoryRuntime>>,
pub long_term_memory_scope_id: Option<String>,
pub tracing_chat_turn: Option<Arc<crate::observability::TracingChatTurn>>,
pub request_audit: Option<Arc<crate::WebRequestAudit>>,
pub tool_outcome_recorder: Arc<crate::tool_stats::ToolOutcomeRecorder>,
pub handler_lookup: crate::tool_registry::HandlerLookupTable,
pub sync_default_sandbox_backend: Arc<dyn crate::tool_sandbox::SyncDefaultSandboxBackend>,
pub readonly_tool_ttl_cache: Arc<crate::readonly_tool_ttl_cache::ReadonlyToolTtlCache>,
}
pub(crate) use crate::cm_agent::agent_turn::{
ExecuteToolsBatchOutcome, dedup_readonly_tool_calls_count,
};
pub(super) struct EmitToolResultParams<'a> {
cfg: &'a Arc<AgentConfig>,
tool_outcome_recorder: &'a Arc<crate::tool_stats::ToolOutcomeRecorder>,
control: crate::agent::agent_turn::TurnControlSink<'a>,
tool_result_envelope_v1: bool,
name: &'a str,
args: &'a str,
id: &'a str,
result: String,
reflection_inject: Option<serde_json::Value>,
envelope_ctx: Option<ToolEnvelopeContext<'a>>,
}
pub(crate) fn sse_sender_closed(out: Option<&mpsc::Sender<String>>) -> bool {
out.is_some_and(|tx| tx.is_closed())
}
async fn abort_tool_batch_if_sse_closed(
out: Option<&mpsc::Sender<String>>,
reason: &'static str,
encoder: &dyn SseEncoder,
) -> bool {
if !sse_sender_closed(out) {
return false;
}
info!(target: LOG_TARGET, "{reason}");
emit_sse_tool_running(
out,
false,
"execute_tools::abort_tool_batch tool_running false",
encoder,
)
.await;
true
}
struct ExecuteToolsCommonCtx<'a> {
tool_calls: &'a [ToolCall],
per_coord: &'a mut PerCoordinator,
messages: &'a mut Vec<Message>,
cfg: &'a Arc<AgentConfig>,
effective_working_dir: &'a Path,
workspace_is_set: bool,
read_file_turn_cache: Option<Arc<crate::read_file_turn_cache::ReadFileTurnCache>>,
workspace_changelist: Option<&'a Arc<WorkspaceChangelist>>,
control: crate::agent::agent_turn::TurnControlSink<'a>,
tool_result_envelope_v1: bool,
web_tool_ctx: Option<&'a tool_registry::WebToolRuntime>,
mcp_turn: Option<&'a crate::mcp::McpTurnHandle>,
request_chrome_trace: Option<Arc<crate::request_chrome_trace::RequestTurnTrace>>,
step_executor_constraint: Option<PlanStepExecutorKind>,
tools_defs_full: &'a [Tool],
turn_allow: Option<&'a HashSet<String>>,
long_term_memory: Option<Arc<LongTermMemoryRuntime>>,
long_term_memory_scope_id: Option<String>,
tracing_chat_turn: Option<Arc<crate::observability::TracingChatTurn>>,
request_audit: Option<Arc<crate::WebRequestAudit>>,
tool_outcome_recorder: Arc<crate::tool_stats::ToolOutcomeRecorder>,
handler_lookup: crate::tool_registry::HandlerLookupTable,
sync_default_sandbox_backend: Arc<dyn crate::tool_sandbox::SyncDefaultSandboxBackend>,
readonly_tool_ttl_cache: Arc<crate::readonly_tool_ttl_cache::ReadonlyToolTtlCache>,
}
async fn per_execute_tools_common(ctx: ExecuteToolsCommonCtx<'_>) -> ExecuteToolsBatchOutcome {
let control = ctx.control.clone();
let out = control.out;
let sse_control_mirror = control.sse_control_mirror.clone();
let sse_encoder = control.sse_encoder.clone();
let force_serial = replay_force_serial_from_env();
emit_sse_tool_running(
out,
true,
"execute_tools::batch tool_running true",
sse_encoder.as_ref(),
)
.await;
let batch_mode = resolve_tool_batch_execution_mode(&ToolBatchModeParams {
force_serial,
workspace_is_set: ctx.workspace_is_set,
handler_lookup: &ctx.handler_lookup,
cfg: ctx.cfg.as_ref(),
tool_calls: ctx.tool_calls,
turn_allow: ctx.turn_allow,
});
let workspace_changed = match batch_mode {
ToolBatchExecutionMode::ParallelReadonlyBatch => {
crate::turn_replay_dump::append_decision_point_event_if_configured(
"tool_execution",
"tool_batch_execution_mode",
"parallel_readonly_batch",
"当前批次满足只读并行条件,采用并行只读批执行以提升吞吐",
serde_json::json!({
"force_serial": force_serial,
"workspace_is_set": ctx.workspace_is_set,
"tool_call_count": ctx.tool_calls.len(),
}),
"current_tool_batch",
None,
);
let outcome = execute_tools_parallel(ctx).await;
if matches!(outcome, ExecuteToolsBatchOutcome::AbortedSse) {
emit_sse_tool_running(
out,
false,
"execute_tools::batch aborted_after_parallel tool_running false",
sse_encoder.as_ref(),
)
.await;
return outcome;
}
false
}
ToolBatchExecutionMode::Serial => {
crate::turn_replay_dump::append_decision_point_event_if_configured(
"tool_execution",
"tool_batch_execution_mode",
"serial",
if force_serial {
"环境变量强制串行执行,关闭并行只读批"
} else {
"当前批次不满足并行条件,回退串行执行"
},
serde_json::json!({
"force_serial": force_serial,
"workspace_is_set": ctx.workspace_is_set,
"tool_call_count": ctx.tool_calls.len(),
}),
"current_tool_batch",
None,
);
if force_serial {
log::info!(
target: LOG_TARGET,
"CM_REPLAY_FORCE_SERIAL enabled: force serial tool execution"
);
crate::turn_replay_dump::append_turn_replay_event_json_if_configured(
"tool_batch_mode",
"force_serial",
Some(&serde_json::json!({
"source": "CM_REPLAY_FORCE_SERIAL",
"parallel_disabled": true
})),
);
}
let mut workspace_changed = false;
let outcome = execute_tools_serial(ctx, &mut workspace_changed).await;
if matches!(outcome, ExecuteToolsBatchOutcome::AbortedSse) {
emit_sse_tool_running(
out,
false,
"execute_tools::batch aborted_after_serial tool_running false",
sse_encoder.as_ref(),
)
.await;
return outcome;
}
workspace_changed
}
};
if let Some(tx) = out
&& workspace_changed
{
let _ = crate::sse::send_string_logged(
tx,
encode_message(SsePayload::WorkspaceChanged {
workspace_changed: true,
}),
"execute_tools::batch workspace_changed",
)
.await;
}
emit_turn_tool_phase_end_sse(out, sse_control_mirror.as_ref(), sse_encoder.as_ref()).await;
if let Some(tx) = out {
crate::sse::send_state_snapshot_sse(tx, serde_json::json!({"phase": "tool_batch_end"}))
.await;
}
emit_sse_tool_running(
out,
false,
"execute_tools::batch tool_running false",
sse_encoder.as_ref(),
)
.await;
ExecuteToolsBatchOutcome::Finished
}
pub(crate) async fn per_execute_tools_web(
tool_calls: &[ToolCall],
per_coord: &mut PerCoordinator,
messages: &mut Vec<Message>,
ctx: WebExecuteCtx<'_>,
) -> ExecuteToolsBatchOutcome {
let WebExecuteCtx {
cfg,
effective_working_dir,
workspace_is_set,
read_file_turn_cache,
control,
web_tool_ctx,
mcp_turn,
workspace_changelist,
request_chrome_trace,
step_executor_constraint,
tools_defs_full,
turn_allow,
long_term_memory,
long_term_memory_scope_id,
tracing_chat_turn,
request_audit,
tool_outcome_recorder,
handler_lookup,
sync_default_sandbox_backend,
readonly_tool_ttl_cache,
} = ctx;
let _tool_trace = request_chrome_trace
.as_ref()
.map(|t| t.enter_section("agent.tools_batch"));
per_execute_tools_common(ExecuteToolsCommonCtx {
tool_calls,
per_coord,
messages,
cfg,
effective_working_dir,
workspace_is_set,
read_file_turn_cache,
workspace_changelist,
control,
tool_result_envelope_v1: cfg.tool_transcript.tool_result_envelope_v1,
web_tool_ctx,
mcp_turn,
request_chrome_trace,
step_executor_constraint,
tools_defs_full,
turn_allow,
long_term_memory,
long_term_memory_scope_id,
tracing_chat_turn,
request_audit,
tool_outcome_recorder,
handler_lookup,
sync_default_sandbox_backend,
readonly_tool_ttl_cache,
})
.await
}