返回 CodeWhale
inline_repl.rs
根目录 / crates / tui / src / core / engine / turn_loop / inline_repl.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_inline_repl_phase(
7 &mut self,
8 turn: &mut TurnContext,
9 tool_policy: &ToolSurfacePolicy,
10 progress: &mut TurnLoopProgress,
11 client: &SharedModelClient,
12 current_text_visible: &str,
13 has_sendable_assistant_content: bool,
14 ) -> PhaseResult<()> {
15 let tool_registry = Some(&tool_policy.registry);
16 // Inline ```repl execution — the normal Agent working kernel.
17 // The kernel is session-scoped: refresh its inspectable context
18 // for this model step, but preserve Python variables/imports
19 // from earlier steps. That keeps the simple `repl` route useful
20 // for sustained work instead of forcing the model through a
21 // separate open/eval/configure control surface.
22
23 // The kernel runs model-written Python, so it answers to the
24 // same command gate as `code_execution`: a narrowed tool
25 // surface (`exec --allowed-tools …`, or plain `exec`'s zero-tool
26 // surface, #6510) must not execute code through a fence.
27 // Plan mode withholds `code_execution` from the catalog, and a
28 // fence is not a way around that: it runs only when the tool is
29 // on this turn's surface, and only after the same approval.
30 let repl_fence_present = has_sendable_assistant_content
31 && crate::repl::sandbox::has_repl_block(current_text_visible);
32 let repl_fence_offered =
33 code_execution_offered(progress.mode, &progress.tool_catalog, tool_policy);
34 let mut repl_fence_skip_reason = (repl_fence_present && !repl_fence_offered)
35 .then(|| "code execution is not available on this turn".to_string());
36 let repl_blocks = if repl_fence_present && repl_fence_offered {
37 crate::repl::sandbox::extract_repl_blocks(current_text_visible)
38 } else {
39 Vec::new()
40 };
41 if !repl_blocks.is_empty() {
42 let approval_id = format!("{}-repl-{}", turn.id, turn.step);
43 repl_fence_skip_reason = if let Some(state) = self.rlm_host.as_ref() {
44 self.rlm_round_refusal(&repl_blocks[0].code, state.gate.as_ref())
45 .await
46 } else {
47 self.repl_fence_blocked_reason(
48 &repl_blocks,
49 "the reply's ```repl block(s) in the session REPL kernel",
50 &approval_id,
51 client.as_ref(),
52 turn,
53 tool_policy,
54 &progress.tool_catalog,
55 tool_registry,
56 &mut progress.active_tool_names,
57 &mut progress.tool_call_budget,
58 progress.mode,
59 progress.fleet_denial_guard.as_ref(),
60 )
61 .await
62 };
63 // Admission may have applied a pending posture change.
64 progress.mode = self.current_mode;
65 if self.turn_wall_clock.exhausted() {
66 let reason = "parent turn deadline exhausted before REPL execution".to_string();
67 repl_fence_skip_reason = Some(reason.clone());
68 progress.turn_error = Some(reason);
69 }
70 }
71 if let Some(reason) = repl_fence_skip_reason.as_deref() {
72 if self.rlm_host.is_some() {
73 return PhaseResult::Return((
74 TurnOutcomeStatus::Failed,
75 Some(format!("RLM code was not run: {reason}")),
76 ));
77 }
78 let _ = self
79 .send_event(Event::status(format!("REPL block not run: {reason}")))
80 .await;
81 }
82 if !repl_blocks.is_empty() && repl_fence_skip_reason.is_none() {
83 let child_deadline =
84 tokio::time::Instant::now() + crate::tools::subagent::DEFAULT_CHILD_WALL_TIME;
85 let repl_deadline = self
86 .nested_work_deadline()
87 .map_or(child_deadline, |parent| parent.min(child_deadline));
88 // A kernel left broken by a dropped turn refuses every round; kill it.
89 drop(self.repl_kernel.take_if(|kernel| kernel.is_broken()));
90 if self.repl_kernel.is_none() {
91 let startup = tokio::select! {
92 biased;
93 () = self.cancel_token.cancelled() => {
94 Err("REPL startup cancelled".into())
95 }
96 result = tokio::time::timeout_at(
97 repl_deadline,
98 crate::repl::runtime::PythonRuntime::new(),
99 ) => result.unwrap_or_else(|_| {
100 Err("parent turn deadline reached during REPL startup".into())
101 }),
102 };
103 self.repl_kernel = match startup {
104 Ok(runtime) => Some(runtime),
105 Err(e) => {
106 let _ = self
107 .send_event(Event::status(format!("REPL init failed: {e}")))
108 .await;
109 progress.turn_error = Some(format!("REPL init failed: {e}"));
110 return PhaseResult::Break;
111 }
112 };
113 }
114
115 if self.rlm_host.is_none() {
116 let kernel_context = self.repl_kernel_context();
117 let refresh_result = tokio::select! {
118 biased;
119 () = self.cancel_token.cancelled() => {
120 Err("REPL context refresh cancelled".into())
121 }
122 result = tokio::time::timeout_at(
123 repl_deadline,
124 self.repl_kernel
125 .as_mut()
126 .expect("REPL kernel initialized above")
127 .replace_context(&kernel_context),
128 ) => result.unwrap_or_else(|_| {
129 Err("parent turn deadline reached during REPL context refresh".into())
130 }),
131 };
132 if let Err(e) = refresh_result {
133 // A broken subprocess cannot be trusted to retain
134 // state. Drop it so a later model step gets a clean,
135 // freshly bootstrapped kernel instead of repeating a
136 // hidden failure.
137 self.repl_kernel = None;
138 let _ = self
139 .send_event(Event::status(format!("REPL context refresh failed: {e}")))
140 .await;
141 progress.turn_error = Some(format!("REPL context refresh failed: {e}"));
142 return PhaseResult::Break;
143 }
144 }
145
146 // Child queries use the same object-safe client as the
147 // root turn. This follows the user-selected provider and
148 // lets deterministic/injected hosts exercise the exact
149 // same kernel contract, rather than quietly dropping
150 // programmatic recursion outside the legacy DeepSeek
151 // client path.
152 //
153 // Depth 0: the approval above covers the code in the
154 // fence, not code a child model writes later. A nested
155 // `rlm(...)` from a fence degrades to a one-shot child
156 // completion (text back to Python) instead of starting a
157 // sub-RLM whose code rounds would run unapproved.
158 let captured = self
159 .rlm_host
160 .as_ref()
161 .map(|state| Arc::clone(&state.caller))
162 .or_else(|| {
163 self.live_tool_context(tool_registry)
164 .and_then(|context| context.rlm_caller.clone())
165 });
166 let bridge = captured.as_deref().map(|caller| {
167 let depth = self.rlm_host.as_ref().map_or(0, |state| match state.mode {
168 crate::core::engine::rlm_host::RlmMode::Completion => 0,
169 crate::core::engine::rlm_host::RlmMode::Recursive { depth_remaining } => {
170 depth_remaining
171 }
172 });
173 let usage = self
174 .rlm_host
175 .as_ref()
176 .map_or_else(crate::rlm::bridge::RlmUsageAccumulator::new, |state| {
177 state.usage.clone()
178 });
179 crate::rlm::RlmBridge::with_usage_accumulator(
180 caller,
181 depth,
182 Duration::from_secs(600),
183 usage,
184 )
185 .with_events(self.tx_event.clone())
186 .with_deadline(Some(repl_deadline))
187 .with_gate(self.rlm_host.as_ref().and_then(|state| state.gate.clone()))
188 });
189 let repl_started = Instant::now();
190
191 let mut final_result: Option<String> = None;
192 let mut kernel_failed = false;
193 let mut empty_cap_hit = false;
194 for (i, block) in repl_blocks.iter().enumerate() {
195 let round_num = i + 1;
196 let _ = self
197 .send_event(Event::status(format!(
198 "REPL round {round_num}: executing..."
199 )))
200 .await;
201
202 // Dropping the cancelled round also stops its owned
203 // RPC/forwarder futures. The ledger below still
204 // accounts completed and pending provider requests;
205 // the kernel is discarded after any round failure.
206 let round_result = tokio::select! {
207 biased;
208 () = self.cancel_token.cancelled() => {
209 Err("REPL execution cancelled".into())
210 }
211 result = tokio::time::timeout_at(
212 repl_deadline,
213 self.repl_kernel
214 .as_mut()
215 .expect("REPL kernel stays alive during a round")
216 .run(&block.code, bridge.as_ref()),
217 ) => result.unwrap_or_else(|_| {
218 Err("REPL execution reached the parent turn deadline".into())
219 }),
220 };
221
222 match round_result {
223 Ok(round) => {
224 if let Some(state) = self.rlm_host.as_mut() {
225 state.total_rpcs = state.total_rpcs.saturating_add(round.rpc_count);
226 let round_output = if round.has_error {
227 format!("stdout:\n{}\nstderr:\n{}", round.stdout, round.stderr)
228 } else {
229 round.stdout.clone()
230 };
231 let stdout = crate::rlm::turn::truncate_text(
232 &round_output,
233 crate::rlm::turn::STDOUT_METADATA_PREVIEW_LEN,
234 );
235 state.trace.push(crate::rlm::turn::RlmRoundTrace {
236 round: state.model_rounds,
237 code_summary: crate::rlm::turn::summarize_code(&block.code),
238 stdout_preview: stdout.clone(),
239 had_error: round.has_error,
240 rpc_count: round.rpc_count,
241 elapsed_ms: u64::try_from(round.elapsed.as_millis())
242 .unwrap_or(u64::MAX),
243 });
244 let feedback = crate::rlm::turn::metadata_text(
245 &state.prompt,
246 state.model_rounds,
247 Some(&block.code),
248 Some(&stdout),
249 );
250 if round.final_value.is_none() {
251 self.add_session_message(
252 self.runtime_text_message_with_turn_metadata(
253 feedback,
254 UserInputProvenance::Runtime,
255 ),
256 )
257 .await;
258 }
259 }
260 if let Some(val) = &round.final_value {
261 let _ = self
262 .send_event(Event::status(format!(
263 "REPL round {round_num}: FINAL result obtained"
264 )))
265 .await;
266 final_result = Some(val.clone());
267 break;
268 }
269
270 // Empty-round guard + provenance (PROMPT-repl-fence-fix.md parts 2 & 3).
271 // Detection stays prompt-only (has_repl_block unchanged) to preserve
272 // saved-transcript replay (tools/rlm.rs kept). Provenance makes clear
273 // the block was the assistant's own; empty rounds get guidance + a
274 // consecutive cap so the model cannot loop forever.
275 let is_empty_round = !round.has_error
276 && round.stdout.trim().is_empty()
277 && round.stderr.trim().is_empty()
278 && round.rpc_count == 0;
279 if is_empty_round {
280 progress.consecutive_empty_repl_rounds =
281 progress.consecutive_empty_repl_rounds.saturating_add(1);
282 let hit_cap = progress.consecutive_empty_repl_rounds >= 3;
283 let feedback = if hit_cap {
284 format!(
285 "[Your emitted ```repl block (round {round_num}) produced no observable output — print something, call a helper, or stop emitting REPL blocks and answer. No output for {consecutive_empty_repl_rounds} consecutive rounds; stopping empty loop]\n[0 child query RPC(s)]",
286 consecutive_empty_repl_rounds =
287 progress.consecutive_empty_repl_rounds
288 )
289 } else {
290 format!(
291 "[Your emitted ```repl block (round {round_num}) produced no observable output — print something, call a helper, or stop emitting REPL blocks and answer]\n[0 child query RPC(s)]"
292 )
293 };
294 self.add_session_message(self.runtime_text_message_with_turn_metadata(
295 feedback,
296 UserInputProvenance::Runtime,
297 ))
298 .await;
299 if hit_cap {
300 empty_cap_hit = true;
301 // Honest stop: do not continue the turn with a lying
302 // "stopping" string. The cap is real.
303 break;
304 }
305 } else if self.rlm_host.is_some() {
306 progress.consecutive_empty_repl_rounds = 0;
307 } else {
308 progress.consecutive_empty_repl_rounds = 0;
309 let provenance_prefix =
310 format!("Your emitted ```repl block (round {round_num}) result:");
311 let feedback = if round.has_error {
312 format!(
313 "{provenance_prefix} error\nstdout:\n{}\nstderr:\n{}",
314 round.stdout, round.stderr
315 )
316 } else {
317 format!(
318 "{provenance_prefix}\n[{} child query RPC(s)]\n{}",
319 round.rpc_count, round.stdout
320 )
321 };
322 self.add_session_message(self.runtime_text_message_with_turn_metadata(
323 feedback,
324 UserInputProvenance::Runtime,
325 ))
326 .await;
327 }
328 }
329 Err(e) => {
330 let _ = self
331 .send_event(Event::status(format!(
332 "REPL round {round_num} failed: {e}"
333 )))
334 .await;
335 self.add_session_message(self.runtime_text_message_with_turn_metadata(
336 format!("[REPL round {round_num} execution failed]\n{e}"),
337 UserInputProvenance::Runtime,
338 ))
339 .await;
340 if self.rlm_host.is_some() {
341 progress.turn_error =
342 Some(if tokio::time::Instant::now() >= repl_deadline {
343 format!("RLM original wall-clock deadline exhausted: {e}")
344 } else {
345 e.clone()
346 });
347 }
348 // A transport error or timeout means Python
349 // may still be executing unknown code. Do not
350 // send another block into that process or
351 // pretend its state is trustworthy.
352 kernel_failed = true;
353 break;
354 }
355 }
356 }
357
358 if kernel_failed {
359 self.repl_kernel = None;
360 if self.rlm_host.is_some() {
361 return PhaseResult::Return((
362 TurnOutcomeStatus::Failed,
363 progress
364 .turn_error
365 .take()
366 .or_else(|| Some("RLM Python execution failed".into())),
367 ));
368 }
369 }
370
371 // Programmatic child calls are real provider work, not
372 // implementation detail. Fold their authoritative usage
373 // into the parent turn exactly once, including failures
374 // after a partial fan-out, so `/cost`, goals, and the
375 // final receipt cannot undercount the working kernel.
376 if self.rlm_host.is_none()
377 && let Some(bridge) = bridge.as_ref()
378 {
379 let snapshot = bridge.usage_snapshot().await;
380 turn.add_usage(&snapshot.usage);
381 let residual_dropped_records = snapshot
382 .dropped_records
383 .saturating_sub(u64::try_from(snapshot.drop_records.len()).unwrap_or(u64::MAX));
384 turn.add_routed_usage_dropped_records(residual_dropped_records);
385 if usage_has_reported_data(&snapshot.usage) {
386 let _ = self
387 .send_event(Event::RoutedTurnUsage {
388 usage: snapshot.usage.clone(),
389 duration_ms: u64::try_from(repl_started.elapsed().as_millis())
390 .unwrap_or(u64::MAX),
391 first_token_ms: None,
392 request_ms: None,
393 })
394 .await;
395 }
396 }
397
398 if let Some(final_val) = final_result {
399 if let Some(state) = self.rlm_host.as_mut() {
400 state.final_answer = Some(final_val.clone());
401 state.termination = Some(crate::rlm::turn::RlmTermination::Final);
402 }
403 // Replace the assistant's text with the FINAL answer.
404 if let Some(last_msg) = self.session.messages.last_mut()
405 && last_msg.role == "assistant"
406 {
407 for block in &mut last_msg.content {
408 if let ContentBlock::Text { text, .. } = block {
409 *text = final_val;
410 break;
411 }
412 }
413 }
414 self.emit_session_updated().await;
415 return PhaseResult::Break;
416 }
417
418 if empty_cap_hit {
419 if let Some(state) = self.rlm_host.as_mut() {
420 state.termination = Some(crate::rlm::turn::RlmTermination::NoCode);
421 return PhaseResult::Return((
422 TurnOutcomeStatus::Failed,
423 Some("RLM: 3 consecutive empty REPL rounds".into()),
424 ));
425 }
426 // Empty cap already fed back with honest "stopping" text
427 // inside the round loop. End the turn now instead of
428 // letting the outer ladder synthesize another provider
429 // request.
430 return PhaseResult::Break;
431 }
432
433 // No FINAL — let the model iterate with the feedback.
434 let _ = self.send_event(Event::status(format!(
435 "Continuing — REPL round feedback (consecutive_empty={consecutive_empty_repl_rounds})", consecutive_empty_repl_rounds = progress.consecutive_empty_repl_rounds
436 )))
437 .await;
438 turn.next_step();
439 return PhaseResult::Retry;
440 }
441
442 PhaseResult::Ready(())
443 }
444 pub(super) async fn continue_rlm_model_step(
445 &mut self,
446 turn: &mut TurnContext,
447 tool_policy: &ToolSurfacePolicy,
448 progress: &mut TurnLoopProgress,
449 client: &SharedModelClient,
450 response: AcceptedModelStep,
451 ) -> PhaseResult<AcceptedModelStep> {
452 let text = response.current_text_visible;
453 let state = self.rlm_host.as_mut().expect("RLM host");
454 state.last_response = text.clone();
455 if state.mode == crate::core::engine::rlm_host::RlmMode::Completion {
456 if text.trim().is_empty() {
457 return PhaseResult::Return((
458 TurnOutcomeStatus::Failed,
459 Some("empty LLM response".into()),
460 ));
461 }
462 state.final_answer = Some(text);
463 state.termination = Some(crate::rlm::turn::RlmTermination::Final);
464 return PhaseResult::Break;
465 }
466 let text_final = crate::rlm::turn::parse_text_final(&text);
467 if state.total_rpcs > 0
468 && let Some(answer) = text_final.as_ref()
469 {
470 state.final_answer = Some(answer.clone());
471 state.termination = Some(crate::rlm::turn::RlmTermination::Final);
472 return PhaseResult::Break;
473 }
474 let Some(code) = crate::rlm::turn::extract_repl_code(&text) else {
475 state.consecutive_no_code = state.consecutive_no_code.saturating_add(1);
476 if state.consecutive_no_code >= crate::rlm::turn::MAX_CONSECUTIVE_NO_CODE {
477 state.termination = Some(crate::rlm::turn::RlmTermination::NoCode);
478 if text_final.is_some() {
479 return PhaseResult::Break;
480 }
481 return PhaseResult::Return((
482 TurnOutcomeStatus::Failed,
483 Some("RLM: model failed to emit ```repl after 3 consecutive rounds".into()),
484 ));
485 }
486 let reminder = "Reminder: emit Python inside a ```repl … ``` fence, inspect `_context` through bounded helpers, and call finalize(value) when done. A prose FINAL before any helper RPC does not demonstrate use of the input.";
487 self.add_session_message(self.runtime_text_message_with_turn_metadata(
488 reminder.to_string(),
489 UserInputProvenance::Runtime,
490 ))
491 .await;
492 turn.next_step();
493 return PhaseResult::Retry;
494 };
495 state.consecutive_no_code = 0;
496 let canonical_fence = format!("```repl\n{code}\n``` ");
497 match self
498 .run_inline_repl_phase(turn, tool_policy, progress, client, &canonical_fence, true)
499 .await
500 {
501 PhaseResult::Ready(()) => PhaseResult::Return((
502 TurnOutcomeStatus::Failed,
503 Some("RLM code was not available on the captured surface".into()),
504 )),
505 PhaseResult::Retry => PhaseResult::Retry,
506 PhaseResult::Break => PhaseResult::Break,
507 PhaseResult::Return(outcome) => PhaseResult::Return(outcome),
508 }
509 }
510
511 async fn rlm_round_refusal(
512 &self,
513 code: &str,
514 gate: Option<&crate::tools::codemode::NestedCallGate>,
515 ) -> Option<String> {
516 use crate::tools::codemode::NestedCallVerdict;
517 let name = tool_catalog::CODE_EXECUTION_TOOL_NAME;
518 let Some(gate) = gate else {
519 return Some("no permission gate is serving this RLM turn".into());
520 };
521 match gate
522 .ask(name.to_string(), serde_json::json!({"code":code}))
523 .await
524 {
525 NestedCallVerdict::Run {
526 name: admitted,
527 input,
528 ..
529 } if admitted == name
530 && input.get("code").and_then(serde_json::Value::as_str) == Some(code) =>
531 {
532 None
533 }
534 NestedCallVerdict::Refused { error, .. } => Some(error.to_string()),
535 _ => Some("the admitted call was not this code".into()),
536 }
537 }
538 }
539
539 lines RUST