返回 CodeWhale
exec_turn_usage.rs
根目录 / crates / tui / tests / integration / exec_turn_usage.rs
1 //! End-to-end shape lock for the per-model-call `turn_usage` event on the
2 //! `codewhale exec --output-format stream-json` stream (#52 / FINISH-0.9.4).
3 //!
4 //! A `wiremock` OpenAI-compatible endpoint stands in for the provider. Two
5 //! cases pin the contract:
6 //!
7 //! - usage reported by the provider -> exactly one `turn_usage` event per
8 //! model call, carrying the reported input/output/reasoning/cache fields,
9 //! and the pre-existing event sequence (`content` … `metadata` → `done`)
10 //! is unchanged for existing consumers;
11 //! - usage absent from the provider stream -> no `turn_usage` event at all
12 //! (honest absence, never fabricated zeros-as-data).
13
14 #![cfg(unix)]
15
16 use std::io::Read;
17 use std::process::{Command, Stdio};
18 use std::time::Duration;
19
20 use serde_json::{Value, json};
21 use tempfile::TempDir;
22 use wait_timeout::ChildExt;
23 use wiremock::matchers::{method, path};
24 use wiremock::{Mock, MockServer, ResponseTemplate};
25
26 const TEST_MODEL: &str = "turn-usage-model";
27 const RUN_TIMEOUT: Duration = Duration::from_secs(60);
28
29 fn sse_chunk(value: Value) -> String {
30 format!(
31 "data: {}\n\n",
32 serde_json::to_string(&value).expect("SSE JSON")
33 )
34 }
35
36 /// Final-answer SSE whose closing chunk reports usage with reasoning and
37 /// DeepSeek-style prompt-cache fields.
38 fn answer_sse_with_usage(answer: &str) -> String {
39 [
40 sse_chunk(json!({
41 "id": "chatcmpl-usage",
42 "object": "chat.completion.chunk",
43 "model": TEST_MODEL,
44 "choices": [{"index": 0, "delta": {"content": answer}, "finish_reason": null}]
45 })),
46 sse_chunk(json!({
47 "id": "chatcmpl-usage",
48 "object": "chat.completion.chunk",
49 "model": TEST_MODEL,
50 "choices": [{"index": 0, "delta": {}, "finish_reason": "stop"}],
51 "usage": {
52 "prompt_tokens": 20,
53 "completion_tokens": 8,
54 "total_tokens": 28,
55 "completion_tokens_details": {"reasoning_tokens": 5},
56 "prompt_cache_hit_tokens": 12,
57 "prompt_cache_miss_tokens": 8
58 }
59 })),
60 "data: [DONE]\n\n".to_string(),
61 ]
62 .join("")
63 }
64
65 /// Final-answer SSE whose provider never reports usage.
66 fn answer_sse_without_usage(answer: &str) -> String {
67 [
68 sse_chunk(json!({
69 "id": "chatcmpl-no-usage",
70 "object": "chat.completion.chunk",
71 "model": TEST_MODEL,
72 "choices": [{"index": 0, "delta": {"content": answer}, "finish_reason": null}]
73 })),
74 sse_chunk(json!({
75 "id": "chatcmpl-no-usage",
76 "object": "chat.completion.chunk",
77 "model": TEST_MODEL,
78 "choices": [{"index": 0, "delta": {}, "finish_reason": "stop"}]
79 })),
80 "data: [DONE]\n\n".to_string(),
81 ]
82 .join("")
83 }
84
85 fn sse_response(body: String) -> ResponseTemplate {
86 ResponseTemplate::new(200)
87 .insert_header("content-type", "text/event-stream")
88 .insert_header("cache-control", "no-cache")
89 .set_body_string(body)
90 }
91
92 fn json_response(value: Value) -> ResponseTemplate {
93 ResponseTemplate::new(200)
94 .insert_header("content-type", "application/json")
95 .set_body_json(value)
96 }
97
98 async fn start_mock_llm(answer_sse: String) -> MockServer {
99 let server = MockServer::start().await;
100
101 Mock::given(method("GET"))
102 .and(path("/v1/models"))
103 .respond_with(json_response(json!({
104 "object": "list",
105 "data": [{ "id": TEST_MODEL, "object": "model" }]
106 })))
107 .mount(&server)
108 .await;
109
110 Mock::given(method("POST"))
111 .and(path("/v1/chat/completions"))
112 .respond_with(sse_response(answer_sse))
113 .mount(&server)
114 .await;
115
116 server
117 }
118
119 fn preserve_host_env(command: &mut Command) {
120 command.env_clear();
121 for key in [
122 "PATH",
123 "PATHEXT",
124 "SystemRoot",
125 "SystemDrive",
126 "WINDIR",
127 "COMSPEC",
128 "TEMP",
129 "TMP",
130 "TERM",
131 "COLORTERM",
132 "LANG",
133 "LC_ALL",
134 ] {
135 if let Some(value) = std::env::var_os(key) {
136 command.env(key, value);
137 }
138 }
139 }
140
141 fn run_exec_stream_json(server: &MockServer) -> Vec<Value> {
142 let stdout = run_exec(
143 server,
144 &[
145 "--auto",
146 "--model",
147 TEST_MODEL,
148 "--output-format",
149 "stream-json",
150 "answer briefly",
151 ],
152 );
153 stdout
154 .lines()
155 .filter(|line| !line.trim().is_empty())
156 .map(|line| {
157 serde_json::from_str(line).unwrap_or_else(|err| {
158 panic!("stream-json line should parse: {err}\nline: {line}\nstdout:\n{stdout}")
159 })
160 })
161 .collect()
162 }
163
164 /// Run `codewhale exec <exec_args>` against `server` and return stdout.
165 fn run_exec(server: &MockServer, exec_args: &[&str]) -> String {
166 let (success, stdout, stderr) = run_exec_unchecked(server, exec_args);
167 assert!(
168 success,
169 "codewhale exec failed\nstdout:\n{stdout}\nstderr:\n{stderr}"
170 );
171 stdout
172 }
173
174 /// Run `codewhale exec <exec_args>` and return whether it exited
175 /// successfully, its stdout and its stderr.
176 fn run_exec_unchecked(server: &MockServer, exec_args: &[&str]) -> (bool, String, String) {
177 run_exec_in_home(server, exec_args, |_| {})
178 }
179
180 /// [`run_exec_unchecked`] with a hook that prepares the isolated `$HOME`
181 /// before the run, e.g. to install an MCP config.
182 fn run_exec_in_home(
183 server: &MockServer,
184 exec_args: &[&str],
185 prepare_home: impl FnOnce(&std::path::Path),
186 ) -> (bool, String, String) {
187 run_exec_with_stdin(server, exec_args, prepare_home, None)
188 }
189
190 /// [`run_exec_in_home`] that optionally pipes `stdin` into the child.
191 fn run_exec_with_stdin(
192 server: &MockServer,
193 exec_args: &[&str],
194 prepare_home: impl FnOnce(&std::path::Path),
195 stdin: Option<Vec<u8>>,
196 ) -> (bool, String, String) {
197 let workspace = TempDir::new().expect("workspace tempdir");
198 let home = TempDir::new().expect("home tempdir");
199
200 let mut command = Command::new(crate::binary::codewhale());
201 preserve_host_env(&mut command);
202 command
203 .current_dir(workspace.path())
204 .arg("--workspace")
205 .arg(workspace.path())
206 .arg("--no-project-config")
207 .arg("exec")
208 .args(exec_args)
209 .env("HOME", home.path())
210 .env("USERPROFILE", home.path())
211 .env("XDG_CONFIG_HOME", home.path().join(".config"))
212 .env("XDG_DATA_HOME", home.path().join(".local").join("share"))
213 .env("XDG_CACHE_HOME", home.path().join(".cache"))
214 .env(
215 "CODEWHALE_CONFIG_PATH",
216 home.path().join(".codewhale").join("config.toml"),
217 )
218 .env(
219 "DEEPSEEK_CONFIG_PATH",
220 home.path().join(".deepseek").join("config.toml"),
221 )
222 .env("DEEPSEEK_API_KEY", "ci-test-key-not-real")
223 .env("DEEPSEEK_BASE_URL", server.uri())
224 .env("CODEWHALE_BASE_URL", server.uri())
225 .env("DEEPSEEK_MODEL", TEST_MODEL)
226 .env("CODEWHALE_MODEL", TEST_MODEL)
227 .env("RUST_LOG", "warn")
228 .stdout(Stdio::piped())
229 .stderr(Stdio::piped());
230
231 std::fs::create_dir_all(home.path().join(".codewhale")).expect("create codewhale config dir");
232 std::fs::create_dir_all(home.path().join(".deepseek")).expect("create deepseek config dir");
233 prepare_home(home.path());
234
235 if stdin.is_some() {
236 command.stdin(Stdio::piped());
237 }
238 let mut child = command.spawn().expect("spawn codewhale exec");
239 let stdin_writer = stdin.map(|bytes| {
240 let mut pipe = child.stdin.take().expect("stdin pipe");
241 std::thread::spawn(move || {
242 use std::io::Write;
243 // The child may exit without reading; a broken pipe is fine.
244 let _ = pipe.write_all(&bytes);
245 })
246 });
247 let stdout_reader = read_pipe_in_background(child.stdout.take().expect("stdout pipe"));
248 let stderr_reader = read_pipe_in_background(child.stderr.take().expect("stderr pipe"));
249
250 let status = match child.wait_timeout(RUN_TIMEOUT).expect("wait for codewhale") {
251 Some(status) => status,
252 None => {
253 let _ = child.kill();
254 let _ = child.wait();
255 let stdout = join_pipe_reader(stdout_reader, "stdout");
256 let stderr = join_pipe_reader(stderr_reader, "stderr");
257 panic!(
258 "codewhale exec timed out after {RUN_TIMEOUT:?}\nstdout:\n{}\nstderr:\n{}",
259 String::from_utf8_lossy(&stdout),
260 String::from_utf8_lossy(&stderr)
261 );
262 }
263 };
264
265 if let Some(writer) = stdin_writer {
266 writer.join().expect("stdin writer thread");
267 }
268 let stdout = join_pipe_reader(stdout_reader, "stdout");
269 let stderr = join_pipe_reader(stderr_reader, "stderr");
270 (
271 status.success(),
272 String::from_utf8_lossy(&stdout).into_owned(),
273 String::from_utf8_lossy(&stderr).into_owned(),
274 )
275 }
276
277 fn read_pipe_in_background<R>(mut reader: R) -> std::thread::JoinHandle<std::io::Result<Vec<u8>>>
278 where
279 R: Read + Send + 'static,
280 {
281 std::thread::spawn(move || {
282 let mut output = Vec::new();
283 reader.read_to_end(&mut output).map(|_| output)
284 })
285 }
286
287 fn join_pipe_reader(
288 handle: std::thread::JoinHandle<std::io::Result<Vec<u8>>>,
289 stream_name: &str,
290 ) -> Vec<u8> {
291 handle
292 .join()
293 .unwrap_or_else(|_| panic!("{stream_name} reader thread panicked"))
294 .unwrap_or_else(|err| panic!("failed to read {stream_name}: {err}"))
295 }
296
297 fn events_of_type<'a>(events: &'a [Value], event_type: &str) -> Vec<&'a Value> {
298 events
299 .iter()
300 .filter(|event| event.get("type").and_then(Value::as_str) == Some(event_type))
301 .collect()
302 }
303
304 #[tokio::test(flavor = "multi_thread")]
305 async fn turn_usage_event_is_emitted_with_reported_fields_and_stream_contract_holds() {
306 let server = start_mock_llm(answer_sse_with_usage("done in one step")).await;
307 let events = run_exec_stream_json(&server);
308
309 // Every event carries the stream schema envelope.
310 for event in &events {
311 assert_eq!(event["schema"], "codewhale.exec-stream");
312 assert_eq!(event["schema_version"], 1);
313 }
314
315 // Exactly one per-call usage receipt, numbered from 1.
316 let usage_events = events_of_type(&events, "turn_usage");
317 assert_eq!(
318 usage_events.len(),
319 1,
320 "expected one turn_usage event: {events:#?}"
321 );
322 let usage = usage_events[0];
323 assert_eq!(usage["turn"], 1);
324 assert_eq!(usage["input_tokens"], 20);
325 assert_eq!(usage["output_tokens"], 8);
326 assert_eq!(usage["reasoning_tokens"], 5);
327 assert_eq!(usage["prompt_cache_hit_tokens"], 12);
328 assert_eq!(usage["prompt_cache_miss_tokens"], 8);
329 assert!(
330 usage["duration_ms"].as_u64().is_some(),
331 "duration_ms must be a non-negative integer: {usage}"
332 );
333 // Fields the provider did not report are omitted, not zero-filled.
334 let usage_object = usage.as_object().expect("turn_usage object");
335 for absent in ["prompt_cache_write_tokens", "reasoning_replay_tokens"] {
336 assert!(
337 !usage_object.contains_key(absent),
338 "{absent} must be omitted when unreported: {usage}"
339 );
340 }
341
342 // The usage receipt lands after the model output it accounts for and
343 // before the terminal receipts.
344 let types: Vec<&str> = events
345 .iter()
346 .filter_map(|event| event.get("type").and_then(Value::as_str))
347 .collect();
348 let content_pos = types.iter().position(|t| *t == "content");
349 let usage_pos = types.iter().position(|t| *t == "turn_usage");
350 assert!(
351 content_pos.is_some_and(|c| usage_pos.is_some_and(|u| c < u)),
352 "turn_usage must follow the content it accounts for: {types:?}"
353 );
354
355 // Existing consumers' terminal contract is unchanged: `metadata`
356 // immediately precedes exactly one trailing `done`.
357 assert_eq!(types.last(), Some(&"done"), "stream must end with done");
358 assert_eq!(
359 types.get(types.len() - 2),
360 Some(&"metadata"),
361 "metadata must immediately precede done: {types:?}"
362 );
363 assert_eq!(
364 events_of_type(&events, "done").len(),
365 1,
366 "exactly one done event"
367 );
368 let metadata = events_of_type(&events, "metadata");
369 assert_eq!(metadata.len(), 1, "exactly one metadata event");
370 // The terminal receipt still carries the cumulative usage.
371 assert_eq!(metadata[0]["meta"]["input_tokens"], 20);
372 assert_eq!(metadata[0]["meta"]["output_tokens"], 8);
373 assert_eq!(metadata[0]["meta"]["reasoning_tokens"], 5);
374 }
375
376 #[tokio::test(flavor = "multi_thread")]
377 async fn turn_usage_event_is_skipped_when_provider_reports_no_usage() {
378 let server = start_mock_llm(answer_sse_without_usage("quiet answer")).await;
379 let events = run_exec_stream_json(&server);
380
381 assert!(
382 events_of_type(&events, "turn_usage").is_empty(),
383 "no turn_usage event without provider-reported usage: {events:#?}"
384 );
385
386 // The rest of the stream contract still holds.
387 let types: Vec<&str> = events
388 .iter()
389 .filter_map(|event| event.get("type").and_then(Value::as_str))
390 .collect();
391 assert!(types.contains(&"content"), "content missing: {types:?}");
392 assert_eq!(types.last(), Some(&"done"), "stream must end with done");
393 assert_eq!(types.get(types.len() - 2), Some(&"metadata"));
394 }
395
396 /// #6510: plain `exec` bypassed the Engine — text mode sent no system prompt,
397 /// `--json` sent an inline "coding assistant" line — so the output format
398 /// changed the model's instructions and nothing was logged. Both now run one
399 /// Engine turn with the one base prompt and no tool catalog, and `--json`
400 /// keeps its documented one-shot receipt fields.
401 #[tokio::test(flavor = "multi_thread")]
402 async fn plain_exec_runs_one_engine_turn_under_one_prompt_authority() {
403 let server = start_mock_llm(answer_sse_with_usage("pong")).await;
404
405 let text = run_exec(&server, &["--model", TEST_MODEL, "answer briefly"]);
406 assert_eq!(text.trim(), "pong", "text mode prints the answer: {text}");
407
408 let json_stdout = run_exec(
409 &server,
410 &["--json", "--model", TEST_MODEL, "answer briefly"],
411 );
412 let receipt: Value = serde_json::from_str(&json_stdout)
413 .unwrap_or_else(|err| panic!("--json receipt should parse: {err}\n{json_stdout}"));
414 assert_eq!(receipt["mode"], "one-shot");
415 assert_eq!(receipt["model"], TEST_MODEL);
416 assert_eq!(receipt["success"], true);
417 assert_eq!(receipt["output"], "pong");
418 assert_eq!(receipt["usage"]["input_tokens"], 20);
419 assert_eq!(receipt["usage"]["output_tokens"], 8);
420 assert!(
421 receipt["tools"].as_array().is_some_and(Vec::is_empty),
422 "{receipt}"
423 );
424
425 let bodies = chat_bodies(&server).await;
426 assert_eq!(bodies.len(), 2, "one model call per run: {bodies:#?}");
427 for body in &bodies {
428 let messages = body["messages"].to_string();
429 assert!(
430 messages.contains("You are Codewhale, an agent working alongside the user"),
431 "both formats must carry the one base prompt: {messages}"
432 );
433 assert!(
434 !messages.contains("You are a coding assistant. Give concise"),
435 "the old --json-only instruction must be gone: {messages}"
436 );
437 assert!(
438 body.get("tools")
439 .is_none_or(|tools| tools.as_array().is_some_and(Vec::is_empty)),
440 "plain exec offers no tools: {body}"
441 );
442 }
443 }
444
445 /// Chat-completions bodies the mock received, in order.
446 async fn chat_bodies(server: &MockServer) -> Vec<Value> {
447 server
448 .received_requests()
449 .await
450 .expect("request recording")
451 .into_iter()
452 .filter(|request| request.url.path() == "/v1/chat/completions")
453 .map(|request| serde_json::from_slice(&request.body).expect("request body JSON"))
454 .collect()
455 }
456
457 /// #6510: a limit is not a tool grant. `--max-turns`, `--disallowed-tools`
458 /// and `--append-system-prompt` used to put plain exec on the full tool
459 /// catalog, so `exec --max-turns 1 "hi"` became a tool-using agent.
460 #[tokio::test(flavor = "multi_thread")]
461 async fn plain_exec_limits_do_not_grant_tools() {
462 let server = start_mock_llm(answer_sse_with_usage("pong")).await;
463
464 for flags in [
465 &["--max-turns", "1"][..],
466 &["--disallowed-tools", "exec_shell"],
467 &["--append-system-prompt", "Be brief."],
468 ] {
469 let mut args = flags.to_vec();
470 args.extend(["--model", TEST_MODEL, "answer briefly"]);
471 let text = run_exec(&server, &args);
472 assert_eq!(text.trim(), "pong", "{flags:?}: {text}");
473 }
474
475 let bodies = chat_bodies(&server).await;
476 assert_eq!(bodies.len(), 3, "one model call per run: {bodies:#?}");
477 for body in &bodies {
478 assert!(
479 body.get("tools")
480 .is_none_or(|tools| tools.as_array().is_some_and(Vec::is_empty)),
481 "a limit flag must not offer tools: {body}"
482 );
483 }
484 }
485
486 /// #6510: plain exec now runs on the Engine, so an output-limit stop follows
487 /// the Engine's one policy: the partial answer is kept, the model is asked to
488 /// continue, and the run succeeds with the whole answer. The old direct call
489 /// printed the partial answer and failed.
490 #[tokio::test(flavor = "multi_thread")]
491 async fn plain_exec_continues_past_an_output_limit_stop() {
492 let server = MockServer::start().await;
493 Mock::given(method("GET"))
494 .and(path("/v1/models"))
495 .respond_with(json_response(json!({
496 "object": "list",
497 "data": [{ "id": TEST_MODEL, "object": "model" }]
498 })))
499 .mount(&server)
500 .await;
501 let truncated = [
502 sse_chunk(json!({
503 "id": "chatcmpl-cut",
504 "object": "chat.completion.chunk",
505 "model": TEST_MODEL,
506 "choices": [{"index": 0, "delta": {"content": "first half"}, "finish_reason": null}]
507 })),
508 sse_chunk(json!({
509 "id": "chatcmpl-cut",
510 "object": "chat.completion.chunk",
511 "model": TEST_MODEL,
512 "choices": [{"index": 0, "delta": {}, "finish_reason": "length"}]
513 })),
514 "data: [DONE]\n\n".to_string(),
515 ]
516 .join("");
517 Mock::given(method("POST"))
518 .and(path("/v1/chat/completions"))
519 .respond_with(sse_response(truncated))
520 .up_to_n_times(1)
521 .with_priority(1)
522 .mount(&server)
523 .await;
524 Mock::given(method("POST"))
525 .and(path("/v1/chat/completions"))
526 .respond_with(sse_response(answer_sse_without_usage(" second half")))
527 .with_priority(2)
528 .mount(&server)
529 .await;
530
531 let json_stdout = run_exec(
532 &server,
533 &["--json", "--model", TEST_MODEL, "answer briefly"],
534 );
535 let receipt: Value = serde_json::from_str(&json_stdout)
536 .unwrap_or_else(|err| panic!("--json receipt should parse: {err}\n{json_stdout}"));
537 assert_eq!(receipt["mode"], "one-shot", "{receipt}");
538 assert_eq!(receipt["success"], true, "{receipt}");
539 let output = receipt["output"].as_str().unwrap_or_default();
540 assert!(
541 output.contains("first half") && output.contains("second half"),
542 "{receipt}"
543 );
544
545 let bodies = chat_bodies(&server).await;
546 assert_eq!(bodies.len(), 2, "one continuation request: {bodies:#?}");
547 assert!(
548 bodies[1]["messages"]
549 .to_string()
550 .contains("stopped generation at its output limit"),
551 "the continuation names the truncation: {}",
552 bodies[1]["messages"]
553 );
554 }
555
556 /// #6510 review: a model that keeps stopping at its output limit must not be
557 /// re-asked until the turn wall clock. Plain exec caps its continuations at a
558 /// small default (8 model steps) unless `--max-turns` sets another, ends the
559 /// run as failed with the partial answer, and never injects the agent wrap-up
560 /// notices (soft landing, final report) into a one-shot answer.
561 #[tokio::test(flavor = "multi_thread")]
562 async fn plain_exec_bounds_output_limit_continuations() {
563 let always_truncated = [
564 sse_chunk(json!({
565 "id": "chatcmpl-cut",
566 "object": "chat.completion.chunk",
567 "model": TEST_MODEL,
568 "choices": [{"index": 0, "delta": {"content": "again"}, "finish_reason": null}]
569 })),
570 sse_chunk(json!({
571 "id": "chatcmpl-cut",
572 "object": "chat.completion.chunk",
573 "model": TEST_MODEL,
574 "choices": [{"index": 0, "delta": {}, "finish_reason": "length"}]
575 })),
576 "data: [DONE]\n\n".to_string(),
577 ]
578 .join("");
579
580 for (flags, expected_requests) in [(&[][..], 8usize), (&["--max-turns", "2"][..], 2)] {
581 let server = start_mock_llm(always_truncated.clone()).await;
582 let mut args = flags.to_vec();
583 args.extend(["--json", "--model", TEST_MODEL, "answer briefly"]);
584 let (success, stdout, stderr) = run_exec_unchecked(&server, &args);
585 assert!(
586 !success,
587 "{flags:?}: a capped run must not exit 0\n{stderr}"
588 );
589 let receipt: Value = serde_json::from_str(&stdout)
590 .unwrap_or_else(|err| panic!("--json receipt should parse: {err}\n{stdout}"));
591 assert_eq!(receipt["mode"], "one-shot", "{receipt}");
592 assert_eq!(receipt["success"], false, "{receipt}");
593 assert!(
594 receipt["output"]
595 .as_str()
596 .is_some_and(|output| output.contains("again")),
597 "the partial answer is kept: {receipt}"
598 );
599
600 let bodies = chat_bodies(&server).await;
601 assert_eq!(bodies.len(), expected_requests, "{flags:?}: {bodies:#?}");
602 for body in &bodies {
603 let messages = body["messages"].to_string();
604 assert!(
605 !messages.contains("Step budget soft landing")
606 && !messages.contains("Write your final report now"),
607 "{flags:?}: no agent wrap-up notice in a one-shot answer: {messages}"
608 );
609 }
610 }
611 }
612
613 /// A stdio MCP server that reads `initialize`, closes its stdin, answers, and
614 /// stays alive: the client's next write (`notifications/initialized`) always
615 /// hits a pipe with no reader. That is the shape of the real failure, where a
616 /// plugin's MCP server whose `node` could not start made `exec --auto` die of
617 /// SIGPIPE (exit 141) before printing anything.
618 const STDIN_CLOSING_MCP_SERVER: &str = r#"IFS= read -r _line; exec 0<&-; printf '%s\n' '{"jsonrpc":"2.0","id":"1","result":{"protocolVersion":"2024-11-05","serverInfo":{"name":"stdin-closer","version":"1.0.0"},"capabilities":{"tools":{}}}}'; sleep 10"#;
619
620 #[tokio::test(flavor = "multi_thread")]
621 async fn exec_auto_survives_an_mcp_server_that_closes_its_stdin() {
622 let server = start_mock_llm(answer_sse_without_usage("ok")).await;
623 let (success, stdout, stderr) = run_exec_in_home(
624 &server,
625 &["--auto", "--model", TEST_MODEL, "Reply with exactly: ok"],
626 |home| {
627 let config = json!({
628 "mcpServers": {
629 "stdin-closer": {
630 "command": "sh",
631 "args": ["-c", STDIN_CLOSING_MCP_SERVER],
632 // Connect at session start, as the Computer Use
633 // plugin's server does, instead of on first use.
634 "required": true
635 }
636 }
637 });
638 let path = home.join(".codewhale").join("mcp.json");
639 std::fs::write(&path, config.to_string()).expect("write mcp.json");
640 },
641 );
642 assert!(
643 success,
644 "exec --auto must not die when an MCP peer closes its pipe\nstdout:\n{stdout}\nstderr:\n{stderr}"
645 );
646 assert_eq!(stdout.trim(), "ok", "stderr:\n{stderr}");
647 }
648
649 /// Plain `exec` offers no tools, yet DeepSeek can still answer with nothing
650 /// but a DSML tool call. The markup must never reach stdout, and the run must
651 /// fail once with the real reason instead of re-requesting the same call and
652 /// ending on "the provider response was incomplete".
653 #[tokio::test(flavor = "multi_thread")]
654 async fn one_shot_exec_strips_deepseek_dsml_and_points_at_auto() {
655 let dsml = "<||DSML|| calls>\n<||DSML|| invoke name=\"read_file\">\n<||DSML|| parameter name=\"path\" string=\"true\">note.txt</||DSML|| parameter>\n</||DSML|| invoke>\n</||DSML|| calls>\n";
656 let server = start_mock_llm(answer_sse_without_usage(dsml)).await;
657 let (success, stdout, stderr) = run_exec_unchecked(
658 &server,
659 &[
660 "--model",
661 TEST_MODEL,
662 "Read note.txt and tell me its contents.",
663 ],
664 );
665 assert!(
666 !success,
667 "a tool-call-only answer is not an answer\nstdout:\n{stdout}\nstderr:\n{stderr}"
668 );
669 assert!(
670 !stdout.contains("DSML") && !stdout.contains("read_file"),
671 "raw tool-call markup reached stdout: {stdout:?}"
672 );
673 assert!(
674 stderr.contains("--auto") && stderr.contains("offers no tools"),
675 "the user is told the task needs tools: {stderr}"
676 );
677 assert_eq!(
678 chat_bodies(&server).await.len(),
679 1,
680 "a zero-tool text call must not be re-requested\nstderr:\n{stderr}"
681 );
682 }
683
684 /// Last user message text of a chat-completions body, without the
685 /// `<turn_meta>` block the engine appends.
686 fn last_user_text(body: &Value) -> String {
687 let message = body["messages"]
688 .as_array()
689 .expect("messages array")
690 .iter()
691 .rev()
692 .find(|message| message["role"] == "user")
693 .expect("a user message");
694 let text = match &message["content"] {
695 Value::String(text) => text.clone(),
696 other => other.to_string(),
697 };
698 match text.split_once("\n<turn_meta>") {
699 Some((prompt, _)) => prompt.to_string(),
700 None => text,
701 }
702 }
703
704 /// #6688: the prompt the model receives comes from `--prompt-file <PATH>` or
705 /// `--prompt-file -` (stdin), past argv's 128 KiB per-argument ceiling, while
706 /// a positional `-` stays literal text (cloud dispatch passes job prompts
707 /// verbatim as argv) and never reads stdin.
708 #[tokio::test(flavor = "multi_thread")]
709 async fn exec_prompt_reaches_the_model_from_prompt_file_and_stdin() {
710 let server = start_mock_llm(answer_sse_with_usage("pong")).await;
711 let dir = TempDir::new().expect("prompt tempdir");
712 let file_body = format!("FILE-PROMPT {}", "x".repeat(200 * 1024));
713 let file = dir.path().join("prompt.txt");
714 std::fs::write(&file, &file_body).expect("write prompt file");
715 let file_arg = file.to_str().expect("utf-8 path");
716
717 let (ok, stdout, stderr) = run_exec_with_stdin(
718 &server,
719 &["--model", TEST_MODEL, "--prompt-file", file_arg],
720 |_| {},
721 None,
722 );
723 assert!(
724 ok,
725 "--prompt-file <PATH>\nstdout:\n{stdout}\nstderr:\n{stderr}"
726 );
727 let (ok, stdout, stderr) = run_exec_with_stdin(
728 &server,
729 &["--model", TEST_MODEL, "--prompt-file", "-"],
730 |_| {},
731 Some(b"STDIN-PROMPT answer briefly".to_vec()),
732 );
733 assert!(ok, "--prompt-file -\nstdout:\n{stdout}\nstderr:\n{stderr}");
734 let (ok, stdout, stderr) = run_exec_with_stdin(
735 &server,
736 &["--model", TEST_MODEL, "-"],
737 |_| {},
738 Some(b"STDIN-MUST-NOT-BE-READ".to_vec()),
739 );
740 assert!(ok, "positional -\nstdout:\n{stdout}\nstderr:\n{stderr}");
741
742 let bodies = chat_bodies(&server).await;
743 assert_eq!(bodies.len(), 3, "one model call per run");
744 assert_eq!(last_user_text(&bodies[0]), file_body);
745 assert_eq!(last_user_text(&bodies[1]), "STDIN-PROMPT answer briefly");
746 assert_eq!(last_user_text(&bodies[2]), "-");
747 }
748
748 lines RUST