返回 CodeWhale
model_step.rs
根目录 / crates / tui / src / core / engine / turn_loop / model_step.rs
1 //! One private phase of the existing Engine turn loop.
2
3 use super::*;
4 use anyhow::anyhow;
5
6 impl Engine {
7 pub(super) async fn run_model_step(
8 &mut self,
9 turn: &mut TurnContext,
10 progress: &mut TurnLoopProgress,
11 client: &SharedModelClient,
12 prepared: PreparedModelStep,
13 ) -> PhaseResult<AcceptedModelStep> {
14 let PreparedModelStep {
15 request: stream_request,
16 zero_tool_turn,
17 fleet_report_response,
18 } = prepared;
19 let child_job = self.child_job();
20 if let Some(job) = child_job.as_ref() {
21 if let Err(error) = job.project(&self.session.messages, job.steps()).await {
22 return PhaseResult::Return((
23 TurnOutcomeStatus::Failed,
24 Some(format!(
25 "child checkpoint failed before dispatch: {error:#}"
26 )),
27 ));
28 }
29 if job.authority.runtime.cancel_token.is_cancelled() {
30 return PhaseResult::Return((TurnOutcomeStatus::Interrupted, None));
31 }
32 }
33 // Session metrics: the model call is measured from this dispatch
34 // instant (connection setup included), and time-to-first-token is
35 // the gap to the first content-bearing stream event.
36 let request_dispatched_at = Instant::now();
37 self.turn_heartbeat.enter(
38 super::turn_heartbeat::TurnPhase::AwaitingModel,
39 Some(format!(
40 "{} / {}",
41 self.api_provider_identity
42 .as_ref()
43 .map_or("unavailable", |identity| {
44 identity
45 .compatibility()
46 .map_or(identity.key.as_str(), |row| row.label)
47 }),
48 stream_request.model
49 )),
50 Some(awaiting_model_bound(&self.config)),
51 );
52 let mut rlm_dispatch = None;
53 let request_cancel = self.cancel_token.clone();
54 let mut child_dispatch = None;
55 let child_reporting = self.child_report();
56 let child_logical_step = child_job.as_ref().map(|job| {
57 if child_reporting {
58 job.steps().saturating_add(1)
59 } else {
60 turn.steps_used().saturating_add(1)
61 }
62 });
63 let retry_observation = self.request_retry_observation();
64 let transport_retries = retry_observation.retries.clone();
65 let stream_result = tokio::select! {
66 biased;
67 () = request_cancel.cancelled() => {
68 turn.stop_diagnostics.transport_retries = turn.stop_diagnostics.transport_retries
69 .saturating_add(transport_retries.load(std::sync::atomic::Ordering::Relaxed));
70 let _ = self.send_event(Event::status("Request cancelled")).await;
71 return PhaseResult::Return((TurnOutcomeStatus::Interrupted, None));
72 }
73 result = crate::llm_client::observe_request_retries(Some(retry_observation), async {
74 rlm_dispatch = match self.rlm_provider_request(client, &stream_request).await {
75 Ok(request) => request,
76 Err(error) => return Err(anyhow!(error)),
77 };
78 turn.stop_diagnostics.model_requests_started = turn
79 .stop_diagnostics
80 .model_requests_started
81 .saturating_add(1);
82 if let Some(job) = child_job.as_ref() {
83 let child = &job.authority;
84 let route = client.effective_route_envelope(&stream_request.model, chrono::Utc::now());
85 let source = format!("child:{}:turn:{}:request:{}:dispatch:{}", child.owner_agent_id, turn.id, turn.stop_diagnostics.model_requests_started, route.dispatched_at.timestamp_nanos_opt().unwrap_or_default());
86 child_dispatch = Some(job.dispatched(child_logical_step.expect("captured child step"), source, route));
87 }
88 if let Some(job) = child_job.as_ref() {
89 tokio::time::timeout(
90 job.authority.runtime.step_api_timeout,
91 client.create_message_stream(stream_request.clone()),
92 ).await.unwrap_or_else(|_| Err(anyhow::Error::new(LlmError::Timeout(
93 job.authority.runtime.step_api_timeout,
94 ))))
95 } else {
96 client.create_message_stream(stream_request.clone()).await
97 }
98 }) => result,
99 };
100 turn.stop_diagnostics.transport_retries = turn
101 .stop_diagnostics
102 .transport_retries
103 .saturating_add(transport_retries.load(std::sync::atomic::Ordering::Relaxed));
104 let stream = match stream_result {
105 Ok(s) => {
106 if let Some(job) = child_job.as_ref() {
107 job.provider_responded();
108 }
109 progress.context_recovery_attempts = 0;
110 // A model has the question now; a later credential
111 // failure in this turn (a token expiring mid-turn, say)
112 // must not take it back (#6566).
113 turn.unanswered_user_message = None;
114 s
115 }
116 Err(e) => {
117 if let Some(dispatch) = rlm_dispatch.as_mut() {
118 dispatch.settle_open_error(&e).await;
119 }
120 if let Some(job) = child_job.as_ref() {
121 job.provider_refused(&e);
122 job.response_settled(false);
123 if let Some(dispatch) = child_dispatch.as_mut() {
124 dispatch.settle_open_error(&e).await;
125 }
126 drop(child_dispatch.take());
127 }
128 if self.child_report() {
129 return PhaseResult::Return((
130 TurnOutcomeStatus::Failed,
131 Some("bounded Core report request failed".into()),
132 ));
133 }
134 // Replacement is permitted only if Core can still retract
135 // this exact unsent user message. A runtime note or rewrite
136 // cannot be silently replayed as the original task.
137 if let Some(mark) = turn.unanswered_user_message
138 && self.can_retract_unanswered_user_message(mark)
139 {
140 match self.admit_child_first_request_replacement(&e).await {
141 Ok(true) => match self.retract_child_unsent_message(mark).await {
142 Ok(true) => {
143 turn.unanswered_user_message = None;
144 self.emit_session_updated().await;
145 return PhaseResult::Return((
146 TurnOutcomeStatus::Failed,
147 Some(codewhale_config::persistence::redact_secrets(&format!(
148 "{e:#}; approved captured replacement queued"
149 ))),
150 ));
151 }
152 Ok(false) => {
153 self.clear_child_pending_route();
154 }
155 Err(error) => {
156 return PhaseResult::Return((
157 TurnOutcomeStatus::Failed,
158 Some(codewhale_config::persistence::redact_secrets(&format!(
159 "{e:#}; unsent child checkpoint failed: {error:#}"
160 ))),
161 ));
162 }
163 },
164 Err(error) => {
165 return PhaseResult::Return((
166 TurnOutcomeStatus::Failed,
167 Some(codewhale_config::persistence::redact_secrets(&format!(
168 "{e:#}; {error:#}"
169 ))),
170 ));
171 }
172 Ok(false) => {}
173 }
174 }
175 if let Some(recovery) = self.recover_child_request(progress, &e).await {
176 return recovery;
177 }
178 // Recovery/classification keeps its existing input. Expanding
179 // diagnostics must not introduce another model request.
180 let message = self.decorate_auth_error_message(e.to_string());
181 if self.rlm_host.is_none()
182 && is_context_length_error_message(&message)
183 && progress.context_recovery_attempts < MAX_CONTEXT_RECOVERY_ATTEMPTS
184 && self
185 .recover_context_overflow(
186 client.as_ref(),
187 stream_request.tools.as_deref(),
188 "provider context-length rejection",
189 turn,
190 )
191 .await
192 {
193 progress.context_recovery_attempts =
194 progress.context_recovery_attempts.saturating_add(1);
195 return PhaseResult::Retry;
196 }
197 if self.rlm_host.is_none()
198 && is_image_input_rejection_message(&message)
199 && self.active_route_capabilities.image_input != CapabilityState::Unsupported
200 && !progress.image_rejection_recovered
201 {
202 progress.image_rejection_recovered = true;
203 // This path tells the user itself; the resend must
204 // not announce the same omission a second time.
205 progress.image_omission_notified = true;
206 self.active_route_capabilities.image_input = CapabilityState::Unsupported;
207 crate::logging::warn(format!(
208 "model {} rejected image content; resending with images replaced by text",
209 self.session.model
210 ));
211 let status = codewhale_localization::tr(
212 codewhale_localization::resolve_locale(&self.config.locale_tag),
213 codewhale_localization::MessageId::ImageInputRejectedResent,
214 )
215 .replace("{model}", &self.session.model);
216 let _ = self.send_event(Event::status(status)).await;
217 return PhaseResult::Retry;
218 }
219 let display_message = self.decorate_auth_error_message(
220 initial_stream_error_user_message(&self.config.locale_tag, &e),
221 );
222 // Classified from the error's types across its whole
223 // context chain, before `e` moves into the envelope: an
224 // adapter's outer context must not hide a connect error,
225 // and a provider's HTTP rejection must not pass for one
226 // because its body text mentions a timeout (#6711).
227 let open_transport_failure = crate::client::is_stream_open_transport_failure(&e);
228 let mut envelope =
229 crate::error_taxonomy::envelope_for_llm_error(e, message.clone());
230 // #6699: the request never became a stream (connect
231 // failure, response-header stall). The transport layer
232 // already spent its own retries; re-issue the identical
233 // request from here through the same bounded resume
234 // budget every other stream failure spends. Nothing
235 // streamed, so there is no fragment to keep or discard,
236 // and no error event is emitted for an attempt that is
237 // retried — an exhausted budget falls through to the
238 // normal failure below. Only a failure with no response
239 // headers qualifies; a provider rejection never does.
240 if self.rlm_host.is_none()
241 && open_transport_failure
242 && !self.cancel_token.is_cancelled()
243 && let Some(attempt) = progress.stream_retry_budget.authorize()
244 {
245 turn.stop_diagnostics.stream_resumes =
246 turn.stop_diagnostics.stream_resumes.saturating_add(1);
247 let _ = self.send_retry_status(format!(
248 "Retry attempt: stream-open {attempt}/{}; connection failed before response headers",
249 progress.stream_retry_budget.limit()
250 )).await;
251 crate::logging::warn(format!(
252 "Stream failed to open (attempt {attempt}/{}); retrying request: {message}",
253 progress.stream_retry_budget.limit()
254 ));
255 return PhaseResult::Retry;
256 }
257 if open_transport_failure && progress.stream_retry_budget.spent() > 0 {
258 let _ = self.send_retry_status(format!(
259 "Retry exhaustion: stream-open stopped after {} retries; connection failed before response headers",
260 progress.stream_retry_budget.spent()
261 )).await;
262 }
263 envelope.message = display_message.clone();
264 // #6566: no model saw the question. Take it back out of
265 // the session before reporting, so the next request does
266 // not send it twice and a resumed session does not show
267 // it twice; the code tells the host to hand the text back.
268 if envelope.category == ErrorCategory::Authentication
269 && let Some(mark) = turn.unanswered_user_message.take()
270 {
271 match self.retract_child_unsent_message(mark).await {
272 Ok(true) => {
273 envelope.code =
274 crate::error_taxonomy::CREDENTIAL_REJECTED_UNSENT_CODE.to_string();
275 self.emit_session_updated().await;
276 }
277 Ok(false) => {}
278 Err(error) => {
279 return PhaseResult::Return((
280 TurnOutcomeStatus::Failed,
281 Some(format!(
282 "{display_message}; unsent child checkpoint failed: {error:#}"
283 )),
284 ));
285 }
286 }
287 }
288 progress.turn_error = Some(display_message);
289 let _ = self.send_event(Event::error(envelope)).await;
290 return PhaseResult::Return((
291 TurnOutcomeStatus::Failed,
292 progress.turn_error.take(),
293 ));
294 }
295 };
296 let StreamOutcome {
297 current_text_raw,
298 current_text_visible,
299 current_thinking,
300 current_thinking_signature,
301 current_thinking_state,
302 mut tool_uses,
303 usage,
304 usage_reported,
305 stop_reason,
306 pending_message_complete,
307 last_text_index,
308 stream_errors,
309 terminal_stream_error,
310 pending_steers,
311 pending_resume,
312 stream_start,
313 first_token_at,
314 request_dispatched_at,
315 stream_error,
316 } = self
317 .process_stream(
318 client.as_ref(),
319 stream,
320 &stream_request,
321 request_dispatched_at,
322 progress.stream_retry_budget.spent(),
323 &mut turn.stop_diagnostics,
324 )
325 .await;
326 self.turn_heartbeat.enter(
327 super::turn_heartbeat::TurnPhase::Preparing,
328 None,
329 Some(super::turn_heartbeat::PREPARING_PHASE_BOUND),
330 );
331 if let Some(dispatch) = rlm_dispatch.as_mut() {
332 dispatch
333 .settle(&usage, stream_error.is_none() && stream_errors == 0)
334 .await;
335 }
336 // C02-05: a response whose stream failed — a provider error frame,
337 // a transport error, a stall, or a cap — is not a complete
338 // response. Unless the retry below re-issues the request, nothing
339 // it collected may execute or continue the turn.
340 let response_stream_failed = stream_error.is_some();
341 progress.turn_error = progress.turn_error.take().or(stream_error.clone());
342 turn.stop_diagnostics
343 .observe_provider_response(stop_reason.as_deref(), tool_uses.len());
344 // Counts and terminal metadata only: never log messages, tool
345 // arguments, credentials, or raw provider bodies.
346 tracing::debug!(
347 target: "provider_response_diagnostics",
348 model_request = turn.stop_diagnostics.model_requests_started,
349 prepared_output_limit_tokens = stream_request.max_tokens,
350 finish_reason = ?turn.stop_diagnostics.last_provider_finish_reason,
351 reported_usage = usage_reported,
352 input_tokens = usage.input_tokens,
353 output_tokens = usage.output_tokens,
354 cached_input_tokens = ?usage.prompt_cache_hit_tokens,
355 reasoning_tokens = ?usage.reasoning_tokens,
356 decoded_tool_calls = tool_uses.len(),
357 visible_text_chars = current_text_visible.chars().count(),
358 "parent model response settled"
359 );
360 // These belong to post-stream response assembly, not stream
361 // consumption: blocks are built from the completed stream state,
362 // and truncation is derived from its terminal stop reason below.
363 let mut content_blocks: Vec<ContentBlock> = Vec::new();
364 let mut output_limit_truncated: Option<String> = None;
365
366 // Account for every provider response before deciding whether to
367 // retry or accept it. A terminal stop reason followed by a
368 // transport error is still a billed, incomplete response; it must
369 // not be discarded and re-issued.
370 if let Some(dispatch) = child_dispatch.as_mut() {
371 dispatch
372 .settle(
373 &usage,
374 !response_stream_failed && !self.cancel_token.is_cancelled(),
375 )
376 .await;
377 }
378 drop(child_dispatch.take());
379 if let Some(job) = child_job.as_ref() {
380 job.response_settled(!response_stream_failed);
381 }
382 turn.add_parent_usage(&usage);
383 turn.note_parent_prompt_len(self.session.messages.len());
384 self.session.latest_parent_input_tokens = turn.latest_parent_input_tokens;
385 if usage_reported {
386 let _ = self
387 .send_event(Event::TurnUsage {
388 max_output_tokens: turn.max_output_tokens.map(|_| stream_request.max_tokens),
389 usage: usage.clone(),
390 duration_ms: u64::try_from(stream_start.elapsed().as_millis())
391 .unwrap_or(u64::MAX),
392 first_token_ms: first_token_at.map(|at| {
393 u64::try_from(
394 at.saturating_duration_since(request_dispatched_at)
395 .as_millis(),
396 )
397 .unwrap_or(u64::MAX)
398 }),
399 request_ms: Some(
400 u64::try_from(request_dispatched_at.elapsed().as_millis())
401 .unwrap_or(u64::MAX),
402 ),
403 })
404 .await;
405 }
406
407 if let Some(state) = self.rlm_host.as_mut() {
408 state.usage.record_nested_event(serde_json::json!({"run_id": state.run_id, "depth_remaining": match state.mode { rlm_host::RlmMode::Completion => 0, rlm_host::RlmMode::Recursive { depth_remaining } => depth_remaining }, "kind":"code", "content":current_text_visible})).await;
409 if is_incomplete_stop_reason(stop_reason.as_deref())
410 || stream_error.is_some()
411 || stream_errors > 0
412 || !tool_uses.is_empty()
413 {
414 let error = if !tool_uses.is_empty() {
415 "RLM provider emitted a structured tool call outside its pure Python surface"
416 .to_string()
417 } else {
418 format!(
419 "Model response incomplete: {}",
420 stream_error
421 .as_deref()
422 .unwrap_or_else(|| stop_reason_detail(stop_reason.as_deref()))
423 )
424 };
425 self.add_interrupted_assistant_text(&current_text_visible)
426 .await;
427 return PhaseResult::Return((TurnOutcomeStatus::Failed, Some(error)));
428 }
429 // Rejected fragments remain in the interrupted Session/code receipt,
430 // but cannot replace the last admitted answer or execute FINAL/REPL.
431 state.last_response = current_text_visible.clone();
432 }
433 if response_stream_failed && self.child_host.is_some() && !self.child_report() {
434 let error = if progress
435 .turn_error
436 .as_deref()
437 .is_some_and(|reason| reason.contains("child step API request timed out"))
438 {
439 anyhow::Error::new(LlmError::Timeout(
440 child_job
441 .as_ref()
442 .expect("captured child")
443 .authority
444 .runtime
445 .step_api_timeout,
446 ))
447 } else {
448 anyhow!(
449 "{}",
450 progress
451 .turn_error
452 .as_deref()
453 .unwrap_or("child stream failed")
454 )
455 };
456 if !terminal_stream_error
457 && let Some(recovery) = self.recover_child_request(progress, &error).await
458 {
459 return recovery;
460 }
461 }
462 if !response_stream_failed {
463 progress.child_request_retries = Default::default();
464 }
465
466 // The injected Child transport owns its response grammar, including
467 // an explicitly configured Custom wire or an approved replacement.
468 let protocol = self
469 .child_request_protocol()
470 .or_else(|| {
471 self.active_route_endpoint
472 .as_ref()
473 .map(|endpoint| endpoint.protocol)
474 })
475 .unwrap_or(codewhale_config::provider::WireFormat::ChatCompletions);
476 // No tool observation or replayable tool history is published before
477 // this whole-response admission. Usage and visible text remain real.
478 if let Err(error) = crate::client::validate_tool_call_ids_for_protocol(
479 protocol,
480 tool_uses.iter().map(|tool| tool.id.as_str()),
481 ) {
482 self.add_interrupted_assistant_text(&current_text_visible)
483 .await;
484 return PhaseResult::Return((TurnOutcomeStatus::Failed, Some(error.to_string())));
485 }
486
487 if self.cancel_token.is_cancelled() {
488 let _ = self.send_event(Event::status("Request cancelled")).await;
489 self.add_interrupted_assistant_text(&current_text_visible)
490 .await;
491 return PhaseResult::Return((TurnOutcomeStatus::Interrupted, None));
492 }
493
494 if is_incomplete_stop_reason(stop_reason.as_deref()) {
495 let reason = stop_reason_detail(stop_reason.as_deref());
496 if self.child_host.is_none()
497 && is_output_limit_stop_reason(stop_reason.as_deref())
498 && stream_errors == 0
499 {
500 // Degrade, don't kill the turn — but only when the stream
501 // finished cleanly. A `max_tokens` stop followed by a
502 // transport error is a billed incomplete response: charge
503 // it and fail closed instead of continuing into a second
504 // request. A generation limit on a complete stream is a
505 // normal provider outcome, not an unrecoverable error:
506 // accept whatever complete tool call or content was
507 // produced and continue. The truncation is surfaced as a
508 // bounded observation after the partial assistant message
509 // is committed (and, for a tool-call response, after the
510 // tool result is appended) so the transcript stays
511 // well-formed.
512 crate::logging::warn(format!(
513 "Model output truncated: provider stop reason `{reason}`; accepting partial response and continuing the turn."
514 ));
515 output_limit_truncated = Some(reason.to_string());
516 // Fall through to the normal content/tool dispatch below.
517 } else {
518 self.settle_unadmitted_tool_calls(&tool_uses, &incomplete_tool_result(reason))
519 .await;
520 // Do not emit MessageComplete: hosts must retain the visible
521 // fragment as interrupted/failed rather than recording it as
522 // a completed assistant item.
523 self.add_interrupted_assistant_text(&current_text_visible)
524 .await;
525 let error = if self.child_host.is_some() {
526 crate::tools::subagent::incomplete_subagent_response_failure(
527 stop_reason.as_deref(),
528 )
529 } else {
530 format!(
531 "Model response incomplete: provider stop reason `{reason}`; no complete response or tool call was accepted."
532 )
533 };
534 crate::logging::warn(&error);
535 return PhaseResult::Return((TurnOutcomeStatus::Failed, Some(error)));
536 }
537 }
538
539 // #103 Phase 3 — transparent retry. The inner loop above bails
540 // when reqwest yields chunk decode errors three times in a row;
541 // most of the time those are recoverable proxy / HTTP/2 issues
542 // and the request can simply be re-issued. Re-issue silently up
543 // to MAX_STREAM_RETRIES, but only when the stream produced
544 // nothing actionable — if any tool call landed or text was
545 // streamed, ship the partial state to the rest of the turn
546 // pipeline so we don't double-bill the user by re-running it.
547 // The post-content exceptions to that rule are the #2990
548 // sleep-resume and the mid-stream network-drop resumes: those
549 // discard the uncommitted fragment unless an operator watched
550 // visible text land (see `StreamResume::InteractiveNetworkDrop`).
551 //
552 // The resume itself is typed state, consumed here by value, so
553 // one drop schedules exactly one retry; and no resume path
554 // appends a synthetic user message to the persisted
555 // conversation — the retried request is the persisted
556 // conversation re-issued, nothing else.
557 if self.child_report() {
558 let refusal = if !tool_uses.is_empty() {
559 Some("provider returned a tool call; no call executed")
560 } else if output_limit_truncated.is_some() {
561 Some("report did not finish: provider output limit")
562 } else if response_stream_failed {
563 Some("provider call failed; no complete report accepted")
564 } else if current_text_visible.trim().is_empty() {
565 Some("report did not finish: empty response")
566 } else {
567 None
568 };
569 if let Some(refusal) = refusal {
570 return PhaseResult::Return((
571 TurnOutcomeStatus::Failed,
572 Some(format!("bounded Core report refused: {refusal}")),
573 ));
574 }
575 }
576 let stream_died_with_nothing = stream_errors > 0
577 && !terminal_stream_error
578 && tool_uses.is_empty()
579 && current_text_visible.trim().is_empty()
580 && current_thinking.trim().is_empty()
581 && !pending_message_complete;
582 let pending_resume = match pending_resume {
583 Some(resume) => Some(resume),
584 None if stream_died_with_nothing => Some(StreamResume::NoContentStreamDeath),
585 None => None,
586 };
587 if self.rlm_host.is_none()
588 && self.child_host.is_none()
589 && !terminal_stream_error
590 && let Some(resume) = pending_resume
591 && let Some(attempt) = progress.stream_retry_budget.authorize()
592 {
593 let limit = progress.stream_retry_budget.limit();
594 turn.stop_diagnostics.stream_resumes =
595 turn.stop_diagnostics.stream_resumes.saturating_add(1);
596 let reason = match resume {
597 StreamResume::AfterSleep => "system sleep interrupted the stream",
598 StreamResume::HeadlessNetworkDrop => "network drop; incomplete reply discarded",
599 StreamResume::InteractiveNetworkDrop => {
600 "network drop; continuing the admitted reply"
601 }
602 StreamResume::NoContentStreamDeath => {
603 "stream ended before a response was completed"
604 }
605 };
606 let _ = self
607 .send_retry_status(format!(
608 "Retry attempt: stream-resume {attempt}/{limit}; {reason}"
609 ))
610 .await;
611 match resume {
612 StreamResume::AfterSleep => {
613 crate::logging::warn(format!(
614 "Resuming after system sleep (attempt {attempt}/{limit}); discarding partial output and retrying request"
615 ));
616 // Finalize any partially-rendered assistant cell so
617 // the retried stream renders fresh instead of
618 // appending to the pre-sleep fragment.
619 if pending_message_complete {
620 let index = last_text_index.unwrap_or(0);
621 let _ = self.send_event(Event::MessageComplete { index }).await;
622 }
623 }
624 StreamResume::HeadlessNetworkDrop => {
625 crate::logging::warn(format!(
626 "Resuming headless turn after mid-stream network drop (attempt {attempt}/{limit}); discarding partial output and retrying request"
627 ));
628 }
629 StreamResume::InteractiveNetworkDrop => {
630 // Commit the partial assistant message so the retried
631 // request sees the prefix as already delivered. Build
632 // the blocks inline; the outer `content_blocks`
633 // variable is still empty at this point and will be
634 // rebuilt on the next round.
635 let mut resume_blocks: Vec<ContentBlock> = Vec::new();
636 // A wire-only placeholder must not ride into the
637 // retry prefix as stored reasoning either.
638 let thinking_is_placeholder_only =
639 crate::client::is_reasoning_replay_placeholder(&current_thinking);
640 if (!current_thinking.is_empty() && !thinking_is_placeholder_only)
641 || current_thinking_state.is_some()
642 {
643 resume_blocks.push(ContentBlock::Thinking {
644 thinking: current_thinking.clone(),
645 signature: current_thinking_signature.clone(),
646 state: current_thinking_state.clone(),
647 });
648 }
649 if !current_text_visible.is_empty() {
650 resume_blocks.push(ContentBlock::Text {
651 text: current_text_visible.clone(),
652 cache_control: None,
653 });
654 }
655 for tool in &tool_uses {
656 resume_blocks.push(ContentBlock::ToolUse {
657 execution_id: Some(tool.execution_id.clone()),
658 id: tool.id.clone(),
659 name: tool.name.clone(),
660 input: tool.input.clone(),
661 caller: tool.caller.clone(),
662 thought_signature: tool.thought_signature.clone(),
663 });
664 }
665 let has_sendable_assistant_content = resume_blocks.iter().any(|block| {
666 matches!(
667 block,
668 ContentBlock::Text { .. } | ContentBlock::ToolUse { .. }
669 )
670 });
671 if !has_sendable_assistant_content {
672 // Thinking-only drop: nothing visible streamed, so
673 // nothing is preserved and nothing is committed.
674 // The re-issued request is identical to the one
675 // that died. Neither the log line nor the status
676 // copy may claim a partial reply was preserved —
677 // that claim is what minted the fake `[runtime]`
678 // user turn in session 1589c05d.
679 crate::logging::warn(format!(
680 "Resuming interactive turn after mid-stream network drop (attempt {attempt}/{limit}); only hidden reasoning streamed — no partial reply to preserve, retrying request"
681 ));
682 } else {
683 crate::logging::warn(format!(
684 "Resuming interactive turn after mid-stream network drop (attempt {attempt}/{limit}); preserving partial reply and retrying request"
685 ));
686 // Finalize the partial text cell so the UI stops
687 // streaming and the retried content lands in a
688 // fresh cell instead of appending to an
689 // unfinished one.
690 if let Some(index) = last_text_index {
691 let _ = self.send_event(Event::MessageComplete { index }).await;
692 }
693 // Persist the fragment the operator already saw —
694 // exactly one assistant cell for it, and no
695 // synthetic user turn after it. The retried
696 // request therefore ends with this fragment, which
697 // is the provider-neutral "continue from here"
698 // contract; retry status receipts stay outside the
699 // session messages and provider request history.
700 // They never manufacture a user turn.
701 self.add_session_message(Message {
702 role: Role::Assistant,
703 content: resume_blocks,
704 })
705 .await;
706 }
707 }
708 StreamResume::NoContentStreamDeath => {
709 crate::logging::warn(format!(
710 "Stream died with no content (attempt {attempt}/{limit}); retrying request"
711 ));
712 }
713 }
714 // Don't preserve the per-stream `turn_error` — we're
715 // about to retry, and a successful retry should not
716 // surface the transient error as the turn outcome.
717 progress.turn_error = None;
718 return PhaseResult::Retry;
719 }
720 if pending_resume.is_some() {
721 if progress.stream_retry_budget.spent() > 0 {
722 let _ = self
723 .send_retry_status(format!(
724 "Retry exhaustion: stream-resume stopped after {} retries; stream interrupted",
725 progress.stream_retry_budget.spent()
726 ))
727 .await;
728 }
729 crate::logging::warn(format!(
730 "Stream retry budget exhausted ({} attempts); failing turn",
731 progress.stream_retry_budget.spent()
732 ));
733 } else if stream_errors == 0 {
734 if progress.stream_retry_budget.spent() > 0 && pending_message_complete {
735 let _ = self
736 .send_retry_status(format!(
737 "Retry recovery: stream recovered after {} retries",
738 progress.stream_retry_budget.spent()
739 ))
740 .await;
741 }
742 // Healthy round → reset retry budget so we don't carry over
743 // state from a previous bad round.
744 progress.stream_retry_budget.reset();
745 }
746
747 let mut final_text = current_text_visible.clone();
748 if tool_uses.is_empty() && tool_parser::has_tool_call_markers(&current_text_raw) {
749 let parsed = tool_parser::parse_tool_calls(&current_text_raw);
750 if parsed.tool_calls.len() > super::streaming::MAX_TOOL_CALLS_PER_RESPONSE {
751 let envelope = super::streaming::tool_call_limit_error();
752 let error = envelope.message.clone();
753 turn.stop_diagnostics.last_response_tool_calls_suppressed =
754 Some(parsed.tool_calls.len());
755 let _ = self.send_stream_event(Event::error(envelope)).await;
756 self.add_interrupted_assistant_text(&current_text_visible)
757 .await;
758 return PhaseResult::Return((TurnOutcomeStatus::Failed, Some(error)));
759 }
760 final_text = parsed.clean_text;
761 for call in parsed.tool_calls {
762 tool_uses.push(ToolUseState {
763 execution_id: self.new_tool_execution_id(),
764 id: call.id,
765 name: call.name,
766 input: call.args,
767 caller: None,
768 thought_signature: None,
769 input_buffer: String::new(),
770 input_parse_error: None,
771 });
772 }
773 if let Err(error) = crate::client::validate_tool_call_ids_for_protocol(
774 protocol,
775 tool_uses.iter().map(|tool| tool.id.as_str()),
776 ) {
777 self.add_interrupted_assistant_text(&current_text_visible)
778 .await;
779 return PhaseResult::Return((TurnOutcomeStatus::Failed, Some(error.to_string())));
780 }
781 }
782
783 // C02-05: the one admission authority for a failed stream. Calls
784 // collected before the failure are announced with an explicit
785 // not-executed result and never reach planning, approval or a
786 // handler; the visible text is kept as an interrupted fragment
787 // and no tool_use enters history, so the transcript stays paired.
788 // A retry never gets here: it `continue`d above with this batch
789 // dropped, and the re-issued request streams its own calls.
790 if response_stream_failed && !tool_uses.is_empty() {
791 let error = progress
792 .turn_error
793 .clone()
794 .unwrap_or_else(|| "provider stream failed".to_string());
795 self.settle_unadmitted_tool_calls(
796 &tool_uses,
797 &stream_failed_tool_result(&summarize_text(&error, 200)),
798 )
799 .await;
800 self.add_interrupted_assistant_text(&current_text_visible)
801 .await;
802 turn.stop_diagnostics.last_response_tool_calls_suppressed = Some(tool_uses.len());
803 return PhaseResult::Return((TurnOutcomeStatus::Failed, Some(error)));
804 }
805
806 for tool in &tool_uses {
807 let _ = self
808 .send_event(Event::ToolCallStarted {
809 id: tool.execution_id.clone(),
810 model_call: Some(tool.model_call()),
811 name: tool.name.clone(),
812 input: final_tool_input(tool),
813 })
814 .await;
815 }
816
817 // Persist only reasoning the provider actually emitted. Some chat
818 // wires require a non-empty `reasoning_content` field when an
819 // assistant message carries tool calls; the route serializer adds
820 // that compatibility value to the outgoing JSON only. Persisting
821 // it here leaked an invented "(reasoning omitted)" block into the
822 // transcript and every provider-neutral session replay.
823 let thinking_is_placeholder_only =
824 crate::client::is_reasoning_replay_placeholder(&current_thinking);
825 if (!current_thinking.is_empty() && !thinking_is_placeholder_only)
826 || current_thinking_state.is_some()
827 {
828 content_blocks.push(ContentBlock::Thinking {
829 thinking: current_thinking.clone(),
830 signature: current_thinking_signature.clone(),
831 state: current_thinking_state.clone(),
832 });
833 }
834
835 // A worker may cooperate with the strategy notice by immediately
836 // reporting its blocker. No intervening useful work means that
837 // report must not become a false Completed result.
838 let fleet_no_progress_report = fleet_report_response
839 || tool_uses.is_empty()
840 && progress
841 .fleet_denial_guard
842 .as_ref()
843 .is_some_and(FleetDenialGuard::awaiting_strategy_change);
844
845 // A protocol-level tool stop promises a call, unlike ordinary
846 // text that merely describes an intended action. Keep that
847 // distinction factual; never synthesize a tool or another request.
848 if tool_uses.is_empty()
849 && !fleet_no_progress_report
850 && progress.turn_error.is_none()
851 && matches!(stop_reason.as_deref(), Some("tool_calls" | "tool_use"))
852 {
853 turn.stop_diagnostics.reason = Some(TurnStopReason::ProviderToolCallMissing);
854 turn.stop_diagnostics.last_response_tool_calls_suppressed = Some(0);
855 self.add_interrupted_assistant_text(&current_text_visible)
856 .await;
857 let reason = stop_reason.as_deref().expect("matched tool stop");
858 return PhaseResult::Return((
859 TurnOutcomeStatus::Failed,
860 Some(
861 codewhale_localization::tr(
862 codewhale_localization::resolve_locale(&self.config.locale_tag),
863 codewhale_localization::MessageId::ProviderToolCallMissing,
864 )
865 .replace("{reason}", reason),
866 ),
867 ));
868 }
869
870 for tool in &mut tool_uses {
871 let Some(schema) = progress
872 .tool_catalog
873 .iter()
874 .find(|candidate| candidate.name == tool.name)
875 .map(|candidate| &candidate.input_schema)
876 else {
877 continue;
878 };
879 normalize_schema_json_containers(&mut tool.input, schema);
880 }
881
882 // A zero-tool turn (plain `exec`) has no tool channel, yet a model
883 // can still answer with nothing but a tool call written as text
884 // (DeepSeek's DSML). The stream filter strips the markup and leaves
885 // at most whitespace, which is not an answer: persisting it would
886 // end the run "successfully" on a blank line, and re-requesting
887 // only reproduces the call. It is failed once below, by name.
888 let zero_tool_text_call = zero_tool_turn
889 && tool_uses.is_empty()
890 && final_text.trim().is_empty()
891 && contains_fake_tool_wrapper(&current_text_raw);
892 if !final_text.is_empty() && !zero_tool_text_call {
893 content_blocks.push(ContentBlock::Text {
894 text: final_text,
895 cache_control: None,
896 });
897 }
898 for tool in &tool_uses {
899 content_blocks.push(ContentBlock::ToolUse {
900 execution_id: Some(tool.execution_id.clone()),
901 id: tool.id.clone(),
902 name: tool.name.clone(),
903 input: tool.input.clone(),
904 caller: tool.caller.clone(),
905 thought_signature: tool.thought_signature.clone(),
906 });
907 }
908
909 if pending_message_complete {
910 let index = last_text_index.unwrap_or(0);
911 let _ = self.send_event(Event::MessageComplete { index }).await;
912 }
913
914 // RLM is a structured tool call (`rlm_query`) handled by the
915 // normal tool dispatch path; inline ```repl blocks (paper §2)
916 // are executed below when tool_uses is empty.
917 // DeepSeek chat API rejects assistant messages that contain only
918 // Keep thinking for UI stream events, but persist only sendable
919 // assistant turns in the conversation state.
920 let has_sendable_assistant_content = content_blocks.iter().any(|block| {
921 matches!(
922 block,
923 ContentBlock::Text { .. } | ContentBlock::ToolUse { .. }
924 )
925 });
926 let has_provider_reasoning = content_blocks.iter().any(|block| {
927 matches!(
928 block,
929 ContentBlock::Thinking {
930 thinking,
931 state,
932 ..
933 } if !thinking.trim().is_empty() || state.is_some()
934 )
935 });
936
937 // Issue #1727: did this turn produce ONLY a reasoning/thinking
938 // block — empty content, no tool calls (e.g. gpt-oss via ollama's
939 // harmony→OpenAI shim mapping to `reasoning_content`)? We do NOT
940 // surface anything here: after this point the same turn can still
941 // CONTINUE for pending steers (~below) or sub-agent completions,
942 // and emitting now would show a spurious "turn ended" notice right
943 // before the turn resumes. Capture the fact and decide later, at
944 // the point the turn is certain to be finishing with no sendable
945 // content (see the `tool_uses.is_empty()` tail).
946 let no_sendable_assistant_content = !has_sendable_assistant_content;
947
948 // Add assistant message to session
949 if has_sendable_assistant_content {
950 self.add_session_message(Message {
951 role: Role::Assistant,
952 content: content_blocks,
953 })
954 .await;
955 }
956
957 // C02-05: a failed stream that carried no tool call keeps its text
958 // (above) and ends the turn with the stream's error. It never
959 // runs a ```repl fence from the failed text, and never authorizes
960 // another request for an output-limit continuation, a steer, a
961 // sub-agent completion or a goal continuation.
962 if response_stream_failed {
963 return PhaseResult::Break;
964 }
965
966 PhaseResult::Ready(AcceptedModelStep {
967 current_text_visible,
968 tool_uses,
969 pending_steers,
970 output_limit_truncated,
971 zero_tool_turn,
972 zero_tool_text_call,
973 fleet_report_response,
974 fleet_no_progress_report,
975 has_sendable_assistant_content,
976 has_provider_reasoning,
977 no_sendable_assistant_content,
978 stop_reason,
979 stream_errors,
980 prepared_output_tokens: stream_request.max_tokens,
981 })
982 }
983 async fn recover_child_request(
984 &mut self,
985 progress: &mut TurnLoopProgress,
986 error: &anyhow::Error,
987 ) -> Option<PhaseResult<AcceptedModelStep>> {
988 use crate::tools::subagent::engine::ChildRequestRecovery;
989 let job = self.child_job()?;
990 if self.child_report() {
991 return None;
992 }
993 let decision = progress
994 .child_request_retries
995 .decide(&job.authority.runtime, error)?;
996 match decision {
997 ChildRequestRecovery::Interrupted {
998 checkpoint_reason,
999 message,
1000 } => {
1001 self.child_host
1002 .as_mut()
1003 .expect("captured child recovery")
1004 .request_stop_reason = Some(checkpoint_reason);
1005 Some(PhaseResult::Return((
1006 TurnOutcomeStatus::Interrupted,
1007 Some(
1008 job.authority
1009 .runtime
1010 .client
1011 .redact_model_bound_text(&message),
1012 ),
1013 )))
1014 }
1015 ChildRequestRecovery::Retry { delay, note } => {
1016 let _ = self
1017 .send_event(Event::status(
1018 job.authority.runtime.client.redact_model_bound_text(&note),
1019 ))
1020 .await;
1021 tokio::select! {
1022 biased;
1023 () = self.cancel_token.cancelled() => Some(PhaseResult::Return((TurnOutcomeStatus::Interrupted, None))),
1024 () = job.authority.runtime.cancel_token.cancelled() => Some(PhaseResult::Return((TurnOutcomeStatus::Interrupted, None))),
1025 () = tokio::time::sleep(delay) => {
1026 progress.turn_error = None;
1027 Some(PhaseResult::Retry)
1028 }
1029 }
1030 }
1031 }
1032 }
1033 }
1034
1034 lines RUST