返回 CodeWhale
preparation.rs
根目录 / crates / tui / src / core / engine / turn_loop / preparation.rs
1 //! One private phase of the existing Engine turn loop.
2
3 use super::*;
4
5 impl Engine {
6 pub(super) async fn prepare_model_step(
7 &mut self,
8 turn: &mut TurnContext,
9 tool_policy: &ToolSurfacePolicy,
10 progress: &mut TurnLoopProgress,
11 client: &SharedModelClient,
12 inspection_surface: Option<&crate::tool_inspection::ToolSurfaceContext>,
13 ) -> PhaseResult<PreparedModelStep> {
14 if self.rlm_host.is_some() {
15 return self
16 .prepare_rlm_model_step(turn, client, inspection_surface)
17 .await;
18 }
19 let strict_tool_mode = tool_policy.strict_tool_mode;
20
21 if self.cancel_token.is_cancelled() {
22 let _ = self.send_event(Event::status("Request cancelled")).await;
23 return PhaseResult::Return((TurnOutcomeStatus::Interrupted, None));
24 }
25 self.turn_heartbeat.enter(
26 super::turn_heartbeat::TurnPhase::Preparing,
27 None,
28 Some(super::turn_heartbeat::PREPARING_PHASE_BOUND),
29 );
30 self.refresh_boot_mcp_catalog(
31 tool_policy,
32 &mut progress.tool_catalog,
33 &mut progress.active_tool_names,
34 )
35 .await;
36 self.record_mcp_server_instructions(&progress.tool_catalog)
37 .await;
38 self.record_current_extension_prompt_contributions().await;
39
40 // R1: the cumulative per-turn wall-clock budget. Checked at the
41 // provider-request boundary so a turn that runs out of time stops
42 // before authorizing another billable request, and every tool
43 // result already produced stays in the transcript. Hitting it is
44 // never a clean success — the turn ends `Failed` with the limit
45 // named, matching how the step ceiling below reports.
46 if let Some(error) = self.turn_wall_clock_exhausted_error() {
47 let _ = self.send_event(Event::status(error.clone())).await;
48 return PhaseResult::Return((TurnOutcomeStatus::Failed, Some(error)));
49 }
50
51 if self.apply_pending_runtime_authority().await {
52 if let Some(guard) = progress.fleet_denial_guard.as_mut() {
53 guard.reset();
54 turn.stop_diagnostics
55 .permission_denial_rounds_without_progress = 0;
56 }
57 progress.mode = self.current_mode;
58 }
59
60 let mut accepted_steer = false;
61 while let Some(pending) = self.next_turn_steer() {
62 if pending.content.trim().is_empty() {
63 // Nothing to deliver; dropping `pending` settles it.
64 continue;
65 }
66 let steer = pending.commit().trim().to_string();
67 accepted_steer = true;
68 self.session
69 .working_set
70 .observe_user_message(&steer, &self.session.workspace);
71 self.add_session_message(self.user_text_message_with_turn_metadata(steer.clone()))
72 .await;
73 let _ = self
74 .send_event(Event::status(format!(
75 "Steer input accepted: {}",
76 summarize_text(&steer, 120)
77 )))
78 .await;
79 }
80 if accepted_steer && let Some(guard) = progress.fleet_denial_guard.as_mut() {
81 guard.reset();
82 turn.stop_diagnostics
83 .permission_denial_rounds_without_progress = 0;
84 }
85
86 // Child agents can finish while the parent model is still taking
87 // tool steps. Surface queued completions before the next provider
88 // request so the parent can use them immediately instead of
89 // discovering them only when it eventually emits no more tools or
90 // the idle handler starts a separate follow-up turn.
91 if !self.is_acp_turn() {
92 self.drain_subagent_completion_events("queued").await;
93 }
94
95 // The pinned system + tools prefix is frozen for the session:
96 // recomposing it here from disk on every tool step is exactly what
97 // kills DeepSeek's KV prefix cache once the agent writes a file
98 // (the project pack listing changes -> the system hash changes ->
99 // the next same-turn request is a full miss). Header changes come
100 // only from explicit ops (`/model`, mode, goal, session sync),
101 // which refresh under a declared reason. Volatile facts the model
102 // must see mid-turn (LSP diagnostics, steer input, subagent
103 // completions) are appended to history above, never spliced into
104 // the frozen prefix.
105 // A zero-tool turn (plain `exec`) spends extra steps only on
106 // output-limit continuations. It has no work to wrap up or report
107 // on, so the agent wrap-up notices below would only bend a
108 // one-shot answer; at the limit it ends honestly instead.
109 let zero_tool_turn = progress.tool_catalog.is_empty();
110 if let Some(job) = self.child_job() {
111 if job.authority.runtime.cancel_token.is_cancelled() {
112 return PhaseResult::Return((TurnOutcomeStatus::Interrupted, None));
113 }
114 if !self.child_report() {
115 if let Some(notice) = job.pacing_notice() {
116 self.add_session_message(self.user_text_message_with_turn_metadata(notice))
117 .await;
118 }
119 if job
120 .work_deadline
121 .is_some_and(|deadline| Instant::now() >= deadline)
122 {
123 job.stop_for_budget("child wall-time work budget exhausted");
124 return PhaseResult::Return((TurnOutcomeStatus::Interrupted, None));
125 }
126 if turn.at_max_steps() {
127 job.stop_for_budget("child model-step work budget exhausted");
128 return PhaseResult::Return((
129 TurnOutcomeStatus::Failed,
130 Some("child model-step work budget exhausted".into()),
131 ));
132 }
133 }
134 }
135 // A1 soft landing: with a finite step budget, once ~80% of it is
136 // spent tell the model once to stop exploring and write its final
137 // report. Savings proved out by the grok-style parity work (ops
138 // A1): a step-faithful harness ends mid-report far too often.
139 if self.child_host.is_none()
140 && !zero_tool_turn
141 && !turn.stop_diagnostics.soft_landing_sent
142 && let Some(step_limit) = turn.step_limit()
143 && step_limit > 0
144 && turn.steps_used() >= ((step_limit as f32 * 0.8).floor() as u32).max(1)
145 {
146 turn.stop_diagnostics.soft_landing_sent = true;
147 let notice = format!(
148 "Step budget soft landing: you have used about {}% of your {} step budget ({}). Stop exploring; write your final, complete report now, in final form, with evidence.",
149 80,
150 turn.max_steps,
151 turn.budget_source.key_label(),
152 );
153 self.add_session_message(self.user_text_message_with_turn_metadata(notice))
154 .await;
155 let _ = self
156 .send_event(Event::status(
157 "Soft landing: wrap up with your final report",
158 ))
159 .await;
160 }
161
162 if turn.at_max_steps() {
163 turn.stop_diagnostics.reason = Some(TurnStopReason::StepBudgetExhausted);
164 if self.child_host.is_none()
165 && progress.step_budget_exhaustion_is_terminal
166 && !progress.final_report_sent
167 && !zero_tool_turn
168 {
169 // A2 report-on-exhaustion: the budget died while the model
170 // still owes work. Never finish silently — grant exactly
171 // one final provider turn to write a bounded report, then
172 // let the natural no-tool termination close the turn.
173 progress.final_report_sent = true;
174 turn.budget_exhausted_final_report = true;
175 let notice = format!(
176 "Your model-step budget was exhausted (limit: {}, {}). You cannot continue working. Write your final report now: what you did, what you proved or found, what remains, and exact evidence. This is your last turn.",
177 turn.max_steps,
178 turn.budget_source.key_label(),
179 );
180 self.add_session_message(self.user_text_message_with_turn_metadata(notice))
181 .await;
182 let _ = self
183 .send_event(Event::status(
184 "Model budget exhausted — final report requested",
185 ))
186 .await;
187 } else if !progress.step_budget_exhaustion_is_terminal {
188 return PhaseResult::Break;
189 } else {
190 let error = format!(
191 "Maximum model steps reached before completion (limit: {}, {})",
192 turn.max_steps,
193 turn.budget_source.key_label(),
194 );
195 let _ = self.send_event(Event::status(error.clone())).await;
196 return PhaseResult::Return((TurnOutcomeStatus::Failed, Some(error)));
197 }
198 }
199
200 // A tool-producing response can spend the remaining goal budget
201 // before this loop reaches the no-tool continuation check below.
202 // Stop at the provider-request boundary so tool results remain in
203 // the transcript, but no additional model request is authorized.
204 // GoalState remains untouched here: the outer turn bookkeeping
205 // records this usage once, then the normal cross-turn reconciler
206 // publishes the terminal Blocked projection.
207 // Token budget is advisory (unbounded) — surface telemetry but don't break.
208 // Like grokbuild/kimicode, only verifier completion/block or backstop ends the run.
209 if let Some(snapshot) = self.goal_snapshot_with_current_turn_usage(&turn.usage)
210 && let Some(budget) = snapshot.token_budget
211 && snapshot.tokens_used >= u64::from(budget)
212 {
213 let _ = self.send_event(Event::status(format!(
214 "Goal over token budget ({} / {budget} tokens) — continuing (unbounded); verify or /goal clear when done.",
215 snapshot.tokens_used
216 )))
217 .await;
218 }
219
220 let active_tools = active_tools_for_request(
221 &progress.tool_catalog,
222 &progress.active_tool_names,
223 strict_tool_mode,
224 );
225 // The report already has a bounded text-only captured request. It
226 // cannot authorize an extra paid summarizer or mutate its evidence.
227 if !self.child_report() {
228 match self
229 .run_auto_compaction(
230 client.as_ref(),
231 active_tools.as_deref(),
232 turn,
233 &mut progress.auto_compaction_suppressed,
234 )
235 .await
236 {
237 AutoCompactionStep::Proceed => {}
238 AutoCompactionStep::Restart => return PhaseResult::Retry,
239 AutoCompactionStep::EndTurn(status, error) => {
240 return PhaseResult::Return((status, error));
241 }
242 }
243 }
244
245 // The guard measures what the compaction gate measures: the honest
246 // estimate, lifted to the provider's last bill plus the growth
247 // since it. The ×1.5-inflated overflow estimate, compared against
248 // the honest ceiling, refused at two
249 // thirds of the budget, and emergency compaction — which targets
250 // the honest budget — could never satisfy it (#6374). A request
251 // the estimate still undercounts is rejected by the provider and
252 // takes the bounded context-length recovery below.
253 let estimated_input = turn
254 .live_input_tokens_for_compaction(
255 &self.session.messages,
256 self.session.system_prompt.as_ref(),
257 self.session.latest_parent_input_tokens,
258 )
259 .and_then(|tokens| usize::try_from(tokens).ok())
260 .unwrap_or(0);
261 if let Some(budget) = route_context_budget_for_route(
262 self.api_provider,
263 &self.session.model,
264 self.active_route_limits,
265 estimated_input,
266 ) {
267 let input_budget = usize::try_from(budget.input_budget_ceiling).unwrap_or(usize::MAX);
268 let triggered = estimated_input > input_budget;
269 let output_ceiling =
270 crate::route_budget::output_ceiling_source(self.api_provider, &self.session.model);
271 let route_input_limit =
272 crate::route_budget::route_input_limit_tokens(self.active_route_limits);
273 let input_ceiling_source =
274 route_input_limit.map_or("window-minus-output-headroom", |limit| {
275 if u64::from(limit) <= budget.input_budget_ceiling {
276 "route-declared-input-limit"
277 } else {
278 "window-minus-output-headroom"
279 }
280 });
281 tracing::debug!(
282 target: "context_budget",
283 provider = self.api_provider.as_str(),
284 model = %self.session.model,
285 resolved_route_window_tokens = budget.window_tokens,
286 resolved_model_output_ceiling_tokens = ?output_ceiling.clamp_tokens(),
287 resolved_model_output_ceiling_source = output_ceiling.as_str(),
288 effective_request_output_cap_tokens = crate::route_budget::effective_max_output_tokens_for_turn(
289 self.api_provider,
290 &self.session.model,
291 self.active_route_limits,
292 turn.max_output_tokens,
293 ),
294 reserved_response_headroom_tokens = budget.output_cap_tokens,
295 safety_headroom_tokens = crate::context_budget::CONTEXT_HEADROOM_TOKENS,
296 resolved_route_input_limit_tokens = ?route_input_limit,
297 estimated_input_tokens = estimated_input,
298 input_budget_ceiling_tokens = budget.input_budget_ceiling,
299 input_budget_ceiling_source = input_ceiling_source,
300 remaining_input_budget_tokens = budget.available_input_tokens,
301 compaction_trigger_tokens = budget.compaction_trigger_tokens,
302 trigger = if triggered { "preflight-token-budget" } else { "none" },
303 "resolved route context budget"
304 );
305 if triggered {
306 if progress.context_recovery_attempts >= MAX_CONTEXT_RECOVERY_ATTEMPTS {
307 let message = context_overflow_exhausted_message(
308 self.config.terminal_chrome_enabled,
309 turn.stop_diagnostics.emergency_compaction_attempts,
310 estimated_input,
311 input_budget,
312 );
313 progress.turn_error = Some(message.clone());
314 let _ = self
315 .send_event(Event::error(ErrorEnvelope::context_overflow(message)))
316 .await;
317 return PhaseResult::Return((
318 TurnOutcomeStatus::Failed,
319 progress.turn_error.take(),
320 ));
321 }
322
323 if self
324 .recover_context_overflow(
325 client.as_ref(),
326 active_tools.as_deref(),
327 "preflight token budget",
328 turn,
329 )
330 .await
331 {
332 progress.context_recovery_attempts =
333 progress.context_recovery_attempts.saturating_add(1);
334 return PhaseResult::Retry;
335 }
336 if self.cancel_token.is_cancelled() {
337 return PhaseResult::Return((TurnOutcomeStatus::Interrupted, None));
338 }
339 // One failure, one true sentence (experience mark 2): a
340 // provider that refused the recovery request is the
341 // cause, and a history with nothing to summarize is a
342 // window problem, not a failed compaction.
343 if let Some(rejection) = turn.context_recovery_rejection.take() {
344 let display_message = self.decorate_auth_error_message(
345 initial_stream_error_user_message(&self.config.locale_tag, &rejection),
346 );
347 let mut envelope = crate::error_taxonomy::envelope_for_llm_error(
348 rejection,
349 display_message.clone(),
350 );
351 envelope.message = display_message.clone();
352 let _ = self.send_event(Event::error(envelope)).await;
353 return PhaseResult::Return((TurnOutcomeStatus::Failed, Some(display_message)));
354 }
355 let message = if crate::compaction::has_compactable_history(&self.session.messages)
356 {
357 "The request still exceeds this model's context budget and automatic recovery did not complete. The conversation is saved; retry or choose a larger context route.".to_string()
358 } else {
359 let prefix_tokens = crate::compaction::estimate_input_tokens_for_pressure(
360 &[],
361 self.session.system_prompt.as_ref(),
362 );
363 super::context::context_does_not_fit_message(
364 self.config.terminal_chrome_enabled,
365 self.api_provider == crate::config::ProviderKind::Ollama,
366 &self.session.model,
367 estimated_input,
368 input_budget,
369 prefix_tokens,
370 )
371 };
372 let _ = self
373 .send_event(Event::error(ErrorEnvelope::context_overflow(
374 message.clone(),
375 )))
376 .await;
377 return PhaseResult::Return((TurnOutcomeStatus::Failed, Some(message)));
378 }
379 }
380
381 // #136: drain any LSP diagnostics collected since the last
382 // request and inject them as a synthetic user message so the
383 // model sees compile errors before its next reasoning step.
384 self.flush_pending_lsp_diagnostics().await;
385
386 // Build the request. Tool selection goes through the same
387 // helper that seeded this turn and that `/preview-request`
388 // reports, so a deferred tool activated mid-turn is reflected
389 // identically in both places.
390 // Resolve `auto` reasoning_effort to a concrete tier (#663).
391 let effective_reasoning_effort = resolve_auto_effort(
392 self.session.reasoning_effort.as_deref(),
393 self.api_provider,
394 &self.api_config.active_route_base_url(),
395 &self.config.model,
396 );
397
398 // Check prefix-cache stability before building the request.
399 // This detects system-prompt or tool-set drift that would
400 // invalidate DeepSeek's KV prefix cache for this turn.
401 // Sends an event on EVERY check so the TUI can maintain
402 // its own counter for the stable-checks tally.
403 let declared_change = self.session.pending_prefix_change_reason.take();
404 if let Some(pm) = self.session.prefix_stability.as_mut() {
405 let system_text = codewhale_core::prefix_cache::system_prompt_text(
406 self.session.system_prompt.as_ref(),
407 );
408 let tools_ref: Option<&[codewhale_models::Tool]> = active_tools.as_deref();
409 let outcome = pm.check(&system_text, tools_ref, declared_change.as_deref());
410 // C5: request N's prefix may only diverge from N-1 across a
411 // DECLARED change. An undeclared drift means the pinned header
412 // moved without stamping a context update — the failure that
413 // silently kills the provider cache while stability claims
414 // still read well. The first check initializes the pin, so it
415 // is exempt.
416 #[cfg(debug_assertions)]
417 if pm.check_count() > 1
418 && declared_change.is_none()
419 && let codewhale_core::prefix_cache::PrefixCheck::Drift { change }
420 | codewhale_core::prefix_cache::PrefixCheck::Repinned { change, .. } = &outcome
421 {
422 debug_assert!(
423 false,
424 "prefix drift without a declared change (C5): the {} changed but no context update was stamped",
425 change.label()
426 );
427 }
428 let pinned_hash = pm
429 .pinned_fingerprint()
430 .map(|fp| fp.combined_sha256.clone())
431 .unwrap_or_default();
432 let stability_pct = (pm.stability_ratio() * 100.0).round() as u32;
433 let pin_reason = pm.pin_reason().unwrap_or_default().to_string();
434 let last_miss_reason = pm.last_miss_reason().unwrap_or_default().to_string();
435 let context_updates = pm.context_update_count();
436 let event = match outcome {
437 codewhale_core::prefix_cache::PrefixCheck::Stable => Event::PrefixCacheChange {
438 description: String::new(),
439 system_prompt_changed: false,
440 tools_changed: false,
441 stability_pct,
442 changed: false,
443 pinned_combined_hash: pinned_hash,
444 pin_reason,
445 last_miss_reason,
446 context_updates,
447 },
448 codewhale_core::prefix_cache::PrefixCheck::Repinned { reason, change } => {
449 // A declared header change re-pins under a logged
450 // reason: the miss is expected and attributable.
451 tracing::debug!(
452 target: "prefix_cache",
453 reason = %reason,
454 "prefix re-pinned: {}",
455 change.description()
456 );
457 Event::PrefixCacheChange {
458 description: format!("{reason} — {}", change.description()),
459 system_prompt_changed: change.system_changed,
460 tools_changed: change.tools_changed,
461 stability_pct,
462 changed: true,
463 pinned_combined_hash: pinned_hash,
464 pin_reason,
465 last_miss_reason,
466 context_updates,
467 }
468 }
469 codewhale_core::prefix_cache::PrefixCheck::Drift { change } => {
470 // Undeclared drift: the pin is kept so the same prefix
471 // keeps counting as a miss until an explicit op moves
472 // it. This should not happen after the mid-loop
473 // refresh removal — if it does it is a real bug.
474 tracing::warn!(
475 target: "prefix_cache",
476 "undeclared prefix drift (pin held): {}",
477 change.description()
478 );
479 Event::PrefixCacheChange {
480 description: format!("drift — {}", change.description()),
481 system_prompt_changed: change.system_changed,
482 tools_changed: change.tools_changed,
483 stability_pct,
484 changed: true,
485 pinned_combined_hash: pinned_hash,
486 pin_reason,
487 last_miss_reason,
488 context_updates,
489 }
490 }
491 };
492 let _ = self.send_event(event).await;
493 }
494
495 // Three-zone prefix contract (#2264): freeze baseline on first
496 // turn, verify against it on subsequent turns. Operates alongside
497 // PrefixStabilityManager as an independent diagnostic layer.
498 // Phase 3: emit a one-shot 'frozen' event on first turn.
499 // Drift is logged (tracing::debug!) but not re-emitted —
500 // PrefixStabilityManager already reports the change above.
501 let system_text =
502 codewhale_core::prefix_cache::system_prompt_text(self.session.system_prompt.as_ref());
503 let current_tools: &[codewhale_models::Tool] = active_tools.as_deref().unwrap_or_default();
504
505 match &self.session.frozen_prefix {
506 Some(frozen) => {
507 if let Err(drift) = frozen.verify(&system_text, current_tools) {
508 // Report drift; never replace the frozen baseline. The
509 // original freeze is the byte prefix the provider cache
510 // is keyed on — re-freezing here would make `/cache`
511 // look stable while the provider cache is already dead.
512 // A declared header change is re-pinned through the
513 // PrefixStabilityManager path above under a logged
514 // reason; the three-zone baseline stays put.
515 tracing::debug!(
516 target: "prefix_cache",
517 "three-zone drift (baseline held): {drift}"
518 );
519 }
520 }
521 None => {
522 let pinned =
523 PinnedPrefix::new(self.session.system_prompt.as_ref(), current_tools.to_vec());
524 let frozen = pinned.freeze();
525 let _ = self
526 .send_event(Event::PrefixCacheChange {
527 description: format!("frozen: {}", frozen.short_id()),
528 system_prompt_changed: false,
529 tools_changed: false,
530 stability_pct: 100,
531 changed: false,
532 pinned_combined_hash: frozen.hash().to_string(),
533 pin_reason: "initial".to_string(),
534 last_miss_reason: String::new(),
535 context_updates: 0,
536 })
537 .await;
538 self.session.frozen_prefix = Some(frozen);
539 }
540 }
541
542 let fleet_report_response = progress
543 .fleet_denial_guard
544 .as_ref()
545 .is_some_and(FleetDenialGuard::report_only);
546 // `take` is what keeps this request-scoped: the nudge is spent
547 // here and never reaches `self.session.messages`.
548 let request_nudge = progress.reasoning_only_nudge.take();
549 let mut request = prepare_primary_turn_request(PrimaryTurnRequest {
550 model: self.session.model.clone(),
551 messages: {
552 let mut messages = self.messages_with_turn_metadata();
553 if let Some(nudge) = request_nudge.as_ref() {
554 messages.push(self.runtime_text_message_with_turn_metadata(
555 nudge.clone(),
556 UserInputProvenance::Runtime,
557 ));
558 }
559 messages
560 },
561 max_tokens: crate::route_budget::effective_max_output_tokens_for_turn(
562 self.api_provider,
563 &self.session.model,
564 self.active_route_limits,
565 turn.max_output_tokens,
566 ),
567 system: self.session.system_prompt.clone(),
568 tools: active_tools.clone(),
569 tool_choice: if active_tools.is_some() {
570 if fleet_report_response || turn.budget_exhausted_final_report {
571 // Keep the pinned tool prefix; only this request's
572 // choice changes. Admission below also enforces this
573 // if a provider ignores the report-only request
574 // (C02-10: the step-budget final report included).
575 Some(json!("none"))
576 } else if strict_tool_mode {
577 Some(json!("required"))
578 } else {
579 Some(json!({ "type": "auto" }))
580 }
581 } else {
582 None
583 },
584 reasoning_effort: effective_reasoning_effort,
585 });
586 if let Some(report) = self
587 .child_host
588 .as_ref()
589 .and_then(|state| state.report.as_ref())
590 {
591 request.system = Some(report.system.clone());
592 request.messages = report.messages.clone();
593 request.max_tokens = report.output_tokens;
594 request.tools = None;
595 request.tool_choice = None;
596 }
597 if turn.max_output_tokens.is_some() {
598 request.max_tokens = request
599 .max_tokens
600 .min(client.effective_max_output_tokens(&self.session.model));
601 }
602 // Normalize images against the route this request is actually
603 // going to. Session history keeps the real image so that switching
604 // to a vision-capable model later makes it visible again; only the
605 // outbound copy is rewritten, and it is rewritten to text that says
606 // why rather than being dropped.
607 let fresh_images = crate::image_attach::images_since_last_user_prompt(&request.messages);
608 let stripped_images = crate::image_attach::strip_images_when_unsupported(
609 &mut request.messages,
610 self.active_route_capabilities.image_input,
611 &self.session.model,
612 );
613 if stripped_images > 0 {
614 crate::logging::warn(format!(
615 "{stripped_images} image block(s) replaced with text: model {} does not accept image input",
616 self.session.model
617 ));
618 if fresh_images > 0 && !progress.image_omission_notified {
619 progress.image_omission_notified = true;
620 let status = codewhale_localization::tr(
621 codewhale_localization::resolve_locale(&self.config.locale_tag),
622 codewhale_localization::MessageId::ImageInputOmitted,
623 )
624 .replace("{model}", &self.session.model)
625 .replace("{count}", &fresh_images.to_string());
626 let _ = self.send_event(Event::status(status)).await;
627 }
628 }
629 let tool_request_snapshot =
630 crate::tool_inspection::ToolInspectionSnapshot::from_prepared_request_with_surface(
631 &turn.id,
632 turn.step,
633 request.tools.as_deref(),
634 inspection_surface,
635 );
636 turn.last_request_snapshot = Some(tool_request_snapshot.clone());
637 turn.stop_diagnostics.route_context_window_tokens = self
638 .active_route_limits
639 .and_then(|limits| limits.context_tokens);
640
641 turn.stop_diagnostics.last_prepared_output_limit_tokens = Some(request.max_tokens);
642
643 // Stream the response. Keep the request around (cloned into the
644 // first call) so we can resend it on a transparent retry below
645 // when the wire dies before any content was streamed (#103).
646 let stream_request = request;
647 // Superfast Decision Gate (shadow mode, off by default). When
648 // SUPERFAST_ENABLED is set, this spawns a detached task that asks
649 // a small System One decision model about the user's turn and only
650 // logs the recommendation. It never changes routing, never skips
651 // the model call below, and never waits on the decision call.
652 // Fired only on the first model request of the turn, where the raw
653 // user message decides intent. See `crate::superfast`.
654 if turn.step == 0 {
655 // Detached on purpose: dropping the handle does not cancel it.
656 drop(crate::superfast::spawn_shadow_gate(
657 &self.api_config,
658 &stream_request.messages,
659 self.config.compaction.runtime_cost_owner.as_deref(),
660 &self.cancel_token,
661 ));
662 }
663 let _ = self
664 .send_event(Event::ToolRequestSnapshot {
665 snapshot: tool_request_snapshot,
666 })
667 .await;
668 if let Some(nudge) = request_nudge {
669 // The "Continuing — " prefix classifies the receipt as
670 // internal: durable clients keep it, collapsed.
671 let _ = self
672 .send_event(Event::status(format!(
673 "{REQUEST_NUDGE_RECEIPT_PREFIX}{nudge}"
674 )))
675 .await;
676 }
677 if let Some(mut route) = turn.pending_route.take() {
678 if let Some(billing) = route.billing.as_mut() {
679 // Freeze the exact provider-live row at CodeWhale's
680 // pre-permit application-dispatch boundary. This is an
681 // admission contract, not provider invoice-time evidence;
682 // a later cancellation/preparation failure has no usage
683 // and therefore contributes no usage cost.
684 let dispatched_at = chrono::Utc::now();
685 billing.dispatched_at = dispatched_at;
686 billing.provider_live_pricing = u64::try_from(dispatched_at.timestamp())
687 .ok()
688 .and_then(|dispatched_at_unix| {
689 crate::client::main_turn_pricing_quote_at(
690 self.codewhale_client.as_ref(),
691 route.provider,
692 &route.provider_identity,
693 &route.model,
694 billing.endpoint_fingerprint.as_deref()?,
695 dispatched_at_unix,
696 )
697 });
698 }
699 let _ = self
700 .send_event(Event::RouteDispatched {
701 turn_id: turn.id.clone(),
702 route,
703 })
704 .await;
705 }
706
707 PhaseResult::Ready(PreparedModelStep {
708 request: stream_request,
709 zero_tool_turn,
710 fleet_report_response,
711 })
712 }
713 /// The same pure primary request builder and the same dispatch phase,
714 /// with a whole bounded RLM history and no autonomous recovery producer.
715 async fn prepare_rlm_model_step(
716 &mut self,
717 turn: &mut TurnContext,
718 client: &SharedModelClient,
719 inspection_surface: Option<&crate::tool_inspection::ToolSurfaceContext>,
720 ) -> PhaseResult<PreparedModelStep> {
721 let state = self.rlm_host.as_ref().expect("RLM host");
722 if let Err(error) = state.caller.validate_live() {
723 return PhaseResult::Return((TurnOutcomeStatus::Failed, Some(error.to_string())));
724 }
725 if self.cancel_token.is_cancelled() {
726 return PhaseResult::Return((
727 TurnOutcomeStatus::Interrupted,
728 Some("RLM evaluation cancelled".into()),
729 ));
730 }
731 if let Some(error) = self.turn_wall_clock_exhausted_error() {
732 return PhaseResult::Return((TurnOutcomeStatus::Failed, Some(error)));
733 }
734 if turn.at_max_steps() {
735 let state = self.rlm_host.as_mut().expect("RLM host");
736 state.termination = Some(crate::rlm::turn::RlmTermination::Exhausted);
737 return PhaseResult::Return((
738 TurnOutcomeStatus::Failed,
739 Some(format!(
740 "RLM exhausted {} iterations without FINAL",
741 crate::rlm::turn::MAX_RLM_ITERATIONS
742 )),
743 ));
744 }
745 self.record_current_extension_prompt_contributions().await;
746 let estimated = turn
747 .live_input_tokens_for_compaction(
748 &self.session.messages,
749 self.session.system_prompt.as_ref(),
750 self.session.latest_parent_input_tokens,
751 )
752 .unwrap_or(0);
753 if route_context_budget_for_route(
754 self.api_provider,
755 &self.session.model,
756 self.active_route_limits,
757 usize::try_from(estimated).unwrap_or(usize::MAX),
758 )
759 .is_some_and(|budget| estimated > budget.input_budget_ceiling)
760 {
761 return PhaseResult::Return((TurnOutcomeStatus::Failed, Some("RLM whole round history exceeds the captured route context budget; history was not truncated".into())));
762 }
763 let mut request = prepare_primary_turn_request(PrimaryTurnRequest {
764 model: self.session.model.clone(),
765 messages: self.messages_with_turn_metadata(),
766 max_tokens: crate::route_budget::effective_max_output_tokens_for_turn(
767 self.api_provider,
768 &self.session.model,
769 self.active_route_limits,
770 turn.max_output_tokens,
771 ),
772 system: self.session.system_prompt.clone(),
773 tools: None,
774 tool_choice: None,
775 reasoning_effort: resolve_auto_effort(
776 self.session.reasoning_effort.as_deref(),
777 self.api_provider,
778 &self.api_config.active_route_base_url(),
779 &self.config.model,
780 ),
781 });
782 request.max_tokens = request
783 .max_tokens
784 .min(client.effective_max_output_tokens(&self.session.model));
785 let snapshot =
786 crate::tool_inspection::ToolInspectionSnapshot::from_prepared_request_with_surface(
787 &turn.id,
788 turn.step,
789 None,
790 inspection_surface,
791 );
792 turn.last_request_snapshot = Some(snapshot.clone());
793 turn.stop_diagnostics.last_prepared_output_limit_tokens = Some(request.max_tokens);
794 let _ = self
795 .send_event(Event::ToolRequestSnapshot { snapshot })
796 .await;
797 PhaseResult::Ready(PreparedModelStep {
798 request,
799 zero_tool_turn: true,
800 fleet_report_response: false,
801 })
802 }
803 }
804
804 lines RUST