返回 CodeWhale
test_cases_17.rs
根目录 / crates / tui / src / core / engine / tests / test_cases_17.rs
1 #[test]
2 fn stream_read_scenario() {
3 // Scenario consolidation of: stream_read_error_message_explains_retry_before_output, stream_read_error_message_explains_no_replay_after_output
4 // from stream_read_error_message_explains_retry_before_output
5 {
6 let message = super::stream_read_error_user_message(
7 "Stream read error: error decoding response body",
8 false,
9 );
10
11 assert!(message.contains("Provider stream connection dropped"));
12 assert!(message.contains("No output had streamed yet"));
13 assert!(message.contains("retry automatically"));
14 assert!(message.contains("Stream read error: error decoding response body"));
15 }
16 // from stream_read_error_message_explains_no_replay_after_output
17 {
18 let message = super::stream_read_error_user_message(
19 "Stream read error: error decoding response body",
20 true,
21 );
22
23 assert!(message.contains("Provider stream connection dropped"));
24 assert!(message.contains("Some output had already streamed"));
25 assert!(message.contains("risking duplicated output"));
26 assert!(message.contains("Stream read error: error decoding response body"));
27 assert_eq!(
28 crate::error_taxonomy::classify_error_message(&message),
29 crate::error_taxonomy::ErrorCategory::Network
30 );
31 }
32 }
33
34 #[test]
35 fn stream_retry_budget_caps_transparent_retries_at_two() {
36 // Case 4 from issue #103: after MAX_TRANSPARENT_STREAM_RETRIES attempts
37 // we stop trying transparently and let the outer error path surface.
38 // (The outer per-turn `stream_retry_attempts` retry is a separate layer
39 // and is still in effect at the whole-turn level.)
40 assert!(
41 super::should_transparently_retry_stream(
42 false,
43 super::MAX_TRANSPARENT_STREAM_RETRIES - 1,
44 super::MAX_TRANSPARENT_STREAM_RETRIES,
45 false,
46 ),
47 "one short of the cap should still retry"
48 );
49 assert!(
50 !super::should_transparently_retry_stream(
51 false,
52 super::MAX_TRANSPARENT_STREAM_RETRIES,
53 super::MAX_TRANSPARENT_STREAM_RETRIES,
54 false,
55 ),
56 "at the cap, no further transparent retries"
57 );
58 assert!(
59 !super::should_transparently_retry_stream(
60 false,
61 super::MAX_TRANSPARENT_STREAM_RETRIES + 5,
62 super::MAX_TRANSPARENT_STREAM_RETRIES,
63 false,
64 ),
65 "well past the cap, definitely no transparent retries"
66 );
67 }
68
69 // === #2990 sleep-resume policy ================================================
70
71 #[test]
72 fn sleep_gap_requires_wallclock_to_outrun_monotonic_clock() {
73 use std::time::Duration;
74 // No divergence: ordinary network failure, clocks agree.
75 assert!(
76 !super::sleep_gap_detected(Duration::from_secs(30), Duration::from_secs(30)),
77 "equal elapsed times must not register as a sleep gap"
78 );
79 // Divergence below the threshold: NTP slew / scheduling jitter.
80 assert!(
81 !super::sleep_gap_detected(Duration::from_secs(5), Duration::from_secs(14)),
82 "9s of divergence is below the 10s threshold"
83 );
84 // Divergence above the threshold: the host was suspended.
85 assert!(
86 super::sleep_gap_detected(Duration::from_secs(5), Duration::from_secs(16)),
87 "11s of divergence must register as a sleep gap"
88 );
89 // Wall clock went backwards (NTP step): saturating_sub → zero gap.
90 assert!(
91 !super::sleep_gap_detected(Duration::from_secs(60), Duration::from_secs(5)),
92 "wall clock behind monotonic must never register as a sleep gap"
93 );
94 }
95
96 #[test]
97 fn sleep_resume_scenario() {
98 // Scenario consolidation of: sleep_resume_retries_even_after_content_streamed, sleep_resume_requires_a_detected_gap, sleep_resume_respects_budget_and_cancellation
99 // from sleep_resume_retries_even_after_content_streamed
100 {
101 // The whole point of #2990: unlike the #103 transparent retry, a
102 // detected sleep gap retries regardless of streamed content — the
103 // partial output predates the sleep and the user was not watching.
104 assert!(
105 super::should_resume_after_sleep(true, 0, super::MAX_STREAM_RETRIES, false),
106 "detected sleep with full budget must resume"
107 );
108 assert!(
109 super::should_resume_after_sleep(
110 true,
111 super::MAX_STREAM_RETRIES - 1,
112 super::MAX_STREAM_RETRIES,
113 false
114 ),
115 "detected sleep one short of the budget must still resume"
116 );
117 }
118 // from sleep_resume_requires_a_detected_gap
119 {
120 // Without a sleep gap this layer stays out of the way entirely, so the
121 // deliberate no-retry-after-content policy for ordinary flakes (#103)
122 // is preserved.
123 assert!(
124 !super::should_resume_after_sleep(false, 0, super::MAX_STREAM_RETRIES, false),
125 "no sleep gap → never resume via this layer"
126 );
127 }
128 // from sleep_resume_respects_budget_and_cancellation
129 {
130 assert!(
131 !super::should_resume_after_sleep(
132 true,
133 super::MAX_STREAM_RETRIES,
134 super::MAX_STREAM_RETRIES,
135 false
136 ),
137 "budget exhausted → surface the failure instead of looping"
138 );
139 assert!(
140 !super::should_resume_after_sleep(true, 0, super::MAX_STREAM_RETRIES, true),
141 "cancelled turn must not be resumed behind the user's back"
142 );
143 }
144 }
145
146 // === headless mid-stream network-drop resume (v0.9.4 Terminal-Bench P0) ======
147 //
148 // Terminal-Bench 2.1 on the 0.9.4 bundle forfeited tasks when the DeepSeek
149 // stream dropped mid-response ("error decoding response body" after partial
150 // content): the #103 policy surfaced the warning and failed the turn, and
151 // `codewhale exec` exited 1. In a headless host no operator watches the
152 // partial deltas and the fragment is never committed, so the turn loop now
153 // re-issues the request instead (bounded by MAX_STREAM_RETRIES), exactly
154 // like the #2990 sleep-resume.
155
156 #[test]
157 fn network_drop_scenario() {
158 // Scenario consolidation of: network_drop_resume_only_fires_for_headless_hosts, network_drop_resume_requires_network_class_error, network_drop_resume_respects_budget_and_cancellation
159 // from network_drop_resume_only_fires_for_headless_hosts
160 {
161 assert!(
162 super::should_resume_after_network_drop(
163 true,
164 true,
165 0,
166 super::MAX_STREAM_RETRIES,
167 false
168 ),
169 "headless host + network-class drop with budget must resume"
170 );
171 assert!(
172 !super::should_resume_after_network_drop(
173 false,
174 true,
175 0,
176 super::MAX_STREAM_RETRIES,
177 false
178 ),
179 "interactive sessions keep the #103 surface-the-warning policy: \
180 the user saw the partial deltas and replay would render them twice"
181 );
182 }
183 // from network_drop_resume_requires_network_class_error
184 {
185 assert!(
186 !super::should_resume_after_network_drop(
187 true,
188 false,
189 0,
190 super::MAX_STREAM_RETRIES,
191 false
192 ),
193 "non-network failures (model/parse/auth) must never be replayed"
194 );
195 }
196 // from network_drop_resume_respects_budget_and_cancellation
197 {
198 assert!(
199 super::should_resume_after_network_drop(
200 true,
201 true,
202 super::MAX_STREAM_RETRIES - 1,
203 super::MAX_STREAM_RETRIES,
204 false,
205 ),
206 "one short of the budget should still resume"
207 );
208 assert!(
209 !super::should_resume_after_network_drop(
210 true,
211 true,
212 super::MAX_STREAM_RETRIES,
213 super::MAX_STREAM_RETRIES,
214 false
215 ),
216 "budget exhausted → surface the failure instead of looping"
217 );
218 assert!(
219 !super::should_resume_after_network_drop(
220 true,
221 true,
222 0,
223 super::MAX_STREAM_RETRIES,
224 true
225 ),
226 "cancelled turn must not be resumed behind the operator's back"
227 );
228 }
229 }
230
231 // === interactive mid-stream network-drop resume (0.9.4; reworked 0.9.10) =========
232 //
233 // The interactive TUI used to fail the turn when a provider stream dropped
234 // after partial output because the #103 policy treated any post-content error
235 // as terminal. The model now preserves a visible partial reply as a committed
236 // assistant message and re-issues the request. Since 0.9.10 the recovery is
237 // typed engine-internal state (`StreamResume`): no synthetic `[runtime]` user
238 // continuation message is appended, and a thinking-only drop preserves
239 // nothing and never claims it did.
240
241 #[test]
242 fn interactive_network_scenario() {
243 // Scenario consolidation of: interactive_network_drop_resume_only_fires_for_interactive_hosts, interactive_network_drop_resume_requires_partial_content_and_no_tools, interactive_network_drop_resume_requires_network_class_error
244 // from interactive_network_drop_resume_only_fires_for_interactive_hosts
245 {
246 assert!(
247 super::should_resume_interactive_after_network_drop(
248 true,
249 true,
250 true,
251 true,
252 0,
253 super::MAX_STREAM_RETRIES,
254 false
255 ),
256 "interactive TUI + partial text + no tools + budget must resume"
257 );
258 assert!(
259 !super::should_resume_interactive_after_network_drop(
260 false,
261 true,
262 true,
263 true,
264 0,
265 super::MAX_STREAM_RETRIES,
266 false
267 ),
268 "headless hosts must use the headless resume path, not this one"
269 );
270 }
271 // from interactive_network_drop_resume_requires_partial_content_and_no_tools
272 {
273 assert!(
274 !super::should_resume_interactive_after_network_drop(
275 true,
276 true,
277 false,
278 true,
279 0,
280 super::MAX_STREAM_RETRIES,
281 false
282 ),
283 "no streamed content → transparent retry or nothing-streamed path"
284 );
285 assert!(
286 !super::should_resume_interactive_after_network_drop(
287 true,
288 true,
289 true,
290 false,
291 0,
292 super::MAX_STREAM_RETRIES,
293 false
294 ),
295 "in-flight tool calls must never be resumed (side-effect duplication)"
296 );
297 }
298 // from interactive_network_drop_resume_requires_network_class_error
299 {
300 assert!(
301 !super::should_resume_interactive_after_network_drop(
302 true,
303 false,
304 true,
305 true,
306 0,
307 super::MAX_STREAM_RETRIES,
308 false
309 ),
310 "non-network failures must surface normally"
311 );
312 }
313 }
314
315 #[test]
316 fn interactive_network_drop_resume_respects_budget_and_cancellation() {
317 assert!(
318 super::should_resume_interactive_after_network_drop(
319 true,
320 true,
321 true,
322 true,
323 super::MAX_STREAM_RETRIES - 1,
324 super::MAX_STREAM_RETRIES,
325 false,
326 ),
327 "one short of the budget should still resume"
328 );
329 assert!(
330 !super::should_resume_interactive_after_network_drop(
331 true,
332 true,
333 true,
334 true,
335 super::MAX_STREAM_RETRIES,
336 super::MAX_STREAM_RETRIES,
337 false,
338 ),
339 "budget exhausted → surface the failure"
340 );
341 assert!(
342 !super::should_resume_interactive_after_network_drop(
343 true,
344 true,
345 true,
346 true,
347 0,
348 super::MAX_STREAM_RETRIES,
349 true
350 ),
351 "cancelled turn must not resume"
352 );
353 }
354
355 /// Model client whose first `failures` streams emit partial content and then
356 /// die with the network-class read error reqwest reports for a dropped
357 /// chunked-transfer body; later streams complete a normal text turn.
358 /// Exercises the headless mid-stream network-drop resume through the real
359 /// turn loop.
360 struct FlakyNetworkDropModelClient {
361 calls: std::sync::atomic::AtomicUsize,
362 failures: usize,
363 terminal_before_drop: bool,
364 content_before_drop: bool,
365 /// #6699: fail the request itself, before any stream exists, in the
366 /// given shape.
367 open_failure: Option<StreamOpenFailure>,
368 }
369
370 /// #6699/#6711: how a request fails before any stream exists.
371 #[derive(Clone, Copy, Debug)]
372 enum StreamOpenFailure {
373 /// The typed transport error `open_sse_response` returns on a header
374 /// stall.
375 HeaderStall,
376 /// A reqwest connect error under the Anthropic adapter's outer context,
377 /// as an HTTP/1.1-pinned open returns it: the outer message names no
378 /// transport cause.
379 WrappedConnectError,
380 /// The Responses adapter shape: the typed retry-layer network error
381 /// under the adapter's outer context.
382 WrappedTypedNetworkError,
383 /// A provider answered with an HTTP rejection whose body mentions a
384 /// connection reset, as the Anthropic adapter reports it (untyped).
385 HttpRejectionMentioningConnection,
386 }
387
388 impl StreamOpenFailure {
389 async fn error(self) -> anyhow::Error {
390 use anyhow::Context as _;
391 match self {
392 Self::HeaderStall => anyhow::Error::new(crate::llm_client::LlmError::NetworkError(
393 "SSE stream request did not receive response headers after 45s \
394 (HTTP/2 and HTTP/1.1)."
395 .to_string(),
396 )),
397 Self::WrappedConnectError => {
398 // A port that was just released refuses the connection.
399 let listener = std::net::TcpListener::bind("127.0.0.1:0").expect("bind");
400 let addr = listener.local_addr().expect("local addr");
401 drop(listener);
402 let client = crate::tls::reqwest_client_builder()
403 .no_proxy()
404 .build()
405 .expect("client");
406 let send = client
407 .post(format!("http://{addr}/v1/messages"))
408 .send()
409 .await;
410 let err = send
411 .context("Anthropic Messages API request failed")
412 .expect_err("closed port must refuse the connection");
413 assert_eq!(err.to_string(), "Anthropic Messages API request failed");
414 err
415 }
416 Self::WrappedTypedNetworkError => Err::<(), _>(anyhow::Error::new(
417 crate::llm_client::LlmError::NetworkError(
418 "Connection failed: error sending request".to_string(),
419 ),
420 ))
421 .context("Responses API request failed")
422 .expect_err("wrapped network error"),
423 Self::HttpRejectionMentioningConnection => anyhow::anyhow!(
424 "Anthropic API error (HTTP 500 Internal Server Error api_error): \
425 upstream connection reset"
426 ),
427 }
428 }
429 }
430
431 #[async_trait::async_trait]
432 impl crate::core::model_client::ModelClient for FlakyNetworkDropModelClient {
433 fn provider_name(&self) -> &str {
434 "flaky-network"
435 }
436
437 fn model(&self) -> &str {
438 "local-model"
439 }
440
441 async fn create_message(
442 &self,
443 _request: codewhale_models::MessageRequest,
444 ) -> anyhow::Result<codewhale_models::MessageResponse> {
445 anyhow::bail!("flaky-network regression uses the streaming model boundary")
446 }
447
448 async fn create_message_stream(
449 &self,
450 _request: codewhale_models::MessageRequest,
451 ) -> anyhow::Result<crate::llm_client::StreamEventBox> {
452 use crate::llm_client::mock::canned;
453 let call = self
454 .calls
455 .fetch_add(1, std::sync::atomic::Ordering::SeqCst)
456 .saturating_add(1);
457 if call <= self.failures {
458 if let Some(open_failure) = self.open_failure {
459 return Err(open_failure.error().await);
460 }
461 if self.terminal_before_drop {
462 let start_usage = Usage {
463 input_tokens: 31,
464 ..Default::default()
465 };
466 let delta_usage = Usage {
467 output_tokens: 8,
468 ..Default::default()
469 };
470 let mut message_start = canned::message_start("terminal_then_drop");
471 if let StreamEvent::MessageStart { message } = &mut message_start {
472 message.usage = start_usage;
473 }
474 let events: Vec<anyhow::Result<codewhale_models::StreamEvent>> = vec![
475 Ok(message_start),
476 Ok(canned::text_block_start(0)),
477 Ok(canned::text_delta(0, "billed truncated fragment")),
478 Ok(canned::block_stop(0)),
479 Ok(canned::message_delta("max_tokens", Some(delta_usage))),
480 Err(anyhow::anyhow!(
481 "Stream read error: error decoding response body"
482 )),
483 ];
484 return Ok(Box::pin(futures_util::stream::iter(events)));
485 }
486 if !self.content_before_drop {
487 return Ok(Box::pin(futures_util::stream::iter(vec![
488 Ok(canned::message_start("empty_then_drop")),
489 Err(anyhow::anyhow!(
490 "Stream read error: error decoding response body"
491 )),
492 ])));
493 }
494 // Partial content first — this flips `any_content_received` so
495 // the #103 transparent retry cannot fire — then the transport
496 // dies the way the 0.9.4 Terminal-Bench crashes did.
497 let events: Vec<anyhow::Result<codewhale_models::StreamEvent>> = vec![
498 Ok(canned::message_start("flaky_msg")),
499 Ok(canned::text_block_start(0)),
500 Ok(canned::text_delta(
501 0,
502 "partial answer that must be discarded",
503 )),
504 Err(anyhow::anyhow!(
505 "Stream read error: error decoding response body"
506 )),
507 ];
508 return Ok(Box::pin(futures_util::stream::iter(events)));
509 }
510 let events = canned::simple_text_turn("recovered after retry")
511 .into_iter()
512 .map(Ok);
513 Ok(Box::pin(futures_util::stream::iter(events)))
514 }
515
516 async fn health_check(&self) -> anyhow::Result<bool> {
517 Ok(true)
518 }
519 }
520
521 /// Drive one headless (`terminal_chrome_enabled = false`, the exec /
522 /// stream-json posture) turn against the flaky client and collect every
523 /// event through the terminal TurnComplete.
524 async fn run_headless_turn_with_flaky_network(
525 failures: usize,
526 ) -> (std::sync::Arc<FlakyNetworkDropModelClient>, Vec<Event>) {
527 let model = std::sync::Arc::new(FlakyNetworkDropModelClient {
528 calls: std::sync::atomic::AtomicUsize::new(0),
529 failures,
530 terminal_before_drop: false,
531 content_before_drop: true,
532 open_failure: None,
533 });
534 let client: crate::core::model_client::SharedModelClient = model.clone();
535 let config = Config::default();
536 let engine_config = EngineConfig {
537 max_steps: 1,
538 snapshots_enabled: false,
539 subagents_enabled: false,
540 terminal_chrome_enabled: false,
541 ..EngineConfig::default()
542 };
543 let (engine, handle) = Engine::new_with_model_client(engine_config, &config, client);
544 let run_task = tokio::spawn(engine.run());
545
546 handle
547 .send(Op::SendMessage(TurnSpec {
548 max_output_tokens: None,
549 content: "solve the task".to_string(),
550 images: Vec::new(),
551 mode: AppMode::Agent,
552 route: resolved_route_for_test(&config, crate::config::DEFAULT_TEXT_MODEL),
553 compaction: Box::new(CompactionConfig::default()),
554 initial_routed_usage: Box::default(),
555 goal_objective: None,
556 goal_token_budget: None,
557 goal_status: crate::tools::goal::GoalStatus::Active,
558 reasoning_effort: None,
559 reasoning_effort_auto: false,
560 auto_model: false,
561 allow_shell: false,
562 trust_mode: false,
563 auto_approve: false,
564 approval_mode: ApprovalMode::Suggest,
565 translation_enabled: false,
566 allowed_tools: None,
567 dynamic_tools: Vec::new(),
568 hook_executor: None,
569 verbosity: None,
570 provenance: UserInputProvenance::ExternalUser,
571 submission_id: None,
572 }))
573 .await
574 .expect("send flaky-network turn");
575
576 let mut events = Vec::new();
577 loop {
578 let event = tokio::time::timeout(model_turn_event_timeout(), async {
579 handle.rx_event.write().await.recv().await
580 })
581 .await
582 .expect("flaky-network event timeout")
583 .expect("flaky-network event");
584 let terminal = matches!(event, Event::TurnComplete { .. });
585 events.push(event);
586 if terminal {
587 break;
588 }
589 }
590 handle.send(Op::Shutdown).await.expect("shutdown engine");
591 run_task.await.expect("engine task");
592 (model, events)
593 }
594
595 #[tokio::test]
596 async fn headless_turn_retries_mid_stream_network_drop_and_recovers() {
597 let (model, events) = run_headless_turn_with_flaky_network(1).await;
598
599 let terminal = events
600 .iter()
601 .find_map(|event| match event {
602 Event::ToolRequestSnapshot { snapshot } => snapshot.terminal.as_ref(),
603 _ => None,
604 })
605 .expect("terminal diagnostics through existing event authority");
606 assert_eq!(terminal.model_requests_started, 2);
607 assert_eq!(terminal.stream_resumes, 1);
608 assert_eq!(terminal.transparent_stream_retries, 0);
609
610 assert_eq!(
611 model.calls.load(std::sync::atomic::Ordering::SeqCst),
612 2,
613 "the dropped stream must be re-issued exactly once"
614 );
615 let (status, error) = events
616 .iter()
617 .find_map(|event| match event {
618 Event::TurnComplete { status, error, .. } => Some((status, error)),
619 _ => None,
620 })
621 .expect("terminal TurnComplete");
622 assert_eq!(
623 *status,
624 TurnOutcomeStatus::Completed,
625 "a recovered retry must complete the turn: {error:?}"
626 );
627 assert!(error.is_none(), "recovered turn must not report an error");
628 assert!(
629 events.iter().any(|event| matches!(
630 event,
631 Event::Status { message } if message.starts_with("Retry attempt: stream-resume 1/")
632 )),
633 "the first actual resume must remain visible: {events:?}"
634 );
635 assert!(
636 !events
637 .iter()
638 .any(|event| matches!(event, Event::Error { .. })),
639 "a transient drop that the retry recovers must not surface an error event: {events:?}"
640 );
641 // The discarded fragment from the dropped attempt must never reach the
642 // transcript — only the retried turn's content is committed.
643 let transcript_text = events
644 .iter()
645 .filter_map(|event| match event {
646 Event::SessionUpdated { messages, .. } => Some(messages),
647 _ => None,
648 })
649 .flat_map(|messages| messages.iter())
650 .flat_map(|message| message.content.iter())
651 .filter_map(|block| match block {
652 ContentBlock::Text { text, .. } => Some(text.as_str()),
653 _ => None,
654 })
655 .collect::<Vec<_>>()
656 .join("\n");
657 assert!(
658 transcript_text.contains("recovered after retry"),
659 "retried content must be committed: {transcript_text}"
660 );
661 assert!(
662 !transcript_text.contains("partial answer that must be discarded"),
663 "the dropped attempt's fragment must be discarded, not committed: {transcript_text}"
664 );
665 }
666
667 #[tokio::test]
668 async fn terminal_diagnostics_count_transparent_stream_requests_without_extra_snapshots() {
669 let workspace = tempdir().expect("tempdir");
670 let model = std::sync::Arc::new(FlakyNetworkDropModelClient {
671 calls: std::sync::atomic::AtomicUsize::new(0),
672 failures: 1,
673 terminal_before_drop: false,
674 content_before_drop: false,
675 open_failure: None,
676 });
677 let client: crate::core::model_client::SharedModelClient = model.clone();
678 let (mut engine, handle) = Engine::new_with_model_client(
679 deterministic_engine_config(workspace.path()),
680 &Config::default(),
681 client,
682 );
683 let registry = crate::tools::ToolRegistry::new(crate::tools::ToolContext::new(
684 workspace.path().to_path_buf(),
685 ));
686 let surface = test_tool_surface(&engine, registry, None, AppMode::Agent);
687 let mut turn = crate::core::turn::TurnContext::new(4);
688 let (status, error) = engine.run_turn(&mut turn, surface, None, None).await;
689 assert_eq!(status, TurnOutcomeStatus::Completed, "{error:?}");
690 assert_eq!(model.calls.load(std::sync::atomic::Ordering::SeqCst), 2);
691 let snapshot = turn
692 .terminal_request_snapshot(status)
693 .expect("terminal snapshot");
694 let terminal = snapshot.terminal.expect("terminal facts");
695 assert_eq!(terminal.model_requests_started, 2);
696 assert_eq!(terminal.transparent_stream_retries, 1);
697 assert_eq!(terminal.stream_resumes, 0);
698 assert_eq!(terminal.model_step_index, 0);
699 let mut receiver = handle.rx_event.write().await;
700 let events = std::iter::from_fn(|| receiver.try_recv().ok()).collect::<Vec<_>>();
701 assert!(events.iter().any(|event| matches!(event,
702 Event::Status { message } if message.starts_with("Retry attempt: transparent-stream 1/")
703 )));
704 assert!(events.iter().any(|event| matches!(event,
705 Event::Status { message } if message == "Retry recovery: transparent stream recovered after 1 retries"
706 )));
707 let snapshots = events
708 .iter()
709 .filter(|event| matches!(event, Event::ToolRequestSnapshot { .. }))
710 .count();
711 assert_eq!(
712 snapshots, 1,
713 "request construction is distinct from stream retries"
714 );
715 }
716
717 /// #6699: drive one interactive turn whose first `failures` requests fail
718 /// before a stream opens, with `max_resumes` as the configured budget.
719 async fn run_turn_with_stream_open_failures(
720 failures: usize,
721 max_resumes: Option<u32>,
722 open_failure: StreamOpenFailure,
723 ) -> (
724 std::sync::Arc<FlakyNetworkDropModelClient>,
725 TurnOutcomeStatus,
726 Option<String>,
727 crate::tool_inspection::TurnStopDiagnostics,
728 Vec<Event>,
729 ) {
730 let workspace = tempdir().expect("tempdir");
731 let model = std::sync::Arc::new(FlakyNetworkDropModelClient {
732 calls: std::sync::atomic::AtomicUsize::new(0),
733 failures,
734 terminal_before_drop: false,
735 content_before_drop: false,
736 open_failure: Some(open_failure),
737 });
738 let client: crate::core::model_client::SharedModelClient = model.clone();
739 let engine_config = EngineConfig {
740 terminal_chrome_enabled: true,
741 stream_retry_limits: turn_budget::resolve_stream_retry_limits(max_resumes, None, None),
742 ..deterministic_engine_config(workspace.path())
743 };
744 let (mut engine, handle) =
745 Engine::new_with_model_client(engine_config, &Config::default(), client);
746 let registry = crate::tools::ToolRegistry::new(crate::tools::ToolContext::new(
747 workspace.path().to_path_buf(),
748 ));
749 let surface = test_tool_surface(&engine, registry, None, AppMode::Agent);
750 let mut turn = crate::core::turn::TurnContext::new(4);
751 let (status, error) = engine.run_turn(&mut turn, surface, None, None).await;
752 let terminal = turn
753 .terminal_request_snapshot(status)
754 .expect("terminal snapshot")
755 .terminal
756 .expect("terminal facts");
757 let mut rx = handle.rx_event.write().await;
758 let events = std::iter::from_fn(|| rx.try_recv().ok()).collect();
759 (model, status, error, terminal, events)
760 }
761
762 #[tokio::test]
763 async fn stream_open_failure_is_retried_through_the_resume_budget() {
764 let (model, status, error, terminal, events) =
765 run_turn_with_stream_open_failures(1, None, StreamOpenFailure::HeaderStall).await;
766 assert_eq!(status, TurnOutcomeStatus::Completed, "{error:?}");
767 assert_eq!(
768 model.calls.load(std::sync::atomic::Ordering::SeqCst),
769 2,
770 "a request that never opened must be re-issued once"
771 );
772 assert_eq!(terminal.model_requests_started, 2);
773 assert_eq!(terminal.stream_resumes, 1);
774 assert!(
775 !events
776 .iter()
777 .any(|event| matches!(event, Event::Error { .. })),
778 "a recovered open failure must not surface an error event: {events:?}"
779 );
780 }
781
782 #[tokio::test]
783 async fn stream_open_failure_fails_the_turn_once_the_budget_is_spent() {
784 let (model, status, error, terminal, events) =
785 run_turn_with_stream_open_failures(usize::MAX, None, StreamOpenFailure::HeaderStall).await;
786 assert_eq!(status, TurnOutcomeStatus::Failed);
787 assert_eq!(
788 model.calls.load(std::sync::atomic::Ordering::SeqCst),
789 1 + super::MAX_STREAM_RETRIES as usize,
790 "initial attempt plus the default resume budget, then the turn fails"
791 );
792 assert_eq!(terminal.stream_resumes, super::MAX_STREAM_RETRIES);
793 assert!(
794 error
795 .as_deref()
796 .is_some_and(|error| error.contains("did not receive response headers")),
797 "the real transport error must surface: {error:?}"
798 );
799 assert_eq!(
800 events
801 .iter()
802 .filter(|event| matches!(event, Event::Error { .. }))
803 .count(),
804 1,
805 "only the exhausted attempt emits an error event: {events:?}"
806 );
807 }
808
809 #[tokio::test]
810 async fn stream_open_failure_honors_a_configured_resume_budget() {
811 let (model, status, _error, terminal, _events) =
812 run_turn_with_stream_open_failures(usize::MAX, Some(0), StreamOpenFailure::HeaderStall)
813 .await;
814 assert_eq!(status, TurnOutcomeStatus::Failed);
815 assert_eq!(
816 model.calls.load(std::sync::atomic::Ordering::SeqCst),
817 1,
818 "`stream_max_resumes = 0` disables the turn-level retry"
819 );
820 assert_eq!(terminal.stream_resumes, 0);
821 }
822
823 /// #6711: an HTTP/1.1-pinned Anthropic open returns the reqwest connect
824 /// error under an outer context that names no transport cause. The retry gate
825 /// must read the error chain, not the outer message.
826 #[tokio::test]
827 async fn stream_open_failure_behind_adapter_context_is_retried() {
828 for open_failure in [
829 StreamOpenFailure::WrappedConnectError,
830 StreamOpenFailure::WrappedTypedNetworkError,
831 ] {
832 let (model, status, error, terminal, _events) =
833 run_turn_with_stream_open_failures(1, None, open_failure).await;
834 assert_eq!(
835 status,
836 TurnOutcomeStatus::Completed,
837 "{open_failure:?}: {error:?}"
838 );
839 assert_eq!(
840 model.calls.load(std::sync::atomic::Ordering::SeqCst),
841 2,
842 "{open_failure:?}: a wrapped transport failure must be re-issued"
843 );
844 assert_eq!(terminal.stream_resumes, 1, "{open_failure:?}");
845 }
846 }
847
848 /// #6711: a provider that answered with an HTTP rejection is not an open
849 /// failure, even when its body mentions a connection or a timeout. The request
850 /// must not be re-issued through the resume budget.
851 #[tokio::test]
852 async fn stream_open_http_rejection_mentioning_connection_is_not_retried() {
853 let (model, status, _error, terminal, events) = run_turn_with_stream_open_failures(
854 usize::MAX,
855 None,
856 StreamOpenFailure::HttpRejectionMentioningConnection,
857 )
858 .await;
859 assert_eq!(status, TurnOutcomeStatus::Failed);
860 assert_eq!(
861 model.calls.load(std::sync::atomic::Ordering::SeqCst),
862 1,
863 "a provider rejection must fail the turn on the first answer"
864 );
865 assert_eq!(terminal.stream_resumes, 0);
866 assert_eq!(
867 events
868 .iter()
869 .filter(|event| matches!(event, Event::Error { .. }))
870 .count(),
871 1,
872 "{events:?}"
873 );
874 }
875
876 #[tokio::test]
877 async fn terminal_output_limit_followed_by_stream_error_is_charged_and_not_retried() {
878 let model = std::sync::Arc::new(FlakyNetworkDropModelClient {
879 calls: std::sync::atomic::AtomicUsize::new(0),
880 failures: 1,
881 terminal_before_drop: true,
882 content_before_drop: true,
883 open_failure: None,
884 });
885 let client: crate::core::model_client::SharedModelClient = model.clone();
886 let config = Config::default();
887 let engine_config = EngineConfig {
888 max_steps: 1,
889 snapshots_enabled: false,
890 subagents_enabled: false,
891 terminal_chrome_enabled: false,
892 ..EngineConfig::default()
893 };
894 let (engine, handle) = Engine::new_with_model_client(engine_config, &config, client);
895 let run_task = tokio::spawn(engine.run());
896
897 handle
898 .send(Op::SendMessage(TurnSpec {
899 max_output_tokens: None,
900 content: "solve the task".to_string(),
901 images: Vec::new(),
902 mode: AppMode::Agent,
903 route: resolved_route_for_test(&config, crate::config::DEFAULT_TEXT_MODEL),
904 compaction: Box::new(CompactionConfig::default()),
905 initial_routed_usage: Box::default(),
906 goal_objective: None,
907 goal_token_budget: None,
908 goal_status: crate::tools::goal::GoalStatus::Active,
909 reasoning_effort: None,
910 reasoning_effort_auto: false,
911 auto_model: false,
912 allow_shell: false,
913 trust_mode: false,
914 auto_approve: false,
915 approval_mode: ApprovalMode::Suggest,
916 translation_enabled: false,
917 allowed_tools: None,
918 dynamic_tools: Vec::new(),
919 hook_executor: None,
920 verbosity: None,
921 provenance: UserInputProvenance::ExternalUser,
922 submission_id: None,
923 }))
924 .await
925 .expect("send terminal-then-drop turn");
926
927 let mut events = Vec::new();
928 loop {
929 let event = tokio::time::timeout(model_turn_event_timeout(), async {
930 handle.rx_event.write().await.recv().await
931 })
932 .await
933 .expect("terminal-then-drop event timeout")
934 .expect("terminal-then-drop event");
935 let terminal = matches!(event, Event::TurnComplete { .. });
936 events.push(event);
937 if terminal {
938 break;
939 }
940 }
941
942 assert_eq!(
943 model.calls.load(std::sync::atomic::Ordering::SeqCst),
944 1,
945 "a provider-declared terminal response must never be re-issued"
946 );
947 assert!(events.iter().any(|event| matches!(
948 event,
949 Event::TurnUsage { usage, .. }
950 if usage.input_tokens == 31 && usage.output_tokens == 8
951 )));
952 let (status, error) = events
953 .iter()
954 .find_map(|event| match event {
955 Event::TurnComplete { status, error, .. } => Some((status, error)),
956 _ => None,
957 })
958 .expect("terminal TurnComplete");
959 assert_eq!(*status, TurnOutcomeStatus::Failed);
960 assert!(
961 error
962 .as_deref()
963 .is_some_and(|error| error.contains("max_tokens")),
964 "{error:?}"
965 );
966 assert!(!events.iter().any(|event| matches!(
967 event,
968 Event::Status { message } if crate::core::events::is_retry_status_receipt(message)
969 )));
970
971 handle.send(Op::Shutdown).await.expect("shutdown engine");
972 run_task.await.expect("engine task");
973 }
974
975 /// Run one user turn against scripted streams; returns its events and how
976 /// many model requests were made.
977 async fn error_frame_turn_events(turns: Vec<Vec<StreamEvent>>) -> (Vec<Event>, usize) {
978 let model = std::sync::Arc::new(crate::llm_client::mock::MockLlmClient::new(turns));
979 let client: crate::core::model_client::SharedModelClient = model.clone();
980 let config = Config::default();
981 let engine_config = EngineConfig {
982 max_steps: 1,
983 snapshots_enabled: false,
984 subagents_enabled: false,
985 terminal_chrome_enabled: false,
986 ..EngineConfig::default()
987 };
988 let (engine, handle) = Engine::new_with_model_client(engine_config, &config, client);
989 let run_task = tokio::spawn(engine.run());
990 handle
991 .send(Op::SendMessage(TurnSpec {
992 max_output_tokens: None,
993 content: "solve the task".to_string(),
994 images: Vec::new(),
995 mode: AppMode::Agent,
996 route: resolved_route_for_test(&config, crate::config::DEFAULT_TEXT_MODEL),
997 compaction: Box::new(CompactionConfig::default()),
998 initial_routed_usage: Box::default(),
999 goal_objective: None,
1000 goal_token_budget: None,
1001 goal_status: crate::tools::goal::GoalStatus::Active,
1002 reasoning_effort: None,
1003 reasoning_effort_auto: false,
1004 auto_model: false,
1005 allow_shell: false,
1006 trust_mode: false,
1007 auto_approve: false,
1008 approval_mode: ApprovalMode::Suggest,
1009 translation_enabled: false,
1010 allowed_tools: None,
1011 dynamic_tools: Vec::new(),
1012 hook_executor: None,
1013 verbosity: None,
1014 provenance: UserInputProvenance::ExternalUser,
1015 submission_id: None,
1016 }))
1017 .await
1018 .expect("send turn");
1019 let mut events = Vec::new();
1020 loop {
1021 let event = tokio::time::timeout(model_turn_event_timeout(), async {
1022 handle.rx_event.write().await.recv().await
1023 })
1024 .await
1025 .expect("turn event timeout")
1026 .expect("turn event");
1027 let terminal = matches!(event, Event::TurnComplete { .. });
1028 events.push(event);
1029 if terminal {
1030 break;
1031 }
1032 }
1033 handle.send(Op::Shutdown).await.expect("shutdown engine");
1034 run_task.await.expect("engine task");
1035 (events, model.call_count())
1036 }
1037
1038 /// #6795: a transient upstream failure delivered as an error frame inside a
1039 /// successful response, before anything streamed, is retried like every other
1040 /// no-content stream death. A terminal-class frame still fails on the first
1041 /// request.
1042 #[tokio::test]
1043 async fn transient_error_frame_with_no_content_is_retried_and_terminal_frame_is_not() {
1044 use crate::llm_client::mock::canned;
1045 let frame = |message: &str| {
1046 vec![StreamEvent::Error {
1047 error: serde_json::json!({ "message": message }),
1048 }]
1049 };
1050
1051 let (events, requests) = error_frame_turn_events(vec![
1052 frame("Provider returned an empty response"),
1053 canned::simple_text_turn("recovered answer"),
1054 ])
1055 .await;
1056 assert_eq!(requests, 2, "the empty-upstream frame is re-issued once");
1057 assert!(events.iter().any(|event| matches!(event,
1058 Event::Status { message } if message.starts_with("Retry attempt: stream-resume 1/")
1059 )));
1060 assert!(events.iter().any(|event| matches!(event,
1061 Event::Status { message } if message == "Retry recovery: stream recovered after 1 retries"
1062 )));
1063 assert!(
1064 events.iter().any(|event| matches!(
1065 event,
1066 Event::TurnComplete {
1067 status: TurnOutcomeStatus::Completed,
1068 error: None,
1069 ..
1070 }
1071 )),
1072 "a successful retry completes the turn"
1073 );
1074 assert!(
1075 !events
1076 .iter()
1077 .any(|event| matches!(event, Event::Error { .. })),
1078 "a retried frame leaves no error card behind"
1079 );
1080
1081 let (events, requests) = error_frame_turn_events(vec![
1082 frame("Model not exist."),
1083 canned::simple_text_turn("must never be requested"),
1084 ])
1085 .await;
1086 assert_eq!(requests, 1, "a terminal-class frame is never retried");
1087 assert!(events.iter().any(|event| matches!(
1088 event,
1089 Event::TurnComplete {
1090 status: TurnOutcomeStatus::Failed,
1091 ..
1092 }
1093 )));
1094 assert!(events.iter().any(|event| matches!(
1095 event,
1096 Event::Error { envelope, .. }
1097 if envelope.message.contains("Model not exist.") && !envelope.recoverable
1098 )));
1099 }
1100
1101 #[tokio::test]
1102 async fn midstream_error_frame_stops_the_stream_and_drops_trailing_deltas() {
1103 // The reported incident: a provider delivered a chunk-level error
1104 // object ("Model not exist.") mid-stream, and the stream kept parsing
1105 // later frames as deltas — so reasoning/text rendered *after* the
1106 // failure. The `StreamEvent::Error` arm must surface a terminal error
1107 // event, stop consuming the stream, and never forward deltas that
1108 // arrive after the failure frame.
1109 let model = std::sync::Arc::new(crate::llm_client::mock::MockLlmClient::new(vec![vec![
1110 crate::llm_client::mock::canned::message_start("midstream-error"),
1111 crate::llm_client::mock::canned::text_block_start(0),
1112 StreamEvent::Error {
1113 error: serde_json::json!({ "message": "Model not exist." }),
1114 },
1115 // Everything after the error frame must never reach the UI.
1116 crate::llm_client::mock::canned::text_delta(0, "TRAILING-DELTA-AFTER-FAILURE"),
1117 crate::llm_client::mock::canned::block_stop(0),
1118 ]]));
1119 let client: crate::core::model_client::SharedModelClient = model.clone();
1120 let config = Config::default();
1121 let engine_config = EngineConfig {
1122 max_steps: 1,
1123 snapshots_enabled: false,
1124 subagents_enabled: false,
1125 terminal_chrome_enabled: false,
1126 ..EngineConfig::default()
1127 };
1128 let (engine, handle) = Engine::new_with_model_client(engine_config, &config, client);
1129 let run_task = tokio::spawn(engine.run());
1130
1131 handle
1132 .send(Op::SendMessage(TurnSpec {
1133 max_output_tokens: None,
1134 content: "solve the task".to_string(),
1135 images: Vec::new(),
1136 mode: AppMode::Agent,
1137 route: resolved_route_for_test(&config, crate::config::DEFAULT_TEXT_MODEL),
1138 compaction: Box::new(CompactionConfig::default()),
1139 initial_routed_usage: Box::default(),
1140 goal_objective: None,
1141 goal_token_budget: None,
1142 goal_status: crate::tools::goal::GoalStatus::Active,
1143 reasoning_effort: None,
1144 reasoning_effort_auto: false,
1145 auto_model: false,
1146 allow_shell: false,
1147 trust_mode: false,
1148 auto_approve: false,
1149 approval_mode: ApprovalMode::Suggest,
1150 translation_enabled: false,
1151 allowed_tools: None,
1152 dynamic_tools: Vec::new(),
1153 hook_executor: None,
1154 verbosity: None,
1155 provenance: UserInputProvenance::ExternalUser,
1156 submission_id: None,
1157 }))
1158 .await
1159 .expect("send midstream-error turn");
1160
1161 let mut events = Vec::new();
1162 loop {
1163 let event = tokio::time::timeout(model_turn_event_timeout(), async {
1164 handle.rx_event.write().await.recv().await
1165 })
1166 .await
1167 .expect("midstream-error event timeout")
1168 .expect("midstream-error event");
1169 let terminal = matches!(event, Event::TurnComplete { .. });
1170 events.push(event);
1171 if terminal {
1172 break;
1173 }
1174 }
1175
1176 // The failure is surfaced as a typed error event at Error severity.
1177 assert!(events.iter().any(|event| matches!(
1178 event,
1179 Event::Error { envelope, .. }
1180 if envelope.message.contains("Model not exist.")
1181 && envelope.category == crate::error_taxonomy::ErrorCategory::InvalidInput
1182 && envelope.severity == crate::error_taxonomy::ErrorSeverity::Error
1183 )));
1184 // Deltas after the failure frame never reach the UI.
1185 assert!(
1186 !events.iter().any(|event| matches!(
1187 event,
1188 Event::MessageDelta { content, .. } if content.contains("TRAILING-DELTA")
1189 )),
1190 "deltas after a mid-stream failure must be dropped: {events:?}"
1191 );
1192 // The turn fails with the provider's message, not a generic truncation.
1193 let (status, error) = events
1194 .iter()
1195 .find_map(|event| match event {
1196 Event::TurnComplete { status, error, .. } => Some((status, error)),
1197 _ => None,
1198 })
1199 .expect("midstream-error TurnComplete");
1200 assert_eq!(*status, TurnOutcomeStatus::Failed);
1201 assert!(
1202 error
1203 .as_deref()
1204 .is_some_and(|e| e.contains("Model not exist.")),
1205 "{error:?}"
1206 );
1207 // The stream is consumed exactly once — no re-issue of a terminal
1208 // rejection.
1209 assert_eq!(model.call_count(), 1);
1210
1211 handle.send(Op::Shutdown).await.expect("shutdown engine");
1212 run_task.await.expect("engine task");
1213 }
1214
1215 /// Emits a full billed response, then cancels the engine's own token while
1216 /// yielding the final stream event — modeling Esc arriving right after the
1217 /// provider finished charging for the response.
1218 struct CancelAfterTerminalUsageModelClient {
1219 calls: std::sync::atomic::AtomicUsize,
1220 // The engine mints a fresh token per turn; read the live one through the
1221 // engine's shared cell at stream time.
1222 token: std::sync::Mutex<Option<Arc<StdMutex<tokio_util::sync::CancellationToken>>>>,
1223 }
1224
1225 #[async_trait::async_trait]
1226 impl crate::core::model_client::ModelClient for CancelAfterTerminalUsageModelClient {
1227 fn provider_name(&self) -> &str {
1228 "mock"
1229 }
1230
1231 fn model(&self) -> &str {
1232 "mock-model"
1233 }
1234
1235 async fn create_message(
1236 &self,
1237 _request: codewhale_models::MessageRequest,
1238 ) -> anyhow::Result<codewhale_models::MessageResponse> {
1239 anyhow::bail!("unused")
1240 }
1241
1242 async fn create_message_stream(
1243 &self,
1244 _request: codewhale_models::MessageRequest,
1245 ) -> anyhow::Result<crate::llm_client::StreamEventBox> {
1246 use crate::llm_client::mock::canned;
1247
1248 self.calls.fetch_add(1, std::sync::atomic::Ordering::SeqCst);
1249 let shared = self
1250 .token
1251 .lock()
1252 .expect("token cell")
1253 .clone()
1254 .expect("token installed before turn");
1255 let token = shared
1256 .lock()
1257 .unwrap_or_else(std::sync::PoisonError::into_inner)
1258 .clone();
1259 let mut message_start = canned::message_start("cancel_after_usage");
1260 if let StreamEvent::MessageStart { message } = &mut message_start {
1261 message.usage = Usage {
1262 input_tokens: 47,
1263 ..Default::default()
1264 };
1265 }
1266 let events = vec![
1267 message_start,
1268 canned::text_block_start(0),
1269 canned::text_delta(0, "answer the user was billed for"),
1270 canned::block_stop(0),
1271 canned::message_delta(
1272 "end_turn",
1273 Some(Usage {
1274 output_tokens: 9,
1275 ..Default::default()
1276 }),
1277 ),
1278 canned::message_stop(),
1279 ];
1280 let last = events.len() - 1;
1281 let stream = futures_util::stream::iter(events.into_iter().enumerate().map(
1282 move |(index, event)| {
1283 if index == last {
1284 token.cancel();
1285 }
1286 Ok(event)
1287 },
1288 ));
1289 Ok(Box::pin(stream))
1290 }
1291
1292 async fn health_check(&self) -> anyhow::Result<bool> {
1293 Ok(true)
1294 }
1295 }
1296
1297 #[tokio::test]
1298 async fn cancellation_after_terminal_usage_still_charges_the_turn() {
1299 let model = std::sync::Arc::new(CancelAfterTerminalUsageModelClient {
1300 calls: std::sync::atomic::AtomicUsize::new(0),
1301 token: std::sync::Mutex::new(None),
1302 });
1303 let client: crate::core::model_client::SharedModelClient = model.clone();
1304 let config = Config::default();
1305 let workspace = tempdir().expect("tempdir");
1306 let (engine, handle) = Engine::new_with_model_client(
1307 deterministic_engine_config(workspace.path()),
1308 &config,
1309 client,
1310 );
1311 *model.token.lock().expect("token cell") = Some(engine.shared_cancel_token.clone());
1312 let run_task = tokio::spawn(engine.run());
1313
1314 handle
1315 .send(external_user_message_op(
1316 "solve the task",
1317 AppMode::Agent,
1318 &config,
1319 ))
1320 .await
1321 .expect("send cancel-after-usage turn");
1322
1323 let mut events = Vec::new();
1324 loop {
1325 let event = tokio::time::timeout(model_turn_event_timeout(), async {
1326 handle.rx_event.write().await.recv().await
1327 })
1328 .await
1329 .expect("cancel-after-usage event timeout")
1330 .expect("cancel-after-usage event");
1331 let terminal = matches!(event, Event::TurnComplete { .. });
1332 events.push(event);
1333 if terminal {
1334 break;
1335 }
1336 }
1337
1338 assert_eq!(model.calls.load(std::sync::atomic::Ordering::SeqCst), 1);
1339 assert!(
1340 events.iter().any(|event| matches!(
1341 event,
1342 Event::TurnUsage { usage, .. }
1343 if usage.input_tokens == 47 && usage.output_tokens == 9
1344 )),
1345 "billed usage must be accounted even though the turn was cancelled"
1346 );
1347 let status = events
1348 .iter()
1349 .find_map(|event| match event {
1350 Event::TurnComplete { status, .. } => Some(*status),
1351 _ => None,
1352 })
1353 .expect("terminal TurnComplete");
1354 assert_eq!(status, TurnOutcomeStatus::Interrupted);
1355
1356 handle.send(Op::Shutdown).await.expect("shutdown engine");
1357 run_task.await.expect("engine task");
1358 }
1359
1360 /// Drive one interactive (`terminal_chrome_enabled = true`) turn against the
1361 /// flaky client and collect every event through the terminal TurnComplete.
1362 async fn run_interactive_turn_with_flaky_network(
1363 failures: usize,
1364 ) -> (std::sync::Arc<FlakyNetworkDropModelClient>, Vec<Event>) {
1365 let model = std::sync::Arc::new(FlakyNetworkDropModelClient {
1366 calls: std::sync::atomic::AtomicUsize::new(0),
1367 failures,
1368 terminal_before_drop: false,
1369 content_before_drop: true,
1370 open_failure: None,
1371 });
1372 let client: crate::core::model_client::SharedModelClient = model.clone();
1373 let config = Config::default();
1374 let engine_config = EngineConfig {
1375 max_steps: 1,
1376 snapshots_enabled: false,
1377 subagents_enabled: false,
1378 terminal_chrome_enabled: true,
1379 ..EngineConfig::default()
1380 };
1381 let (engine, handle) = Engine::new_with_model_client(engine_config, &config, client);
1382 let run_task = tokio::spawn(engine.run());
1383
1384 handle
1385 .send(Op::SendMessage(TurnSpec {
1386 max_output_tokens: None,
1387 content: "solve the task".to_string(),
1388 images: Vec::new(),
1389 mode: AppMode::Agent,
1390 route: resolved_route_for_test(&config, crate::config::DEFAULT_TEXT_MODEL),
1391 compaction: Box::new(CompactionConfig::default()),
1392 initial_routed_usage: Box::default(),
1393 goal_objective: None,
1394 goal_token_budget: None,
1395 goal_status: crate::tools::goal::GoalStatus::Active,
1396 reasoning_effort: None,
1397 reasoning_effort_auto: false,
1398 auto_model: false,
1399 allow_shell: false,
1400 trust_mode: false,
1401 auto_approve: false,
1402 approval_mode: ApprovalMode::Suggest,
1403 translation_enabled: false,
1404 allowed_tools: None,
1405 dynamic_tools: Vec::new(),
1406 hook_executor: None,
1407 verbosity: None,
1408 provenance: UserInputProvenance::ExternalUser,
1409 submission_id: None,
1410 }))
1411 .await
1412 .expect("send interactive flaky-network turn");
1413
1414 let mut events = Vec::new();
1415 loop {
1416 let event = tokio::time::timeout(model_turn_event_timeout(), async {
1417 handle.rx_event.write().await.recv().await
1418 })
1419 .await
1420 .expect("interactive flaky-network event timeout")
1421 .expect("interactive flaky-network event");
1422 let terminal = matches!(event, Event::TurnComplete { .. });
1423 events.push(event);
1424 if terminal {
1425 break;
1426 }
1427 }
1428 handle.send(Op::Shutdown).await.expect("shutdown engine");
1429 run_task.await.expect("engine task");
1430 (model, events)
1431 }
1432
1433 #[tokio::test]
1434 async fn interactive_turn_preserves_partial_reply_and_recovers_after_network_drop() {
1435 let (model, events) = run_interactive_turn_with_flaky_network(1).await;
1436
1437 assert_eq!(
1438 model.calls.load(std::sync::atomic::Ordering::SeqCst),
1439 2,
1440 "the dropped stream must be re-issued exactly once"
1441 );
1442 let (status, error) = events
1443 .iter()
1444 .find_map(|event| match event {
1445 Event::TurnComplete { status, error, .. } => Some((status, error)),
1446 _ => None,
1447 })
1448 .expect("terminal TurnComplete");
1449 assert_eq!(
1450 *status,
1451 TurnOutcomeStatus::Completed,
1452 "a recovered retry must complete the turn: {error:?}"
1453 );
1454 assert!(error.is_none(), "recovered turn must not report an error");
1455 assert!(
1456 events.iter().any(|event| matches!(
1457 event,
1458 Event::Status { message } if message.starts_with("Retry attempt: stream-resume 1/")
1459 )),
1460 "the first actual resume must remain visible: {events:?}"
1461 );
1462 assert!(
1463 !events
1464 .iter()
1465 .any(|event| matches!(event, Event::Error { .. })),
1466 "a transient drop that the retry recovers must not surface an error event: {events:?}"
1467 );
1468
1469 // The visible fragment must survive as an assistant message, followed by
1470 // the retried assistant content — and nothing else. The recovery is
1471 // typed internal state, so no synthetic `[runtime]` user turn may appear.
1472 let transcript = events
1473 .iter()
1474 .rev()
1475 .find_map(|event| match event {
1476 Event::SessionUpdated { messages, .. } => Some(messages.clone()),
1477 _ => None,
1478 })
1479 .expect("final SessionUpdated");
1480 assert!(
1481 transcript
1482 .iter()
1483 .flat_map(|message| message.content.iter())
1484 .all(|block| match block {
1485 ContentBlock::Text { text, .. } => !text.contains("[runtime]"),
1486 _ => true,
1487 }),
1488 "a retried turn must not insert a synthetic [runtime] user message: {transcript:?}"
1489 );
1490 assert_eq!(
1491 transcript
1492 .iter()
1493 .filter(|message| message.role == "user")
1494 .count(),
1495 1,
1496 "the operator's own turn must be the only user message: {transcript:?}"
1497 );
1498 let assistant_cells = transcript
1499 .iter()
1500 .filter(|message| message.role == "assistant")
1501 .count();
1502 assert_eq!(
1503 assistant_cells, 2,
1504 "preserved fragment + one authoritative continuation: {transcript:?}"
1505 );
1506 let transcript_text = transcript
1507 .iter()
1508 .flat_map(|message| message.content.iter())
1509 .filter_map(|block| match block {
1510 ContentBlock::Text { text, .. } => Some(text.as_str()),
1511 _ => None,
1512 })
1513 .collect::<Vec<_>>()
1514 .join("\n");
1515 assert!(
1516 transcript_text.contains("partial answer that must be discarded"),
1517 "the visible partial reply must be preserved in the session: {transcript_text}"
1518 );
1519 assert_eq!(
1520 transcript_text.matches("recovered after retry").count(),
1521 1,
1522 "exactly one authoritative final answer, not a duplicate: {transcript_text}"
1523 );
1524 }
1525
1526 /// Streams only hidden reasoning and then dies with the network-class read
1527 /// error; later streams complete a normal text turn. This is the shape of the
1528 /// 0.9.10 regression: a thinking-only drop used to persist a synthetic
1529 /// `[runtime]` user message claiming a partial answer had been preserved.
1530 struct ThinkingOnlyDropModelClient {
1531 calls: std::sync::atomic::AtomicUsize,
1532 failures: usize,
1533 }
1534
1535 #[async_trait::async_trait]
1536 impl crate::core::model_client::ModelClient for ThinkingOnlyDropModelClient {
1537 fn provider_name(&self) -> &str {
1538 "flaky-network"
1539 }
1540
1541 fn model(&self) -> &str {
1542 "local-model"
1543 }
1544
1545 async fn create_message(
1546 &self,
1547 _request: codewhale_models::MessageRequest,
1548 ) -> anyhow::Result<codewhale_models::MessageResponse> {
1549 anyhow::bail!("thinking-only drop regression uses the streaming model boundary")
1550 }
1551
1552 async fn create_message_stream(
1553 &self,
1554 _request: codewhale_models::MessageRequest,
1555 ) -> anyhow::Result<crate::llm_client::StreamEventBox> {
1556 use crate::llm_client::mock::canned;
1557 let call = self
1558 .calls
1559 .fetch_add(1, std::sync::atomic::Ordering::SeqCst)
1560 .saturating_add(1);
1561 if call <= self.failures {
1562 // Hidden reasoning only — no text block is ever opened, so nothing
1563 // visible streams before the transport dies. This still flips
1564 // `any_content_received`, which is what routes the drop to the
1565 // interactive resume path instead of the transparent retry.
1566 let events: Vec<anyhow::Result<codewhale_models::StreamEvent>> = vec![
1567 Ok(canned::message_start("thinking_only_msg")),
1568 Ok(StreamEvent::ContentBlockStart {
1569 index: 0,
1570 content_block: codewhale_models::ContentBlockStart::Thinking {
1571 thinking: String::new(),
1572 },
1573 }),
1574 Ok(canned::thinking_delta(
1575 0,
1576 "hidden reasoning that no operator ever saw",
1577 )),
1578 Err(anyhow::anyhow!(
1579 "Stream read error: error decoding response body"
1580 )),
1581 ];
1582 return Ok(Box::pin(futures_util::stream::iter(events)));
1583 }
1584 let events = canned::simple_text_turn("the one authoritative answer")
1585 .into_iter()
1586 .map(Ok);
1587 Ok(Box::pin(futures_util::stream::iter(events)))
1588 }
1589
1590 async fn health_check(&self) -> anyhow::Result<bool> {
1591 Ok(true)
1592 }
1593 }
1594
1595 struct TransportRetryModelClient {
1596 attempts: std::sync::atomic::AtomicUsize,
1597 failures: usize,
1598 }
1599
1600 #[async_trait::async_trait]
1601 impl crate::core::model_client::ModelClient for TransportRetryModelClient {
1602 fn provider_name(&self) -> &str {
1603 "custom"
1604 }
1605 fn model(&self) -> &str {
1606 crate::config::DEFAULT_TEXT_MODEL
1607 }
1608 async fn create_message(
1609 &self,
1610 _: codewhale_models::MessageRequest,
1611 ) -> anyhow::Result<codewhale_models::MessageResponse> {
1612 anyhow::bail!("streaming fixture only")
1613 }
1614 async fn create_message_stream(
1615 &self,
1616 _: codewhale_models::MessageRequest,
1617 ) -> anyhow::Result<crate::llm_client::StreamEventBox> {
1618 crate::llm_client::with_retry(
1619 &crate::llm_client::RetryConfig {
1620 max_retries: 2,
1621 initial_delay: 0.0,
1622 jitter: false,
1623 ..Default::default()
1624 },
1625 || async {
1626 let attempt = self
1627 .attempts
1628 .fetch_add(1, std::sync::atomic::Ordering::SeqCst);
1629 if attempt < self.failures {
1630 Err(crate::llm_client::LlmError::ServerError {
1631 status: 503,
1632 message: "RAW-ENGINE-RETRY-FAILURE".into(),
1633 })
1634 } else {
1635 Ok(())
1636 }
1637 },
1638 None,
1639 )
1640 .await
1641 .map_err(|error| anyhow::Error::new(error.last_error))?;
1642 Ok(Box::pin(futures_util::stream::iter(
1643 crate::llm_client::mock::canned::simple_text_turn("transport recovered answer")
1644 .into_iter()
1645 .map(Ok),
1646 )))
1647 }
1648 async fn health_check(&self) -> anyhow::Result<bool> {
1649 Ok(true)
1650 }
1651 }
1652
1653 #[tokio::test]
1654 async fn engine_transport_retry_receipts_and_counts_match_actual_scripted_attempts() {
1655 for (failures, quiet) in [(1, false), (3, false), (1, true)] {
1656 let workspace = tempdir().unwrap();
1657 let model = std::sync::Arc::new(TransportRetryModelClient {
1658 attempts: std::sync::atomic::AtomicUsize::new(0),
1659 failures,
1660 });
1661 let config = Config {
1662 notifications: Some(crate::config::NotificationsConfig {
1663 quiet,
1664 ..Default::default()
1665 }),
1666 ..Config::default()
1667 };
1668 let (mut engine, handle) = Engine::new_with_model_client(
1669 deterministic_engine_config(workspace.path()),
1670 &config,
1671 model.clone(),
1672 );
1673 let surface = test_tool_surface(
1674 &engine,
1675 crate::tools::ToolRegistry::new(crate::tools::ToolContext::new(
1676 workspace.path().to_path_buf(),
1677 )),
1678 None,
1679 AppMode::Agent,
1680 );
1681 let mut turn = crate::core::turn::TurnContext::new(1);
1682 let (status, error) = engine.run_turn(&mut turn, surface, None, None).await;
1683 let expected_retries = failures.min(2) as u32;
1684 assert_eq!(
1685 model.attempts.load(std::sync::atomic::Ordering::SeqCst),
1686 (expected_retries + 1) as usize
1687 );
1688 assert_eq!(turn.stop_diagnostics.model_requests_started, 1);
1689 assert_eq!(turn.stop_diagnostics.transport_retries, expected_retries);
1690 assert_eq!(turn.stop_diagnostics.stream_resumes, 0);
1691 assert_eq!(turn.stop_diagnostics.transparent_stream_retries, 0);
1692 let mut receiver = handle.rx_event.write().await;
1693 let events = std::iter::from_fn(|| receiver.try_recv().ok()).collect::<Vec<_>>();
1694 if quiet {
1695 assert!(
1696 !events.iter().any(|event| matches!(event,
1697 Event::Status { message } if message.starts_with("Retry attempt:")
1698 )),
1699 "quiet suppresses attempt lines, not actual counts or completed summaries"
1700 );
1701 } else {
1702 for attempt in 1..=expected_retries {
1703 assert!(events.iter().any(|event| matches!(event,
1704 Event::Status { message } if message.starts_with(&format!("Retry attempt: transport {attempt}/2;"))
1705 )));
1706 }
1707 }
1708 assert!(
1709 !events.iter().any(|event| matches!(event,
1710 Event::Status { message } if message.contains("RAW-ENGINE-RETRY-FAILURE")
1711 )),
1712 "remote payload belongs to the original error, never a retry receipt"
1713 );
1714 if failures == 1 {
1715 assert_eq!(status, TurnOutcomeStatus::Completed, "{error:?}");
1716 assert!(events.iter().any(|event| matches!(event,
1717 Event::Status { message } if message == "Retry recovery: transport request recovered after 1 retries"
1718 )));
1719 } else {
1720 assert_eq!(status, TurnOutcomeStatus::Failed);
1721 assert!(error.unwrap().contains("RAW-ENGINE-RETRY-FAILURE"));
1722 assert!(events.iter().any(|event| matches!(event,
1723 Event::Status { message } if message.starts_with("Retry exhaustion: transport request stopped after 2 retries;")
1724 )));
1725 }
1726 assert!(engine.session.messages.iter().all(|message| message.content.iter().all(|block|
1727 !matches!(block, ContentBlock::Text { text, .. } if crate::core::events::is_retry_status_receipt(text))
1728 )), "receipts must never join the provider prompt or session message graph");
1729 }
1730 }
1731
1732 #[tokio::test]
1733 async fn full_retry_receipt_queue_cancellation_does_not_dispatch_an_extra_request() {
1734 let workspace = tempdir().unwrap();
1735 let (engine, _handle) = Engine::new_with_model_client(
1736 deterministic_engine_config(workspace.path()),
1737 &Config::default(),
1738 std::sync::Arc::new(TransportRetryModelClient {
1739 attempts: std::sync::atomic::AtomicUsize::new(0),
1740 failures: 0,
1741 }),
1742 );
1743 while engine.tx_event.try_send(Event::status("occupied")).is_ok() {}
1744 let observation = engine.request_retry_observation();
1745 let count = observation.retries.clone();
1746 let mut calls = 0;
1747 let policy = crate::llm_client::RetryConfig {
1748 max_retries: 2,
1749 initial_delay: 0.0,
1750 jitter: false,
1751 ..Default::default()
1752 };
1753 let outcome = {
1754 let request = crate::llm_client::observe_request_retries(
1755 Some(observation),
1756 crate::llm_client::with_retry(
1757 &policy,
1758 || {
1759 calls += 1;
1760 async {
1761 Err::<(), _>(crate::llm_client::LlmError::ServerError {
1762 status: 503,
1763 message: "private".into(),
1764 })
1765 }
1766 },
1767 None,
1768 ),
1769 );
1770 tokio::pin!(request);
1771 assert!(
1772 futures_util::poll!(request.as_mut()).is_pending(),
1773 "the receipt waits on the one full queue"
1774 );
1775 engine.cancel_token.cancel();
1776 tokio::time::timeout(Duration::from_secs(1), async {
1777 tokio::select! {
1778 biased;
1779 () = engine.cancel_token.cancelled() => None,
1780 result = request.as_mut() => Some(result),
1781 }
1782 })
1783 .await
1784 .unwrap()
1785 };
1786 assert!(outcome.is_none());
1787 assert_eq!(calls, 1);
1788 assert_eq!(
1789 count.load(std::sync::atomic::Ordering::Relaxed),
1790 0,
1791 "a scheduled retry cancelled before dispatch is not spent"
1792 );
1793 }
1794
1794 lines RUST