| 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 |