返回 CodeWhale
tool_batch.rs
根目录 / crates / tui / src / core / engine / turn_loop / tool_batch.rs
1 //! One private phase of the existing Engine turn loop.
2
3 use super::*;
4
5 impl Engine {
6 pub(super) async fn run_tool_batch_phase(
7 &mut self,
8 turn: &mut TurnContext,
9 tool_policy: &ToolSurfacePolicy,
10 progress: &mut TurnLoopProgress,
11 client: &SharedModelClient,
12 response: AcceptedModelStep,
13 ) -> PhaseResult<()> {
14 let tool_registry = Some(&tool_policy.registry);
15 let AcceptedModelStep {
16 current_text_visible,
17 mut tool_uses,
18 mut pending_steers,
19 mut output_limit_truncated,
20 fleet_report_response,
21 fleet_no_progress_report,
22 ..
23 } = response;
24 // A user can change Ask / Auto-Review / Full Access while the
25 // provider is streaming. Apply the newest typed authority before
26 // planning this tool batch; already-running tools are never
27 // retroactively reclassified.
28 let authority_changed_before_tools = self.apply_pending_runtime_authority().await;
29 if authority_changed_before_tools {
30 // A response requested as report-only never acquires execution
31 // authority after it streamed. Reset only after pairing its
32 // suppressed calls; the next response can use the new posture.
33 if !fleet_report_response && let Some(guard) = progress.fleet_denial_guard.as_mut() {
34 guard.reset();
35 turn.stop_diagnostics
36 .permission_denial_rounds_without_progress = 0;
37 }
38 progress.mode = self.current_mode;
39 }
40
41 // Execute tools
42 if self.shared_paused.lock().is_ok_and(|paused| *paused) {
43 let _ = self.send_event(Event::status("Request was Paused")).await;
44 self.add_interrupted_assistant_text(&current_text_visible)
45 .await;
46 return PhaseResult::Return((TurnOutcomeStatus::Interrupted, None));
47 }
48
49 let tool_exec_lock = self.tool_exec_lock.clone();
50 let mcp_pool = if !fleet_report_response
51 && tool_uses.iter().any(|tool| {
52 McpPool::is_mcp_tool(&tool.name) || tool.name == EXECUTE_TOOLS_TOOL_NAME
53 }) {
54 match self.ensure_mcp_pool().await {
55 Ok(pool) => Some(pool),
56 Err(err) => {
57 let _ = self.send_event(Event::status(err.to_string())).await;
58 None
59 }
60 }
61 } else {
62 None
63 };
64
65 // Tool discovery may be the first action after a model request
66 // that overlapped MCP startup. Search the ready catalog now.
67 self.refresh_boot_mcp_catalog(
68 tool_policy,
69 &mut progress.tool_catalog,
70 &mut progress.active_tool_names,
71 )
72 .await;
73 // Parked: per-tool timeouts, approvals, and the UI tool-hang
74 // watchdog own a tool batch's bound.
75 self.turn_heartbeat.enter(
76 super::turn_heartbeat::TurnPhase::Tools,
77 Some(
78 tool_uses
79 .iter()
80 .map(|tool| tool.name.as_str())
81 .collect::<Vec<_>>()
82 .join(", "),
83 ),
84 None,
85 );
86 let PlannedToolCalls {
87 plans,
88 hook_contexts,
89 batch_sandbox_policy,
90 } = self
91 .plan_tool_calls(
92 client.as_ref(),
93 turn,
94 tool_policy,
95 &mut tool_uses,
96 &progress.tool_catalog,
97 tool_registry,
98 &mut progress.active_tool_names,
99 &mut progress.tool_call_budget,
100 progress.mode,
101 progress.fleet_denial_guard.as_ref(),
102 ToolCallSource::Model,
103 )
104 .await;
105
106 let origin_turn_id = turn.id.clone();
107 let mut nested_gate_env = NestedGateEnv {
108 client: client.as_ref(),
109 turn: &mut *turn,
110 tool_policy,
111 tool_call_budget: &mut progress.tool_call_budget,
112 fleet_denial_guard: progress.fleet_denial_guard.as_ref(),
113 authority_changed: false,
114 };
115 let (outcomes, authority_changed_during_tools) = self
116 .execute_planned_tools(
117 plans,
118 &origin_turn_id,
119 &current_text_visible,
120 &mut progress.tool_catalog,
121 &mut progress.active_tool_names,
122 tool_registry,
123 tool_exec_lock,
124 mcp_pool,
125 &batch_sandbox_policy,
126 &mut progress.mode,
127 &mut nested_gate_env,
128 )
129 .await;
130
131 let authority_changed = authority_changed_before_tools || authority_changed_during_tools;
132 self.turn_heartbeat.enter(
133 super::turn_heartbeat::TurnPhase::Preparing,
134 None,
135 Some(super::turn_heartbeat::PREPARING_PHASE_BOUND),
136 );
137 let denial_action = self
138 .process_tool_results(
139 outcomes,
140 turn,
141 &mut progress.tool_catalog,
142 &mut progress.active_tool_names,
143 &hook_contexts,
144 if authority_changed || fleet_report_response {
145 None
146 } else {
147 progress.fleet_denial_guard.as_mut()
148 },
149 )
150 .await;
151
152 let accepted_steer_after_tools = !pending_steers.is_empty();
153 if !pending_steers.is_empty() {
154 for pending in pending_steers.drain(..) {
155 let steer = pending.commit().trim().to_string();
156 self.session
157 .working_set
158 .observe_user_message(&steer, &self.session.workspace);
159 self.add_session_message(self.user_text_message_with_turn_metadata(steer))
160 .await;
161 }
162 }
163
164 if authority_changed || accepted_steer_after_tools {
165 if let Some(guard) = progress.fleet_denial_guard.as_mut() {
166 guard.reset();
167 turn.stop_diagnostics
168 .permission_denial_rounds_without_progress = 0;
169 }
170 } else if fleet_no_progress_report {
171 // Exactly one accepted report response, including empty,
172 // reasoning-only, truncated or tool-producing responses.
173 if self.cancel_token.is_cancelled() {
174 return PhaseResult::Return((TurnOutcomeStatus::Interrupted, None));
175 }
176 turn.stop_diagnostics.last_response_tool_calls_suppressed = Some(tool_uses.len());
177 let error = if turn.budget_exhausted_final_report {
178 // One response can serve both report requests; the
179 // explicit budget retains its existing stop provenance.
180 format!(
181 "Maximum model steps reached before completion (limit: {}, {})",
182 turn.max_steps,
183 turn.budget_source.key_label()
184 )
185 } else {
186 turn.stop_diagnostics.reason = Some(TurnStopReason::NoProgress);
187 FLEET_NO_PROGRESS_STOP.to_string()
188 };
189 let _ = self.send_event(Event::status(error.clone())).await;
190 return PhaseResult::Return((TurnOutcomeStatus::Failed, Some(error)));
191 } else {
192 let notice = match denial_action {
193 FleetDenialAction::Continue => None,
194 FleetDenialAction::SwitchStrategy => {
195 turn.stop_diagnostics.permission_strategy_switches = turn
196 .stop_diagnostics
197 .permission_strategy_switches
198 .saturating_add(1);
199 Some(FLEET_STRATEGY_SWITCH_NOTICE)
200 }
201 FleetDenialAction::FinalReport => {
202 turn.stop_diagnostics.final_report_requested = true;
203 Some(FLEET_FINAL_REPORT_NOTICE)
204 }
205 };
206 if let Some(notice) = notice {
207 // Dynamic guard facts are append-only runtime history;
208 // BASE_PROMPT and the session's pinned prefix stay intact.
209 self.add_session_message(self.runtime_text_message_with_turn_metadata(
210 notice.to_string(),
211 UserInputProvenance::Runtime,
212 ))
213 .await;
214 }
215 }
216
217 // Surface an output-limit truncation after the tool result so the
218 // transcript stays well-formed (a `tool_result` must follow the
219 // assistant `tool_use` directly) and the model can act on it.
220 if let Some(reason) = output_limit_truncated.take() {
221 self.add_session_message(
222 self.runtime_text_message_with_turn_metadata(
223 format!(
224 "[runtime] The provider stopped generation at its output limit (`{reason}`) before completing. Your last response was cut off. Continue from where you left off; do not repeat content already delivered."
225 ),
226 UserInputProvenance::Runtime,
227 ),
228 )
229 .await;
230 }
231
232 // A successful tool step is productive progress, not a runaway
233 // synthetic resume. Declared per-task tool budgets and max_steps
234 // remain the explicit limits for tool-driven work.
235 let _ = self
236 .send_event(Event::status("Continuing — tool results".to_string()))
237 .await;
238 turn.next_step();
239
240 PhaseResult::Ready(())
241 }
242 }
243
243 lines RUST