| 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(¤t_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(¤t_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(¤t_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(¤t_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(¤t_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(¤t_text_raw) { |
| 749 | let parsed = tool_parser::parse_tool_calls(¤t_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(¤t_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(¤t_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(¤t_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(¤t_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(¤t_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(¤t_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(¬e), |
| 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 |