返回 CodeWhale
budget_handback_tests.rs
根目录 / crates / tui / src / tools / subagent / budget_handback_tests.rs
1 use super::tests::{chat_fixture_response, make_worker_spec, stub_runtime};
2 use super::*;
3 use axum::{Json, Router, http::StatusCode, response::IntoResponse, routing::post};
4 use tempfile::{TempDir, tempdir};
5 use tokio::sync::Notify;
6
7 struct Fixture {
8 workspace: TempDir,
9 manager: SharedSubAgentManager,
10 task: Option<JoinHandle<()>>,
11 server: JoinHandle<()>,
12 requests: Arc<std::sync::Mutex<Vec<Value>>>,
13 report_started: Arc<Notify>,
14 release_report: Arc<Notify>,
15 cancel: CancellationToken,
16 completions: mpsc::Receiver<SubAgentCompletion>,
17 mailbox: MailboxReceiver,
18 }
19
20 impl Fixture {
21 async fn finish(&mut self) -> SubAgentResult {
22 tokio::time::timeout(Duration::from_secs(30), self.task.take().unwrap())
23 .await
24 .expect("bounded worker")
25 .expect("worker task");
26 self.manager
27 .read()
28 .await
29 .get_result("report-worker")
30 .unwrap()
31 }
32 }
33
34 impl Drop for Fixture {
35 fn drop(&mut self) {
36 if let Some(task) = &self.task {
37 task.abort();
38 }
39 self.server.abort();
40 }
41 }
42
43 async fn fixture(mode: &'static str, first_tokens: u64, max_steps: u32) -> Fixture {
44 // Only these cases exercise wall/API timeouts. The other cases exercise
45 // budget and report semantics, so leave room for full-suite scheduling.
46 let timeout_case = matches!(mode, "hold" | "timeout" | "work-timeout");
47 // "work-timeout" narrows its own deadline below; its 30s budget sizes the
48 // hand-back reserve and must not undercut that deadline.
49 let wall_time_secs = if matches!(mode, "hold" | "timeout") {
50 5
51 } else {
52 30
53 };
54 let workspace = tempdir().unwrap();
55 fs::write(
56 workspace.path().join("README.md"),
57 "TOOL_EVIDENCE: checksum validation is still missing.\n",
58 )
59 .unwrap();
60 let requests = Arc::new(std::sync::Mutex::new(Vec::new()));
61 let report_started = Arc::new(Notify::new());
62 let release_report = Arc::new(Notify::new());
63 let app = Router::new().route("/{*path}", post({
64 let requests = Arc::clone(&requests);
65 let report_started = Arc::clone(&report_started);
66 let release_report = Arc::clone(&release_report);
67 move |Json(body): Json<Value>| {
68 let requests = Arc::clone(&requests);
69 let report_started = Arc::clone(&report_started);
70 let release_report = Arc::clone(&release_report);
71 async move {
72 let call = {
73 let mut requests = requests.lock().unwrap();
74 requests.push(body.clone());
75 requests.len()
76 };
77 let choice = if call == 1 {
78 json!({"index": 0, "message": {"role": "assistant",
79 "content": if mode == "tool-only" { Value::Null } else { json!("RECORDED_FINDING: checksum validation is missing.") },
80 "tool_calls": [{"id": "read-one", "type": "function", "function": {
81 "name": "read", "arguments": "{\"path\":\"README.md\"}"
82 }}]}, "finish_reason": "tool_calls"})
83 } else {
84 report_started.notify_one();
85 if matches!(mode, "hold" | "timeout" | "work-timeout") { release_report.notified().await; }
86 if mode == "failure" {
87 return (StatusCode::BAD_REQUEST, Json(json!({"error": {"message": "fixture rejection"}}))).into_response();
88 }
89 if mode == "tool" {
90 json!({"index": 0, "message": {"role": "assistant", "content": "REJECTED_REPORT: I wrote report.md.",
91 "tool_calls": [{"id": "must-not-write", "type": "function", "function": {
92 "name": "write_file", "arguments": "{\"path\":\"report.md\",\"content\":\"must not execute\"}"
93 }}]}, "finish_reason": "tool_calls"})
94 } else if mode == "truncated" {
95 json!({"index": 0, "message": {"role": "assistant", "content": "TRUNCATED_REPORT: README evid"}, "finish_reason": "length"})
96 } else {
97 json!({"index": 0, "message": {"role": "assistant", "content":
98 "PARTIAL_REPORT: README evidence identifies missing checksum validation. No report file was produced. Next: implement and verify the checksum check."}, "finish_reason": "stop"})
99 }
100 };
101 let usage = if (call == 1 && matches!(mode, "unknown" | "resume-unknown")) || (call > 1 && mode == "report-unknown") { Value::Null }
102 else if call == 1 { json!({"prompt_tokens": first_tokens.saturating_sub(5), "completion_tokens": 5, "total_tokens": first_tokens}) }
103 else if mode == "resume-unknown" { json!({"prompt_tokens": 10, "completion_tokens": 5, "total_tokens": 15}) }
104 else { json!({"prompt_tokens": 20, "completion_tokens": 10, "total_tokens": 30}) };
105 chat_fixture_response(body.get("stream").and_then(Value::as_bool).unwrap_or(false), json!({"id": format!("handback-{call}"), "model": "deepseek-v4-flash", "choices": [choice], "usage": usage})).into_response()
106 }
107 }
108 }));
109 let listener = tokio::net::TcpListener::bind("127.0.0.1:0").await.unwrap();
110 let address = listener.local_addr().unwrap();
111 let server = tokio::spawn(async move {
112 axum::serve(listener, app).await.unwrap();
113 });
114 let config = crate::config::Config {
115 retry: Some(crate::config::RetryConfig {
116 enabled: Some(false),
117 max_retries: Some(0),
118 initial_delay: Some(0.0),
119 max_delay: Some(0.0),
120 exponential_base: Some(1.0),
121 jitter: None,
122 jitter_factor: None,
123 respect_retry_after: None,
124 }),
125 ..Default::default()
126 }
127 .with_legacy_root(
128 Some("fixture-key".to_string()),
129 Some(format!("http://{address}/v1")),
130 );
131 let manager = Arc::new(RwLock::new(
132 SubAgentManager::new(workspace.path().to_path_buf(), 4)
133 .with_state_path(workspace.path().join(".codewhale/subagents/state.json")),
134 ));
135 let mut spec = make_worker_spec("report-worker", workspace.path().to_path_buf());
136 spec.max_steps = max_steps;
137 spec.runtime_profile.max_steps = max_steps;
138 spec.runtime_profile.wall_time_secs = Some(wall_time_secs);
139 spec.runtime_profile.wall_deadline_ms = Some(epoch_millis_now() + wall_time_secs * 1_000);
140 if mode == "work-timeout" {
141 // A partly consumed original deadline leaves time to persist the
142 // missing-coverage receipt after the in-flight call is abandoned.
143 //
144 // The hand-back window is the reserve `wall_deadlines` carves off the
145 // hard deadline: `wall_time_secs * 100ms`. The shared 5s budget made
146 // that only 500ms, which had to cover the digest artifact, the
147 // unreported-usage state write, the pre-report checkpoint and the
148 // loopback connect before the report reached the server. On Windows
149 // CI it did not fit, so the report fell back without reaching the
150 // server (2 requests, not 3). A 30s budget reserves 3s. The work
151 // deadline lands 6s in: the first call, its tool, and the dispatch of
152 // the second call must all fit before it, and a loaded shared-process
153 // `cargo test --workspace` run overran the earlier 2s (#6698). The
154 // step API timeout below is longer than that, so the in-flight call
155 // is still abandoned by wall time, not by a step timeout.
156 spec.runtime_profile.wall_deadline_ms = Some(epoch_millis_now() + 9_000);
157 }
158 spec.launch_manifest = Some(serde_json::from_value(json!({
159 "owner_session": "root", "child_id": "report-worker", "profile": spec.runtime_profile,
160 "prompt": "Read README.md and produce report.md", "cwd": workspace.path(), "worktree": false,
161 "writable_roots": [], "writable_files": [], "coordination_contracts": [], "deliverables": ["report.md"],
162 "resume_identity": null, "generation": 1, "resume_from_agent_id": null
163 })).unwrap());
164 if mode == "resume-unknown" {
165 spec.launch_manifest.as_mut().unwrap().deliverables.clear();
166 }
167 let mut runtime = stub_runtime();
168 runtime.client = CodewhaleClient::new(&config).unwrap();
169 runtime.api_config = Some(Arc::new(config));
170 runtime.context = ToolContext::new(workspace.path().to_path_buf());
171 runtime.accounting_origin = SubAgentAccountingOrigin::capture(&runtime.context);
172 runtime.manager = Arc::clone(&manager);
173 runtime.worker_profile = spec.runtime_profile.clone();
174 runtime.spawn_depth = spec.spawn_depth;
175 runtime.allow_shell = false;
176 runtime.accept_edits = false;
177 runtime.step_api_timeout = if mode == "timeout" {
178 Duration::from_millis(100)
179 } else if mode == "work-timeout" {
180 // Past the 6s work deadline, and bounding the held hand-back call
181 // no tighter than the 9s hard deadline already does.
182 Duration::from_secs(10)
183 } else if timeout_case {
184 Duration::from_secs(2)
185 } else {
186 Duration::from_secs(10)
187 };
188 let cancel = runtime.cancel_token.clone();
189 let (parent_tx, completions) = mpsc::channel(16);
190 runtime.parent_completion_tx = Some(parent_tx);
191 let (mailbox, mailbox_rx) = Mailbox::new(CancellationToken::new());
192 runtime.mailbox = Some(mailbox);
193 let assignment = SubAgentAssignment::new(
194 "Read README.md, report findings and identify what remains.".to_string(),
195 None,
196 );
197 let (input_tx, input_rx) = mpsc::unbounded_channel();
198 let mut agent = SubAgent::new(
199 "report-worker".to_string(),
200 FleetRole::Scout,
201 assignment.objective.clone(),
202 assignment.clone(),
203 runtime.model.clone(),
204 None,
205 Some(vec!["read_file".to_string()]),
206 input_tx,
207 workspace.path().to_path_buf(),
208 manager.read().await.current_session_boot_id.clone(),
209 );
210 agent.status = SubAgentStatus::Running;
211 {
212 let mut guard = manager.write().await;
213 guard.register_worker_for_session(spec, &runtime.context.state_namespace, None);
214 guard.agents.insert("report-worker".to_string(), agent);
215 }
216 let task = tokio::spawn(run_subagent_task(SubAgentTask {
217 manager_handle: Arc::clone(&manager),
218 runtime,
219 agent_id: "report-worker".to_string(),
220 agent_type: FleetRole::Scout,
221 prompt: assignment.objective.clone(),
222 assignment,
223 allowed_tools: Some(vec!["read_file".to_string()]),
224 fork_context: false,
225 started_at: Instant::now(),
226 max_steps,
227 wall_time: Duration::from_secs(wall_time_secs),
228 wall_ceiling_ms: None,
229 input_rx,
230 launch_gate: None,
231 _foreground_child_registration: None,
232 }));
233 Fixture {
234 workspace,
235 manager,
236 task: Some(task),
237 server,
238 requests,
239 report_started,
240 release_report,
241 cancel,
242 completions,
243 mailbox: mailbox_rx,
244 }
245 }
246
247 #[tokio::test]
248 #[allow(clippy::await_holding_lock)]
249 async fn budget_handback_turn_consolidates_tool_only_work_and_checks_declared_deliverables() {
250 let _retry = crate::retry_status::test_guard();
251 crate::retry_status::clear_rate_limit();
252 let mut fixture = fixture("tool-only", 15, 1).await;
253 let result = fixture.finish().await;
254 assert_eq!(result.status, SubAgentStatus::BudgetExhausted);
255 assert_eq!(result.steps_taken, 2);
256 assert_eq!(result.usage.as_ref().unwrap().total_tokens, Some(45));
257 assert!(result.result.as_deref().unwrap().contains("PARTIAL_REPORT"));
258 assert!(
259 result
260 .checkpoint
261 .as_ref()
262 .unwrap()
263 .messages
264 .iter()
265 .flat_map(|message| &message.content)
266 .any(
267 |block| matches!(block, ContentBlock::ToolResult { tool_use_id, content, .. }
268 if tool_use_id == "read-one" && content.contains("TOOL_EVIDENCE"))
269 ),
270 "the read must execute successfully before reporting: {:?}",
271 result.checkpoint,
272 );
273 let requests = fixture.requests.lock().unwrap();
274 assert_eq!(requests.len(), 2);
275 assert_eq!(requests[0]["model"], "deepseek-v4-flash");
276 assert!(
277 requests[0]["tools"]
278 .as_array()
279 .unwrap()
280 .iter()
281 .any(|tool| tool["function"]["name"] == "read")
282 );
283 assert_eq!(requests[1]["model"], requests[0]["model"]);
284 assert!(requests[1].get("tools").is_none_or(Value::is_null));
285 assert!(requests[1].get("tool_choice").is_none_or(Value::is_null));
286 assert!(
287 requests[1].to_string().contains("TOOL_EVIDENCE"),
288 "report must read actual completed tool output"
289 );
290 assert!(
291 requests[1]["max_tokens"]
292 .as_u64()
293 .or_else(|| requests[1]["max_completion_tokens"].as_u64())
294 .unwrap()
295 <= 1_024
296 );
297 drop(requests);
298 let guard = fixture.manager.read().await;
299 let worker = &guard.worker_records["report-worker"];
300 assert_eq!(worker.verification.status, "deliverable_missing");
301 assert_eq!(worker.verification.deliverables[0].path, "report.md");
302 assert!(
303 !worker.spec.runtime_profile.permissions.write,
304 "report did not widen the Scout's authority"
305 );
306 drop(guard);
307 let completion = fixture.completions.try_recv().unwrap();
308 assert!(completion.payload.contains("budget_exhausted"));
309 assert!(completion.payload.contains("deliverable_missing"));
310 assert!(fixture.completions.try_recv().is_err());
311 }
312
313 /// #6536 — the provider truncates the hand-back report. The deterministic
314 /// digest recorded before that turn stays the deliverable: in the result
315 /// text `agent result` / `agent wait` return, and as a private file.
316 #[tokio::test]
317 #[allow(clippy::await_holding_lock)]
318 async fn budget_handback_truncated_report_leaves_the_digest_as_the_deliverable() {
319 let _retry = crate::retry_status::test_guard();
320 crate::retry_status::clear_rate_limit();
321 let mut fixture = fixture("truncated", 15, 1).await;
322 let result = fixture.finish().await;
323 assert_eq!(fixture.requests.lock().unwrap().len(), 2);
324 assert_eq!(result.status, SubAgentStatus::BudgetExhausted);
325 let text = result.result.as_deref().unwrap();
326 assert!(
327 text.contains("RECORDED_FINDING"),
328 "digest is the result: {text}"
329 );
330 assert!(!text.contains("TRUNCATED_REPORT"), "{text}");
331 assert!(text.contains("did not finish"), "{text}");
332 assert!(text.contains("this child's deliverable"), "{text}");
333
334 let state_root = fixture.manager.read().await.state_root.clone();
335 let artifact = checked_subagent_state_path(
336 &state_root,
337 &Path::new(".codewhale/state/subagent-results").join(format!(
338 "{}.md",
339 crate::hashing::sha256_hex(b"report-worker")
340 )),
341 )
342 .unwrap();
343 let saved = fs::read_to_string(&artifact).expect("digest artifact written");
344 assert!(saved.contains("RECORDED_FINDING"), "{saved}");
345 assert!(!saved.contains("TRUNCATED_REPORT"), "{saved}");
346 assert!(text.contains(&artifact.display().to_string()), "{text}");
347 }
348
349 #[tokio::test]
350 #[allow(clippy::await_holding_lock)]
351 async fn budget_handback_turn_rejects_provider_tools_and_preserves_fallback_verdicts() {
352 let _retry = crate::retry_status::test_guard();
353 crate::retry_status::clear_rate_limit();
354 let mut fixture = fixture("tool", 15, 1).await;
355 let result = fixture.finish().await;
356 assert_eq!(fixture.requests.lock().unwrap().len(), 2);
357 assert_eq!(result.status, SubAgentStatus::BudgetExhausted);
358 assert!(!fixture.workspace.path().join("report.md").exists());
359 assert!(
360 result
361 .result
362 .as_deref()
363 .unwrap()
364 .contains("provider returned a tool call")
365 );
366 assert!(
367 !result
368 .result
369 .as_deref()
370 .unwrap()
371 .contains("REJECTED_REPORT")
372 );
373 assert_eq!(result.usage.as_ref().unwrap().total_tokens, Some(45));
374 let history = &result.checkpoint.as_ref().unwrap().messages;
375 let completed = history
376 .iter()
377 .flat_map(|message| &message.content)
378 .filter_map(|block| match block {
379 ContentBlock::ToolResult { tool_use_id, .. } => Some(tool_use_id.as_str()),
380 _ => None,
381 })
382 .collect::<HashSet<_>>();
383 for block in history.iter().flat_map(|message| &message.content) {
384 match block {
385 ContentBlock::ToolUse { id, .. } => assert!(
386 completed.contains(id.as_str()),
387 "no orphan tool call may survive for replay"
388 ),
389 ContentBlock::ServerToolUse { .. } => {
390 panic!("a rejected server tool call entered history")
391 }
392 ContentBlock::Text { text, .. } => assert!(!text.contains("REJECTED_REPORT")),
393 _ => {}
394 }
395 }
396 assert!(
397 history
398 .iter()
399 .flat_map(|message| &message.content)
400 .any(|block| matches!(block,
401 ContentBlock::Text { text, .. } if text.contains("Host budget hand-back receipt")))
402 );
403 assert_eq!(
404 fixture.manager.read().await.worker_records["report-worker"]
405 .verification
406 .status,
407 "deliverable_missing"
408 );
409 }
410
411 #[tokio::test]
412 #[allow(clippy::await_holding_lock)]
413 async fn budget_handback_turn_failure_and_timeout_do_not_add_worker_retries() {
414 let _retry = crate::retry_status::test_guard();
415 crate::retry_status::clear_rate_limit();
416 for (mode, reason) in [
417 ("failure", "provider call failed"),
418 ("timeout", "report deadline expired"),
419 ] {
420 let mut fixture = fixture(mode, 15, 1).await;
421 let result = fixture.finish().await;
422 assert_eq!(fixture.requests.lock().unwrap().len(), 2);
423 assert_eq!(result.status, SubAgentStatus::BudgetExhausted);
424 assert!(
425 result.result.as_deref().unwrap().contains(reason),
426 "{result:?}"
427 );
428 assert_eq!(result.usage.as_ref().unwrap().total_tokens, Some(15));
429 assert_eq!(
430 fixture.manager.read().await.worker_records["report-worker"].has_unreported_usage,
431 mode == "timeout",
432 "only the timed-out dispatched call establishes missing coverage here",
433 );
434 assert_eq!(
435 fixture.manager.read().await.worker_records["report-worker"]
436 .verification
437 .status,
438 "deliverable_missing"
439 );
440 }
441 }
442
443 #[tokio::test]
444 #[allow(clippy::await_holding_lock)]
445 async fn budget_handback_turn_missing_report_usage_is_not_claimed_as_zero_cost() {
446 let _retry = crate::retry_status::test_guard();
447 crate::retry_status::clear_rate_limit();
448 let mut fixture = fixture("report-unknown", 15, 1).await;
449 let result = fixture.finish().await;
450 assert_eq!(fixture.requests.lock().unwrap().len(), 2);
451 assert_eq!(result.status, SubAgentStatus::BudgetExhausted);
452 assert_eq!(result.usage.as_ref().unwrap().total_tokens, Some(15));
453 let report = result.result.as_deref().unwrap();
454 assert!(report.contains("PARTIAL_REPORT"));
455 assert!(report.contains("only a subtotal, not a zero-cost report"));
456 let host_receipts: Vec<_> = result
457 .checkpoint
458 .as_ref()
459 .unwrap()
460 .messages
461 .iter()
462 .flat_map(|message| &message.content)
463 .filter_map(|block| match block {
464 ContentBlock::Text { text, .. }
465 if text.starts_with("Host budget hand-back receipt:") =>
466 {
467 Some(text)
468 }
469 _ => None,
470 })
471 .collect();
472 assert_eq!(
473 host_receipts.len(),
474 1,
475 "one attributed canonical report outcome"
476 );
477 assert!(host_receipts[0].contains("unreported provider usage remains unknown"));
478 }
479
480 #[tokio::test]
481 #[allow(clippy::await_holding_lock)]
482 async fn budget_handback_inflight_wall_timeout_persists_unreported_usage() {
483 let _retry = crate::retry_status::test_guard();
484 crate::retry_status::clear_rate_limit();
485 let mut fixture = fixture("work-timeout", 15, 4).await;
486 let result = fixture.finish().await;
487 assert_eq!(result.status, SubAgentStatus::BudgetExhausted);
488 assert!(
489 result
490 .result
491 .as_deref()
492 .unwrap()
493 .contains("child wall-time work budget exhausted")
494 );
495 assert_eq!(
496 fixture.requests.lock().unwrap().len(),
497 3,
498 "the reserved hand-back turn still dispatches after an unmeasured in-flight request; its own deadline bounds it"
499 );
500 let manager = fixture.manager.read().await;
501 assert!(manager.worker_records["report-worker"].has_unreported_usage);
502 assert_eq!(
503 manager.worker_records["report-worker"].usage.total_tokens,
504 Some(15)
505 );
506 }
507
508 #[test]
509 fn budget_handback_coverage_marker_is_sticky_without_reclassifying_legacy_or_measured_zero() {
510 let tmp = tempdir().unwrap();
511 let mut manager = SubAgentManager::new(tmp.path().to_path_buf(), 4);
512 for id in ["worker", "unrelated"] {
513 let spec = make_worker_spec(id, tmp.path().to_path_buf());
514 manager.register_worker(spec);
515 }
516 let measured_zero = Usage {
517 prompt_cache_hit_tokens: Some(0),
518 ..Usage::default()
519 };
520 manager.record_worker_usage("worker", "zero", &measured_zero, None);
521 manager.record_worker_usage("unrelated", "unrelated-missing", &Usage::default(), None);
522 assert_eq!(manager.worker_records["worker"].usage.total_tokens, Some(0));
523 assert!(!manager.worker_records["worker"].has_unreported_usage);
524 let mut legacy = serde_json::to_value(&manager.worker_records["worker"]).unwrap();
525 legacy
526 .as_object_mut()
527 .unwrap()
528 .remove("has_unreported_usage");
529 let legacy: AgentWorkerRecord = serde_json::from_value(legacy).unwrap();
530 assert!(!legacy.has_unreported_usage);
531 manager.worker_records.insert("worker".to_string(), legacy);
532 let (_output, lease) = manager.reserve_handback("worker", 500, 1_024).unwrap();
533 drop(lease);
534 manager.record_worker_usage("worker", "missing", &Usage::default(), None);
535 manager.record_worker_usage(
536 "worker",
537 "known-later",
538 &Usage {
539 input_tokens: 10,
540 output_tokens: 5,
541 ..Usage::default()
542 },
543 None,
544 );
545 let saved = serde_json::to_vec(&manager.worker_records).unwrap();
546 manager.worker_records = serde_json::from_slice(&saved).unwrap();
547 assert_eq!(
548 manager.worker_records["worker"].usage.total_tokens,
549 Some(15)
550 );
551 assert!(manager.worker_records["worker"].has_unreported_usage);
552 }
553
554 #[tokio::test]
555 #[allow(clippy::await_holding_lock)]
556 async fn budget_handback_turn_cancellation_wins_once_and_releases_shared_reservation() {
557 let _retry = crate::retry_status::test_guard();
558 crate::retry_status::clear_rate_limit();
559 let mut fixture = fixture("hold", 15, 2).await;
560 tokio::time::timeout(Duration::from_secs(2), fixture.report_started.notified())
561 .await
562 .unwrap();
563 fixture.cancel.cancel();
564 let result = fixture.finish().await;
565 assert_eq!(result.status, SubAgentStatus::Cancelled);
566 assert_eq!(result.usage.as_ref().unwrap().total_tokens, Some(15));
567 assert!(fixture.manager.read().await.worker_records["report-worker"].has_unreported_usage);
568 assert!(
569 fixture
570 .completions
571 .try_recv()
572 .unwrap()
573 .payload
574 .contains("cancelled")
575 );
576 assert!(fixture.completions.try_recv().is_err());
577 assert!(
578 fixture
579 .manager
580 .read()
581 .await
582 .handback_reservations
583 .values()
584 .all(|value| value.upgrade().is_none())
585 );
586 }
587
588 #[tokio::test]
589 #[allow(clippy::await_holding_lock)]
590 async fn budget_handback_turn_cancellation_after_response_preserves_actual_usage() {
591 let _retry = crate::retry_status::test_guard();
592 crate::retry_status::clear_rate_limit();
593 let mut fixture = fixture("hold", 15, 2).await;
594 tokio::time::timeout(Duration::from_secs(2), fixture.report_started.notified())
595 .await
596 .unwrap();
597 let manager = Arc::clone(&fixture.manager);
598 let guard = manager.write().await;
599 fixture.release_report.notify_one();
600 // Runtime billing publishes the decoded response before it waits for the
601 // worker ledger lock. This makes the cancellation seam deterministic.
602 tokio::time::timeout(Duration::from_secs(2), async {
603 loop {
604 let entry = fixture.mailbox.recv().await.unwrap();
605 if let MailboxMessage::TokenUsage {
606 agent_id,
607 source_id,
608 usage,
609 ..
610 } = entry.message
611 && agent_id == "report-worker"
612 && usage_total_tokens(&usage) == 30
613 {
614 assert!(source_id.starts_with("child:report-worker:turn:"));
615 assert!(source_id.contains(":request:1:dispatch:"));
616 break;
617 }
618 }
619 })
620 .await
621 .unwrap();
622 fixture.cancel.cancel();
623 drop(guard);
624 let result = fixture.finish().await;
625 assert_eq!(result.status, SubAgentStatus::Cancelled);
626 assert_eq!(result.usage.as_ref().unwrap().total_tokens, Some(45));
627 assert!(!fixture.manager.read().await.worker_records["report-worker"].has_unreported_usage);
628 assert!(
629 fixture
630 .completions
631 .try_recv()
632 .unwrap()
633 .payload
634 .contains("cancelled")
635 );
636 assert!(fixture.completions.try_recv().is_err());
637 }
638
639 #[test]
640 fn handback_reservation_uses_the_fixed_allowance_and_refuses_a_second_turn() {
641 let tmp = tempdir().unwrap();
642 let mut manager = SubAgentManager::new(tmp.path().to_path_buf(), 4);
643 manager.register_worker(make_worker_spec("w", tmp.path().to_path_buf()));
644 let (output, lease) = manager.reserve_handback("w", 500, 1_024).unwrap();
645 assert_eq!(output, 1_024);
646 assert!(matches!(
647 manager.reserve_handback("w", 500, 1_024),
648 Err(reason) if reason.contains("already in flight")
649 ));
650 drop(lease);
651 assert!(matches!(
652 manager.reserve_handback("w", 8_200, 1_024),
653 Err(reason) if reason.contains("fixed hand-back allowance")
654 ));
655 }
656
657 #[tokio::test]
658 async fn budget_handback_expired_original_deadline_refuses_the_model_call() {
659 let tmp = tempdir().unwrap();
660 let mut runtime = stub_runtime();
661 runtime.context = ToolContext::new(tmp.path().to_path_buf());
662 runtime.manager = Arc::new(RwLock::new(SubAgentManager::new(
663 tmp.path().to_path_buf(),
664 1,
665 )));
666 runtime
667 .manager
668 .write()
669 .await
670 .register_worker(make_worker_spec("expired", tmp.path().to_path_buf()));
671 runtime.worker_profile.wall_deadline_ms = Some(epoch_millis_now().saturating_sub(1));
672 let authority = engine::ChildAuthority::capture(
673 runtime.clone(),
674 FleetRole::Worker,
675 "expired".into(),
676 "report".into(),
677 None,
678 );
679 let job = engine::ChildJob::admitted(
680 authority,
681 SubAgentAssignment::new("report".into(), None),
682 Instant::now(),
683 2,
684 false,
685 None,
686 )
687 .await
688 .unwrap();
689 let outcome =
690 budget_handback::admit_report(&job, &[], "wall-time budget exhausted", 1024).await;
691 assert!(matches!(outcome, Err(ref why) if why.contains("deadline has expired")));
692 assert_eq!(job.steps(), 0, "no model turn was admitted");
693 assert!(!runtime.manager.read().await.worker_records["expired"].has_unreported_usage);
694 assert!(
695 runtime
696 .manager
697 .read()
698 .await
699 .handback_reservations
700 .is_empty()
701 );
702 }
703
704 fn git(root: &Path, args: &[&str]) {
705 let output = std::process::Command::new("git")
706 .arg("-C")
707 .arg(root)
708 .args(args)
709 .output()
710 .expect("git");
711 assert!(
712 output.status.success(),
713 "git {args:?}: {}",
714 String::from_utf8_lossy(&output.stderr)
715 );
716 }
717
718 /// Addendum F4 (fleet-5): cancelling keeps the work. A Stop on a
719 /// write-scoped child appends the same preservation receipt a budget death
720 /// gets, exactly once, and leaves a read-only child's result alone.
721 #[tokio::test]
722 async fn cancel_appends_work_preservation_note_once() {
723 let tmp = tempdir().unwrap();
724 let root = tmp.path();
725 git(root, &["init", "--quiet"]);
726 git(root, &["config", "user.name", "Budget test"]);
727 git(root, &["config", "user.email", "budget@example.invalid"]);
728 fs::write(root.join("src.rs"), "baseline\n").unwrap();
729 git(root, &["add", "--", "src.rs"]);
730 git(root, &["commit", "--quiet", "-m", "baseline"]);
731
732 let manager = Arc::new(RwLock::new(SubAgentManager::new(root.to_path_buf(), 2)));
733 for (agent_id, write) in [("cancel-writer", true), ("cancel-scout", false)] {
734 let mut spec = make_worker_spec(agent_id, root.to_path_buf());
735 spec.runtime_profile.permissions.write = write;
736 let mut guard = manager.write().await;
737 guard.register_worker(spec);
738 let (input_tx, _input_rx) = mpsc::unbounded_channel();
739 let mut agent = SubAgent::new(
740 agent_id.to_string(),
741 FleetRole::Worker,
742 "work that gets stopped".to_string(),
743 SubAgentAssignment {
744 native_preset: None,
745 objective: "edit".to_string(),
746 role: Some("worker".to_string()),
747 },
748 "deepseek-v4-flash".to_string(),
749 None,
750 None,
751 input_tx,
752 root.to_path_buf(),
753 guard.current_session_boot_id.clone(),
754 );
755 agent.task_handle = Some(tokio::spawn(async {
756 tokio::time::sleep(Duration::from_secs(60)).await;
757 }));
758 guard.agents.insert(agent_id.to_string(), agent);
759 }
760 fs::create_dir_all(root.join("scratch")).unwrap();
761 fs::write(root.join("scratch/half-done.rs"), "wip\n").unwrap();
762
763 let stopped = manager.write().await.cancel_agent("cancel-writer").unwrap();
764 let preserved = preserve_cancelled_work(&manager, stopped).await;
765 let text = preserved.result.as_deref().unwrap_or_default();
766 assert!(text.starts_with(CANCELLED_BY_PARENT_RESULT), "{text}");
767 assert!(text.contains("scratch/half-done.rs"), "{text}");
768 let stored = manager.read().await.get_result("cancel-writer").unwrap();
769 assert_eq!(stored.result, preserved.result, "the receipt is persisted");
770
771 // A repeated Stop does not stack a second receipt.
772 let again = manager.write().await.cancel_agent("cancel-writer").unwrap();
773 let again = preserve_cancelled_work(&manager, again).await;
774 assert_eq!(again.result, preserved.result);
775
776 // A read-only child has no baseline: its result stays the plain Stop.
777 let scout = manager.write().await.cancel_agent("cancel-scout").unwrap();
778 let scout = preserve_cancelled_work(&manager, scout).await;
779 assert_eq!(scout.result.as_deref(), Some(CANCELLED_BY_PARENT_RESULT));
780 }
781
782 /// F4: a Stop cascades to descendants, and each write-scoped descendant
783 /// stopped with the parent gets its own receipt; a read-only one does not.
784 #[tokio::test]
785 async fn cancel_receipts_each_writing_descendant_stopped_with_its_parent() {
786 let tmp = tempdir().unwrap();
787 let root = tmp.path();
788 git(root, &["init", "--quiet"]);
789 git(root, &["config", "user.name", "Budget test"]);
790 git(root, &["config", "user.email", "budget@example.invalid"]);
791 fs::write(root.join("src.rs"), "baseline\n").unwrap();
792 git(root, &["add", "--", "src.rs"]);
793 git(root, &["commit", "--quiet", "-m", "baseline"]);
794
795 let manager = Arc::new(RwLock::new(SubAgentManager::new(root.to_path_buf(), 4)));
796 for (agent_id, write, parent) in [
797 ("tree-parent", true, None),
798 ("tree-writer", true, Some("tree-parent")),
799 ("tree-scout", false, Some("tree-parent")),
800 ("tree-stranger", true, None),
801 ] {
802 let mut spec = make_worker_spec(agent_id, root.to_path_buf());
803 spec.runtime_profile.permissions.write = write;
804 spec.parent_run_id = parent.map(str::to_string);
805 let mut guard = manager.write().await;
806 guard.register_worker(spec);
807 let (input_tx, _input_rx) = mpsc::unbounded_channel();
808 let mut agent = SubAgent::new(
809 agent_id.to_string(),
810 FleetRole::Worker,
811 "work that gets stopped".to_string(),
812 SubAgentAssignment {
813 native_preset: None,
814 objective: "edit".to_string(),
815 role: Some("worker".to_string()),
816 },
817 "deepseek-v4-flash".to_string(),
818 None,
819 None,
820 input_tx,
821 root.to_path_buf(),
822 guard.current_session_boot_id.clone(),
823 );
824 agent.task_handle = Some(tokio::spawn(async {
825 tokio::time::sleep(Duration::from_secs(60)).await;
826 }));
827 guard.agents.insert(agent_id.to_string(), agent);
828 }
829 fs::create_dir_all(root.join("scratch")).unwrap();
830 fs::write(root.join("scratch/half-done.rs"), "wip\n").unwrap();
831
832 // The cascade order of `cancel_agent_for_session`: descendants, then the
833 // target. The unrelated writer is stopped too but is not a descendant.
834 let parent = {
835 let mut guard = manager.write().await;
836 for id in ["tree-writer", "tree-scout", "tree-stranger"] {
837 guard.cancel_agent(id).unwrap();
838 }
839 guard.cancel_agent("tree-parent").unwrap()
840 };
841 let parent = preserve_cancelled_work(&manager, parent).await;
842 assert!(
843 parent
844 .result
845 .as_deref()
846 .is_some_and(|text| text.contains("scratch/half-done.rs")),
847 "{:?}",
848 parent.result
849 );
850
851 let guard = manager.read().await;
852 let writer = guard.get_result("tree-writer").unwrap();
853 let writer_text = writer.result.as_deref().unwrap_or_default();
854 assert!(
855 writer_text.starts_with(CANCELLED_BY_PARENT_RESULT),
856 "{writer_text}"
857 );
858 assert!(
859 writer_text.contains("scratch/half-done.rs"),
860 "{writer_text}"
861 );
862 let scout = guard.get_result("tree-scout").unwrap();
863 assert_eq!(scout.result.as_deref(), Some(CANCELLED_BY_PARENT_RESULT));
864 let stranger = guard.get_result("tree-stranger").unwrap();
865 assert_eq!(stranger.result.as_deref(), Some(CANCELLED_BY_PARENT_RESULT));
866 }
867
868 /// #5529: a budget death must name the work the worker left on disk. The
869 /// spawn-time delivery baseline is what makes the inventory attributable to
870 /// this worker rather than the parent's own dirty files.
871 #[tokio::test]
872 async fn run_death_preservation_note_names_surviving_workspace_changes() {
873 let tmp = tempdir().unwrap();
874 let root = tmp.path();
875 git(root, &["init", "--quiet"]);
876 git(root, &["config", "user.name", "Budget test"]);
877 git(root, &["config", "user.email", "budget@example.invalid"]);
878 fs::write(root.join("src.rs"), "baseline\n").unwrap();
879 git(root, &["add", "--", "src.rs"]);
880 git(root, &["commit", "--quiet", "-m", "baseline"]);
881
882 let manager = Arc::new(RwLock::new(SubAgentManager::new(root.to_path_buf(), 2)));
883 let mut spec = make_worker_spec("preserve-worker", root.to_path_buf());
884 spec.runtime_profile.permissions.write = true;
885 manager.write().await.register_worker(spec);
886
887 // The worker's unfinished work lands after the baseline was captured.
888 fs::create_dir_all(root.join("scratch")).unwrap();
889 fs::write(root.join("scratch/leftover.rs"), "wip\n").unwrap();
890
891 let mut runtime = stub_runtime();
892 runtime.manager = Arc::clone(&manager);
893
894 let note =
895 budget_work_preservation_note(&runtime.manager, "preserve-worker", "wall_time_budget")
896 .await
897 .expect("write-scoped worker has a baseline");
898 assert!(
899 note.contains("scratch/leftover.rs"),
900 "note should name the surviving path: {note}"
901 );
902 assert!(note.contains(&root.display().to_string()), "{note}");
903
904 // A read-only worker captured no baseline — there is no file work to
905 // inventory and the note stays absent rather than lying.
906 let mut scout_spec = make_worker_spec("scout-worker", root.to_path_buf());
907 scout_spec.runtime_profile.permissions.write = false;
908 manager.write().await.register_worker(scout_spec);
909 assert!(
910 budget_work_preservation_note(&runtime.manager, "scout-worker", "wall_time_budget")
911 .await
912 .is_none()
913 );
914
915 // A write-scoped worker that changed nothing still gets an explicit
916 // "no changes" receipt instead of silence.
917 let mut clean_spec = make_worker_spec("clean-worker", root.to_path_buf());
918 clean_spec.runtime_profile.permissions.write = true;
919 clean_spec.workspace = root.to_path_buf();
920 let clean_root = tempdir().unwrap();
921 let clean_path = clean_root.path();
922 git(clean_path, &["init", "--quiet"]);
923 git(clean_path, &["config", "user.name", "Budget test"]);
924 git(
925 clean_path,
926 &["config", "user.email", "budget@example.invalid"],
927 );
928 fs::write(clean_path.join("src.rs"), "baseline\n").unwrap();
929 git(clean_path, &["add", "--", "src.rs"]);
930 git(clean_path, &["commit", "--quiet", "-m", "baseline"]);
931 clean_spec.workspace = clean_path.to_path_buf();
932 manager.write().await.register_worker(clean_spec);
933 let note = budget_work_preservation_note(&runtime.manager, "clean-worker", "wall_time_budget")
934 .await
935 .expect("baseline exists");
936 assert!(note.contains("No workspace changes"), "{note}");
937 }
938
939 fn git_out(root: &Path, args: &[&str]) -> String {
940 let output = std::process::Command::new("git")
941 .arg("-C")
942 .arg(root)
943 .args(args)
944 .output()
945 .expect("git");
946 assert!(
947 output.status.success(),
948 "git {args:?}: {}",
949 String::from_utf8_lossy(&output.stderr)
950 );
951 String::from_utf8_lossy(&output.stdout).trim().to_string()
952 }
953
954 /// #6194 item 4 / #5529: on an isolated worktree a budget death commits the
955 /// worker's uncommitted changes as labeled salvage instead of leaving them
956 /// for manual recovery.
957 #[tokio::test]
958 async fn budget_death_checkpoint_commits_uncommitted_work_on_isolated_worktree() {
959 let tmp = tempdir().unwrap();
960 let root = tmp.path();
961 git(root, &["init", "--quiet"]);
962 git(root, &["config", "user.name", "Budget test"]);
963 git(root, &["config", "user.email", "budget@example.invalid"]);
964 fs::write(root.join("src.rs"), "baseline\n").unwrap();
965 git(root, &["add", "--", "src.rs"]);
966 git(root, &["commit", "--quiet", "-m", "baseline"]);
967
968 let manager = Arc::new(RwLock::new(SubAgentManager::new(root.to_path_buf(), 2)));
969 let mut spec = make_worker_spec("checkpoint-worker", root.to_path_buf());
970 spec.runtime_profile.permissions.write = true;
971 spec.launch_manifest = Some(ChildLaunchManifest {
972 owner_session: "root".to_string(),
973 child_id: "checkpoint-worker".to_string(),
974 profile: spec.runtime_profile.clone(),
975 prompt: spec.objective.clone(),
976 cwd: Some(root.display().to_string()),
977 worktree: true,
978 writable_roots: vec![root.display().to_string()],
979 writable_files: Vec::new(),
980 coordination_contracts: Vec::new(),
981 expected_artifact: None,
982 deliverables: Vec::new(),
983 resume_identity: None,
984 generation: 1,
985 resume_from_agent_id: None,
986 });
987 manager.write().await.register_worker(spec);
988 fs::write(root.join("src.rs"), "baseline\nuncommitted fix\n").unwrap();
989 fs::write(root.join("new.rs"), "wip\n").unwrap();
990
991 let mut runtime = stub_runtime();
992 runtime.manager = Arc::clone(&manager);
993 let note =
994 budget_work_preservation_note(&runtime.manager, "checkpoint-worker", "wall_time_budget")
995 .await
996 .expect("note");
997 assert!(
998 note.contains("checkpointed in commit"),
999 "note should name the salvage commit: {note}"
1000 );
1001 let subject = git_out(root, &["log", "--format=%s", "-1"]);
1002 assert!(
1003 subject.starts_with("checkpoint: checkpoint-worker (wall_time_budget)"),
1004 "marker message names the worker and cause: {subject}"
1005 );
1006 assert!(
1007 git_out(root, &["status", "--porcelain=v1", "--"]).is_empty(),
1008 "checkpoint leaves a clean tree"
1009 );
1010 }
1011
1012 /// A shared checkout may hold the parent's or a sibling's dirty files, so no
1013 /// auto-commit happens there — the note keeps the manual-salvage wording.
1014 #[tokio::test]
1015 async fn budget_death_checkpoint_skips_shared_checkout() {
1016 let tmp = tempdir().unwrap();
1017 let root = tmp.path();
1018 git(root, &["init", "--quiet"]);
1019 git(root, &["config", "user.name", "Budget test"]);
1020 git(root, &["config", "user.email", "budget@example.invalid"]);
1021 fs::write(root.join("src.rs"), "baseline\n").unwrap();
1022 git(root, &["add", "--", "src.rs"]);
1023 git(root, &["commit", "--quiet", "-m", "baseline"]);
1024
1025 let manager = Arc::new(RwLock::new(SubAgentManager::new(root.to_path_buf(), 2)));
1026 let mut spec = make_worker_spec("shared-worker", root.to_path_buf());
1027 spec.runtime_profile.permissions.write = true;
1028 manager.write().await.register_worker(spec);
1029 fs::write(root.join("src.rs"), "baseline\nuncommitted fix\n").unwrap();
1030
1031 let mut runtime = stub_runtime();
1032 runtime.manager = Arc::clone(&manager);
1033 let note = budget_work_preservation_note(&runtime.manager, "shared-worker", "wall_time_budget")
1034 .await
1035 .expect("note");
1036 assert!(!note.contains("checkpointed in commit"), "{note}");
1037 assert!(note.contains("survive on disk"), "{note}");
1038 assert!(
1039 !git_out(root, &["status", "--porcelain=v1", "--"]).is_empty(),
1040 "shared checkout stays dirty"
1041 );
1042 }
1043
1044 /// When the worker committed everything itself before death, the note says so
1045 /// instead of claiming a checkpoint or manual salvage.
1046 #[tokio::test]
1047 async fn budget_death_checkpoint_reports_worker_committed_tree() {
1048 let tmp = tempdir().unwrap();
1049 let root = tmp.path();
1050 git(root, &["init", "--quiet"]);
1051 git(root, &["config", "user.name", "Budget test"]);
1052 git(root, &["config", "user.email", "budget@example.invalid"]);
1053 fs::write(root.join("src.rs"), "baseline\n").unwrap();
1054 git(root, &["add", "--", "src.rs"]);
1055 git(root, &["commit", "--quiet", "-m", "baseline"]);
1056
1057 let manager = Arc::new(RwLock::new(SubAgentManager::new(root.to_path_buf(), 2)));
1058 let mut spec = make_worker_spec("tidy-worker", root.to_path_buf());
1059 spec.runtime_profile.permissions.write = true;
1060 spec.launch_manifest = Some(ChildLaunchManifest {
1061 owner_session: "root".to_string(),
1062 child_id: "tidy-worker".to_string(),
1063 profile: spec.runtime_profile.clone(),
1064 prompt: spec.objective.clone(),
1065 cwd: Some(root.display().to_string()),
1066 worktree: true,
1067 writable_roots: vec![root.display().to_string()],
1068 writable_files: Vec::new(),
1069 coordination_contracts: Vec::new(),
1070 expected_artifact: None,
1071 deliverables: Vec::new(),
1072 resume_identity: None,
1073 generation: 1,
1074 resume_from_agent_id: None,
1075 });
1076 manager.write().await.register_worker(spec);
1077 fs::write(root.join("src.rs"), "baseline\nworker fix\n").unwrap();
1078 git(root, &["add", "--", "src.rs"]);
1079 git(root, &["commit", "--quiet", "-m", "worker fix"]);
1080
1081 let mut runtime = stub_runtime();
1082 runtime.manager = Arc::clone(&manager);
1083 let note = budget_work_preservation_note(&runtime.manager, "tidy-worker", "wall_time_budget")
1084 .await
1085 .expect("note");
1086 assert!(note.contains("committed before death"), "{note}");
1087 }
1088
1089 fn assistant_message(content: Vec<ContentBlock>) -> Message {
1090 Message {
1091 role: Role::Assistant,
1092 content,
1093 }
1094 }
1095
1096 fn tool_use(name: &str, input: Value) -> ContentBlock {
1097 ContentBlock::ToolUse {
1098 execution_id: None,
1099 id: format!("call_{name}"),
1100 name: name.to_string(),
1101 input,
1102 caller: None,
1103 thought_signature: None,
1104 }
1105 }
1106
1107 #[test]
1108 fn fallback_partial_text_prefers_last_assistant_text() {
1109 let messages = vec![
1110 assistant_message(vec![ContentBlock::Text {
1111 text: "first".to_string(),
1112 cache_control: None,
1113 }]),
1114 assistant_message(vec![
1115 tool_use("Read", json!({"path": "src/main.rs"})),
1116 ContentBlock::Text {
1117 text: "second".to_string(),
1118 cache_control: None,
1119 },
1120 ]),
1121 ];
1122 assert_eq!(budget_handback::fallback_partial_text(&messages), "second");
1123 }
1124
1125 #[test]
1126 fn fallback_partial_text_digests_thinking_and_tool_calls_without_text() {
1127 let messages = vec![
1128 assistant_message(vec![ContentBlock::Thinking {
1129 thinking: "checking whether the ring slot write precedes the read".to_string(),
1130 signature: None,
1131 state: None,
1132 }]),
1133 assistant_message(vec![tool_use("Read", json!({"path": "ring.rs"}))]),
1134 assistant_message(vec![tool_use(
1135 "Grep",
1136 json!({"pattern": "slot", "path": "ring.rs"}),
1137 )]),
1138 ];
1139 let digest = budget_handback::fallback_partial_text(&messages);
1140 assert!(digest.contains("Tool calls (newest first)"), "{digest}");
1141 assert!(digest.contains("- Grep ring.rs"), "{digest}");
1142 assert!(digest.contains("- Read ring.rs"), "{digest}");
1143 assert!(
1144 digest.find("- Grep").unwrap() < digest.find("- Read").unwrap(),
1145 "{digest}"
1146 );
1147 assert!(digest.contains("unverified"), "{digest}");
1148 assert!(
1149 digest.contains("ring slot write precedes the read"),
1150 "{digest}"
1151 );
1152 }
1153
1154 #[test]
1155 fn fallback_partial_text_caps_tool_entries_and_reports_overflow() {
1156 let messages: Vec<Message> = (0..14)
1157 .map(|i| {
1158 assistant_message(vec![tool_use(
1159 "Read",
1160 json!({"path": format!("file_{i}.rs")}),
1161 )])
1162 })
1163 .collect();
1164 let digest = budget_handback::fallback_partial_text(&messages);
1165 assert!(digest.contains("...and 2 more"), "{digest}");
1166 assert!(!digest.contains("file_0.rs"), "{digest}");
1167 assert!(digest.contains("file_13.rs"), "{digest}");
1168 }
1169
1170 #[test]
1171 fn fallback_partial_text_is_silent_only_when_nothing_was_recorded() {
1172 assert!(budget_handback::fallback_partial_text(&[]).contains("No assistant text was recorded"));
1173 let user_only = vec![Message {
1174 role: Role::User,
1175 content: vec![ContentBlock::Text {
1176 text: "do the thing".to_string(),
1177 cache_control: None,
1178 }],
1179 }];
1180 assert!(
1181 budget_handback::fallback_partial_text(&user_only)
1182 .contains("No assistant text was recorded")
1183 );
1184 }
1185
1186 #[test]
1187 fn budget_repair_only_rewrites_the_newly_synthesized_final_execution() {
1188 for (first, last) in [(Some("first"), Some("last")), (None, None)] {
1189 let mut messages: Vec<Message> = serde_json::from_value(json!([
1190 {"role":"assistant","content":[{"type":"tool_use","id":"reused","execution_id":first,"name":"read","input":{}}]},
1191 {"role":"user","content":[{"type":"tool_result","tool_use_id":"reused","execution_id":first,"content":"earlier completed output"}]},
1192 {"role":"assistant","content":[{"type":"tool_use","id":"reused","execution_id":last,"name":"read","input":{}}]}
1193 ])).unwrap();
1194 let earlier = messages[..2].to_vec();
1195 budget_handback::repair_stopped_tool_calls(&mut messages, "fixture deadline");
1196 assert_eq!(&messages[..2], earlier.as_slice());
1197 assert!(
1198 matches!(&messages[3].content[0], ContentBlock::ToolResult { execution_id, content, .. }
1199 if execution_id.as_deref() == last && content.contains("budget_exhausted"))
1200 );
1201 }
1202 }
1203
1203 lines RUST