| 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 |