| 1 | //! Process-level acceptance for adaptive exact-evidence routing (#4619). |
| 2 | //! |
| 3 | //! "Exact" begins at the common engine routing seam: tool adapters such as |
| 4 | //! Bash intentionally bound their own operating-system stream and annotate |
| 5 | //! that truncation before returning a `ToolResult`. Adaptive evidence binds |
| 6 | //! every byte of that returned result. Root streaming, sequential/deferred |
| 7 | //! completion, and MCP all converge on the same engine seam; sub-agents have a |
| 8 | //! separate call site covered by `tools::subagent::tests`. |
| 9 | |
| 10 | use std::io::Read; |
| 11 | use std::path::{Path, PathBuf}; |
| 12 | use std::process::{Command, Stdio}; |
| 13 | use std::sync::Arc; |
| 14 | use std::sync::atomic::{AtomicUsize, Ordering}; |
| 15 | use std::time::Duration; |
| 16 | |
| 17 | use serde_json::{Value, json}; |
| 18 | use sha2::{Digest, Sha256}; |
| 19 | use tempfile::TempDir; |
| 20 | use wait_timeout::ChildExt; |
| 21 | use wiremock::matchers::{method, path}; |
| 22 | use wiremock::{Mock, MockServer, Request, Respond, ResponseTemplate}; |
| 23 | |
| 24 | const MODEL: &str = "adaptive-evidence-test"; |
| 25 | const SUCCESS_CALL_ID: &str = "call_bash_success"; |
| 26 | const FAILURE_CALL_ID: &str = "call_bash_failure"; |
| 27 | const RETRIEVE_CALL_ID: &str = "call_retrieve_omitted_range"; |
| 28 | const SUCCESS_SENTINEL: &str = "DEEP_SUCCESS_EVIDENCE_SENTINEL_4619"; |
| 29 | const FAILURE_SENTINEL: &str = "DEEP_FAILURE_EVIDENCE_SENTINEL_4619"; |
| 30 | |
| 31 | #[tokio::test(flavor = "multi_thread", worker_threads = 2)] |
| 32 | async fn headless_bash_success_and_failure_are_distinct_bounded_exact_evidence() { |
| 33 | let workspace = TempDir::new().expect("workspace"); |
| 34 | let home = TempDir::new().expect("home"); |
| 35 | |
| 36 | let server = mock_llm().await; |
| 37 | let output = run_exec(workspace.path(), home.path(), &server); |
| 38 | assert!( |
| 39 | output.status.success(), |
| 40 | "exec failed\nstdout:\n{}\nstderr:\n{}", |
| 41 | String::from_utf8_lossy(&output.stdout), |
| 42 | String::from_utf8_lossy(&output.stderr) |
| 43 | ); |
| 44 | let events: Vec<Value> = std::str::from_utf8(&output.stdout) |
| 45 | .expect("UTF-8 stream events") |
| 46 | .lines() |
| 47 | .filter(|line| !line.trim().is_empty()) |
| 48 | .map(|line| serde_json::from_str(line).expect("valid stream event")) |
| 49 | .collect(); |
| 50 | |
| 51 | let requests = server.received_requests().await.expect("recorded requests"); |
| 52 | let success_receipt = |
| 53 | receipt_for(&requests, SUCCESS_CALL_ID).expect("model-visible Bash success receipt"); |
| 54 | let failure_receipt = |
| 55 | receipt_for(&requests, FAILURE_CALL_ID).expect("model-visible Bash failure receipt"); |
| 56 | for (receipt, sentinel) in [ |
| 57 | (&success_receipt, SUCCESS_SENTINEL), |
| 58 | (&failure_receipt, FAILURE_SENTINEL), |
| 59 | ] { |
| 60 | assert!( |
| 61 | receipt.contains("of output omitted"), |
| 62 | "model-facing truncation must state how much was omitted" |
| 63 | ); |
| 64 | assert!( |
| 65 | receipt.contains("full output at"), |
| 66 | "model-facing truncation must name the recovery path" |
| 67 | ); |
| 68 | assert!( |
| 69 | // The footer prints the artifact directory with the platform's |
| 70 | // path separator; compare on a normalized view so Windows |
| 71 | // backslashes don't fail an otherwise-correct footer. |
| 72 | receipt.replace('\\', "/").contains("/artifacts/"), |
| 73 | "the footer names where the omitted bytes live on disk" |
| 74 | ); |
| 75 | // The receipt must name a route the model can take *from this |
| 76 | // receipt*. It previously asserted the opposite — that the footer must |
| 77 | // NOT name `retrieve_tool_result` — which came from #5018's |
| 78 | // "no storage language" pass, not from any shell-vs-tool-result |
| 79 | // distinction: #4619 shipped the footer naming |
| 80 | // `retrieve_tool_result ref=art_<call>`, #5018 replaced the whole |
| 81 | // recovery line with "view full output in the tool details view" (a |
| 82 | // view the model cannot open) and froze that removal as a negative |
| 83 | // assertion, and #5212 restored the artifact path but left the stale |
| 84 | // negative in place. The follow-up probe below settles it empirically |
| 85 | // for *this* receipt: the scripted model reads the ref out of the |
| 86 | // receipt text it was handed, calls `retrieve_tool_result` with it, |
| 87 | // and gets back the exact line the receipt omitted. |
| 88 | assert!( |
| 89 | receipt.contains("retrieve_tool_result"), |
| 90 | "the footer must name the recovery route the model can actually take" |
| 91 | ); |
| 92 | assert!(!receipt.contains("[Exact evidence retained")); |
| 93 | assert!(!receipt.contains(sentinel)); |
| 94 | assert!( |
| 95 | receipt.len() <= 42_000, |
| 96 | "bounded preview must stay within the hybrid 32 KiB head + 8 KiB tail receipt budget, got {} bytes", |
| 97 | receipt.len() |
| 98 | ); |
| 99 | } |
| 100 | assert_ne!(success_receipt, failure_receipt); |
| 101 | |
| 102 | // The model followed the `ref=` the failure receipt named and got the byte |
| 103 | // range the receipt omitted. The mock parsed that ref out of the receipt |
| 104 | // text itself, so this only passes when the footer hands over a ref the |
| 105 | // retrieval tool can resolve in the origin session. |
| 106 | let retrieve_receipt = receipt_for(&requests, RETRIEVE_CALL_ID) |
| 107 | .expect("model-visible retrieve_tool_result receipt"); |
| 108 | assert!( |
| 109 | retrieve_receipt.contains(FAILURE_SENTINEL), |
| 110 | "retrieve_tool_result must return the range the receipt omitted, got: {retrieve_receipt}" |
| 111 | ); |
| 112 | |
| 113 | let artifact_dir = find_artifact_dir(home.path()).expect("origin-session artifacts"); |
| 114 | let payloads = std::fs::read_dir(&artifact_dir) |
| 115 | .expect("artifact directory") |
| 116 | .filter_map(Result::ok) |
| 117 | .map(|entry| entry.path()) |
| 118 | .filter(|path| path.extension().and_then(|ext| ext.to_str()) == Some("txt")) |
| 119 | .count(); |
| 120 | assert_eq!(payloads, 2, "exactly one evidence payload per result"); |
| 121 | |
| 122 | assert_ne!( |
| 123 | quoted_after(&success_receipt, "ref=\""), |
| 124 | quoted_after(&failure_receipt, "ref=\""), |
| 125 | "separate executions retain separate evidence handles" |
| 126 | ); |
| 127 | let success = assert_exact_artifact( |
| 128 | &artifact_dir, |
| 129 | SUCCESS_CALL_ID, |
| 130 | &success_receipt, |
| 131 | &events, |
| 132 | SUCCESS_SENTINEL, |
| 133 | "success", |
| 134 | ); |
| 135 | let failure = assert_exact_artifact( |
| 136 | &artifact_dir, |
| 137 | FAILURE_CALL_ID, |
| 138 | &failure_receipt, |
| 139 | &events, |
| 140 | FAILURE_SENTINEL, |
| 141 | "error", |
| 142 | ); |
| 143 | assert_ne!( |
| 144 | success, failure, |
| 145 | "success and failure bytes must stay distinct" |
| 146 | ); |
| 147 | } |
| 148 | |
| 149 | fn assert_exact_artifact( |
| 150 | artifact_dir: &Path, |
| 151 | provider_id: &str, |
| 152 | receipt: &str, |
| 153 | events: &[Value], |
| 154 | sentinel: &str, |
| 155 | status: &str, |
| 156 | ) -> Vec<u8> { |
| 157 | let handle = quoted_after(receipt, "ref=\"").expect("model-visible artifact handle"); |
| 158 | let call_id = handle.strip_prefix("art_").expect("tool artifact handle"); |
| 159 | uuid::Uuid::parse_str(call_id).expect("artifact belongs to a host execution"); |
| 160 | assert_ne!(call_id, provider_id); |
| 161 | for kind in ["tool_use", "tool_result"] { |
| 162 | let matching = events |
| 163 | .iter() |
| 164 | .filter(|event| event["type"] == kind && event["id"] == call_id) |
| 165 | .collect::<Vec<_>>(); |
| 166 | assert_eq!(matching.len(), 1, "one {kind} for this exact execution"); |
| 167 | let event = matching[0]; |
| 168 | assert_eq!(event["name"], "Bash"); |
| 169 | if kind == "tool_use" { |
| 170 | assert!( |
| 171 | event["input"]["command"] |
| 172 | .as_str() |
| 173 | .unwrap() |
| 174 | .contains(sentinel) |
| 175 | ); |
| 176 | } else { |
| 177 | // Model history trims the preview; stream events preserve it verbatim. |
| 178 | assert_eq!( |
| 179 | event["output"].as_str().expect("stream output").trim(), |
| 180 | receipt |
| 181 | ); |
| 182 | assert_eq!(event["status"], status); |
| 183 | } |
| 184 | } |
| 185 | let exact = |
| 186 | std::fs::read(artifact_dir.join(format!("{handle}.txt"))).expect("exact evidence bytes"); |
| 187 | assert!( |
| 188 | String::from_utf8_lossy(&exact).contains(sentinel), |
| 189 | "deep content omitted from context must remain retrievable" |
| 190 | ); |
| 191 | let metadata: Value = serde_json::from_slice( |
| 192 | &std::fs::read(artifact_dir.join(format!("{handle}.evidence.json"))) |
| 193 | .expect("evidence metadata"), |
| 194 | ) |
| 195 | .expect("valid evidence metadata"); |
| 196 | let digest = Sha256::digest(&exact) |
| 197 | .iter() |
| 198 | .map(|byte| format!("{byte:02x}")) |
| 199 | .collect::<String>(); |
| 200 | assert_eq!(metadata["handle"], handle); |
| 201 | assert_eq!(metadata["call_id"], call_id); |
| 202 | assert_eq!(metadata["tool_name"], "Bash"); |
| 203 | assert_eq!(metadata["digest"], digest); |
| 204 | assert_eq!(metadata["size_bytes"], exact.len() as u64); |
| 205 | assert_eq!(metadata["generation"], 1); |
| 206 | assert_eq!(metadata["redacted"], false); |
| 207 | assert_eq!(metadata["encoding"], "utf-8"); |
| 208 | assert_eq!(metadata["retention_state"], "live"); |
| 209 | assert!( |
| 210 | metadata["origin_session"] |
| 211 | .as_str() |
| 212 | .is_some_and(|id| !id.is_empty()) |
| 213 | ); |
| 214 | exact |
| 215 | } |
| 216 | |
| 217 | async fn mock_llm() -> MockServer { |
| 218 | let server = MockServer::start().await; |
| 219 | Mock::given(method("GET")) |
| 220 | .and(path("/v1/models")) |
| 221 | .respond_with(json_response(json!({ |
| 222 | "object": "list", |
| 223 | "data": [{"id": MODEL, "object": "model"}] |
| 224 | }))) |
| 225 | .mount(&server) |
| 226 | .await; |
| 227 | Mock::given(method("POST")) |
| 228 | .and(path("/v1/chat/completions")) |
| 229 | .respond_with(EvidenceScenario { |
| 230 | requests: Arc::new(AtomicUsize::new(0)), |
| 231 | }) |
| 232 | .mount(&server) |
| 233 | .await; |
| 234 | server |
| 235 | } |
| 236 | |
| 237 | #[derive(Clone)] |
| 238 | struct EvidenceScenario { |
| 239 | requests: Arc<AtomicUsize>, |
| 240 | } |
| 241 | |
| 242 | /// Scripted four-turn scenario. Turn 3 is the empirical half of the design |
| 243 | /// question this test settles: rather than asserting from the outside which |
| 244 | /// recovery route a Bash receipt *should* name, the scripted model reads the |
| 245 | /// route out of the receipt it was actually handed and takes it, so the |
| 246 | /// assertion in the test body observes whether that route returns the bytes |
| 247 | /// the receipt omitted. |
| 248 | impl Respond for EvidenceScenario { |
| 249 | fn respond(&self, request: &Request) -> ResponseTemplate { |
| 250 | let sequence = self.requests.fetch_add(1, Ordering::SeqCst); |
| 251 | let body = request.body_json::<Value>().unwrap_or(Value::Null); |
| 252 | let response = match sequence { |
| 253 | 0 => bash_tool_sse(SUCCESS_CALL_ID, true), |
| 254 | 1 => bash_tool_sse(FAILURE_CALL_ID, false), |
| 255 | 2 => { |
| 256 | let receipt = tool_result_content_for(&body, FAILURE_CALL_ID) |
| 257 | .expect("Bash failure receipt on the third request"); |
| 258 | let reference = quoted_after(receipt, "ref=\"") |
| 259 | .expect("failure receipt must hand over a retrieval ref"); |
| 260 | retrieve_tool_sse(reference, FAILURE_SENTINEL) |
| 261 | } |
| 262 | _ => final_sse(), |
| 263 | }; |
| 264 | assert!( |
| 265 | sequence < 4, |
| 266 | "unexpected extra model request #{sequence}: {body}" |
| 267 | ); |
| 268 | sse_response(response) |
| 269 | } |
| 270 | } |
| 271 | |
| 272 | /// Read the recovery ref the truncation footer hands the model, e.g. the |
| 273 | /// `art_<execution-id>` inside `… call retrieve_tool_result with ref="…"`. |
| 274 | fn quoted_after<'a>(receipt: &'a str, marker: &str) -> Option<&'a str> { |
| 275 | let rest = &receipt[receipt.find(marker)? + marker.len()..]; |
| 276 | rest.get(..rest.find('"')?) |
| 277 | } |
| 278 | |
| 279 | fn run_exec(workspace: &Path, home: &Path, server: &MockServer) -> std::process::Output { |
| 280 | std::fs::create_dir_all(home.join(".codewhale")).expect("config directory"); |
| 281 | std::fs::create_dir_all(home.join(".deepseek")).expect("legacy config directory"); |
| 282 | std::fs::write( |
| 283 | home.join(".codewhale/config.toml"), |
| 284 | "allow_shell = true\n\n[retry]\nenabled = false\n", |
| 285 | ) |
| 286 | .expect("headless test config"); |
| 287 | let mut command = Command::new(crate::binary::codewhale()); |
| 288 | preserve_host_env(&mut command); |
| 289 | command |
| 290 | .current_dir(workspace) |
| 291 | .args(["--workspace", workspace.to_str().expect("workspace utf8")]) |
| 292 | .arg("--no-project-config") |
| 293 | .args([ |
| 294 | "exec", |
| 295 | "--auto", |
| 296 | "--model", |
| 297 | MODEL, |
| 298 | "--output-format", |
| 299 | "stream-json", |
| 300 | ]) |
| 301 | .arg("run both provider-fixtured Bash evidence probes") |
| 302 | .env("HOME", home) |
| 303 | .env("USERPROFILE", home) |
| 304 | .env("XDG_CONFIG_HOME", home.join(".config")) |
| 305 | .env("XDG_DATA_HOME", home.join(".local/share")) |
| 306 | .env("XDG_CACHE_HOME", home.join(".cache")) |
| 307 | .env("CODEWHALE_CONFIG_PATH", home.join(".codewhale/config.toml")) |
| 308 | .env("DEEPSEEK_CONFIG_PATH", home.join(".deepseek/config.toml")) |
| 309 | .env("DEEPSEEK_API_KEY", "ci-test-key-not-real") |
| 310 | .env("DEEPSEEK_BASE_URL", server.uri()) |
| 311 | .env("CODEWHALE_BASE_URL", server.uri()) |
| 312 | .env("DEEPSEEK_MODEL", MODEL) |
| 313 | .env("CODEWHALE_MODEL", MODEL) |
| 314 | // The adaptive evidence lane is an explicit opt-in; this test exists |
| 315 | // to prove it end to end, so the spawned binary runs opted in. |
| 316 | .env("CODEWHALE_ADAPTIVE_OUTPUT_ROUTING", "1") |
| 317 | .env("RUST_LOG", "warn") |
| 318 | .stdout(Stdio::piped()) |
| 319 | .stderr(Stdio::piped()); |
| 320 | run_with_timeout(command, Duration::from_secs(45)) |
| 321 | } |
| 322 | |
| 323 | fn find_artifact_dir(home: &Path) -> Option<PathBuf> { |
| 324 | let sessions = home.join(".codewhale/sessions"); |
| 325 | std::fs::read_dir(sessions) |
| 326 | .ok()? |
| 327 | .filter_map(Result::ok) |
| 328 | .find_map(|entry| { |
| 329 | let path = entry.path().join("artifacts"); |
| 330 | path.is_dir().then_some(path) |
| 331 | }) |
| 332 | } |
| 333 | |
| 334 | fn receipt_for(requests: &[Request], call_id: &str) -> Option<String> { |
| 335 | requests |
| 336 | .iter() |
| 337 | .filter_map(|request| request.body_json::<Value>().ok()) |
| 338 | .find_map(|body| tool_result_content_for(&body, call_id).map(str::to_owned)) |
| 339 | } |
| 340 | |
| 341 | fn tool_result_content_for<'a>(body: &'a Value, call_id: &str) -> Option<&'a str> { |
| 342 | body.get("messages")? |
| 343 | .as_array()? |
| 344 | .iter() |
| 345 | .find(|message| { |
| 346 | message.get("role").and_then(Value::as_str) == Some("tool") |
| 347 | && message.get("tool_call_id").and_then(Value::as_str) == Some(call_id) |
| 348 | })? |
| 349 | .get("content")? |
| 350 | .as_str() |
| 351 | } |
| 352 | |
| 353 | fn bash_tool_sse(call_id: &str, success: bool) -> String { |
| 354 | let (sentinel, prefix) = if success { |
| 355 | (SUCCESS_SENTINEL, "BASH-SUCCESS") |
| 356 | } else { |
| 357 | (FAILURE_SENTINEL, "BASH-FAILURE") |
| 358 | }; |
| 359 | let command = probe_command(sentinel, prefix, success); |
| 360 | tool_call_sse( |
| 361 | call_id, |
| 362 | "Bash", |
| 363 | &json!({"action": "run", "command": command, "timeout_ms": 30_000}), |
| 364 | ) |
| 365 | } |
| 366 | |
| 367 | /// Take the recovery route the receipt named, asking for the omitted range by |
| 368 | /// the sentinel that rides in it. |
| 369 | fn retrieve_tool_sse(reference: &str, query: &str) -> String { |
| 370 | tool_call_sse( |
| 371 | RETRIEVE_CALL_ID, |
| 372 | "retrieve_tool_result", |
| 373 | &json!({"ref": reference, "mode": "query", "query": query}), |
| 374 | ) |
| 375 | } |
| 376 | |
| 377 | fn tool_call_sse(call_id: &str, name: &str, arguments: &Value) -> String { |
| 378 | let arguments = serde_json::to_string(arguments).expect("tool arguments"); |
| 379 | [ |
| 380 | chunk(json!({"id":"tool","object":"chat.completion.chunk","model":MODEL,"choices":[{"index":0,"delta":{"tool_calls":[{"index":0,"id":call_id,"type":"function","function":{"name":name,"arguments":arguments}}]},"finish_reason":null}]})), |
| 381 | chunk(json!({"id":"tool","object":"chat.completion.chunk","model":MODEL,"choices":[{"index":0,"delta":{},"finish_reason":"tool_calls"}],"usage":{"prompt_tokens":10,"completion_tokens":2,"total_tokens":12}})), |
| 382 | "data: [DONE]\n\n".to_string(), |
| 383 | ].join("") |
| 384 | } |
| 385 | |
| 386 | /// Stderr filler line the sentinel rides on, which has to land in a narrow |
| 387 | /// window. `shell_output` bounds each stream to `TRUNCATED_HEAD_BYTES` = |
| 388 | /// 30_000/5 = 6_000 bytes of head plus 24_000 of tail, so anything past ~line |
| 389 | /// 68 of stderr never reaches the artifact at all. Below that, the bounded |
| 390 | /// stdout section (~30.1 KB: head + notice + tail) plus the `STDERR:` |
| 391 | /// separator puts stderr filler line *n* at roughly 30_110 + 86n bytes of the |
| 392 | /// result, so anything before ~line 31 is still inside the preview's 32 KiB |
| 393 | /// head and the receipt would show it. Line 50 sits near the middle of |
| 394 | /// [31, 68] with ~1.6 KB of slack on each side. |
| 395 | /// |
| 396 | /// #5212 wrote 100 here against a "22 KB head bound" that does not exist: the |
| 397 | /// sentinel landed in the stream's own omitted middle, so the artifact never |
| 398 | /// carried it and this test's retrievability assertion failed. It went |
| 399 | /// unnoticed because an earlier assertion in the same loop failed first. |
| 400 | const SENTINEL_LINE: usize = 50; |
| 401 | |
| 402 | /// Shell fixture that emits enough bytes to force exact-evidence routing under |
| 403 | /// the 32_768-token default threshold. The Bash adapter self-bounds each |
| 404 | /// stream to ~30 KB, so a single stream would now fit inside the hybrid |
| 405 | /// 32 KiB + 8 KiB preview budget; the probe therefore fills stdout AND stderr |
| 406 | /// (~60 KB combined) so the envelope still omits a middle range, with the |
| 407 | /// sentinel at [`SENTINEL_LINE`] of stderr. The probe executes through the |
| 408 | /// platform shell — bash on Unix, `cmd /C` on Windows (#1691) — so each |
| 409 | /// platform needs native syntax to exercise the same routing path. |
| 410 | #[cfg(not(windows))] |
| 411 | fn probe_command(sentinel: &str, prefix: &str, success: bool) -> String { |
| 412 | let trailer = if success { "" } else { "; exit 7" }; |
| 413 | let stdout_loop = format!( |
| 414 | "i=0; while [ \"$i\" -lt 2800 ]; do printf '{prefix}-%04d-xxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxx\\n' \"$i\"; i=$((i + 1)); done" |
| 415 | ); |
| 416 | let stderr_loop = format!( |
| 417 | "j=0; while [ \"$j\" -lt 2800 ]; do if [ \"$j\" -eq {SENTINEL_LINE} ]; then printf '%s\\n' '{sentinel}'; fi; printf '{prefix}-ERR-%04d-xxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxx\\n' \"$j\"; j=$((j + 1)); done" |
| 418 | ); |
| 419 | format!("{stdout_loop}; {{ {stderr_loop}; }} >&2{trailer}") |
| 420 | } |
| 421 | |
| 422 | /// PowerShell syntax: on Windows the shell dispatcher prefers `pwsh.exe`, |
| 423 | /// then in-box `powershell.exe`, only falling back to `cmd.exe` when no |
| 424 | /// PowerShell exists at all. Single quotes only — the payload is passed to |
| 425 | /// `-Command` as one argv string, and four or more double quotes would push |
| 426 | /// it onto the temp-`-File` path for no benefit. The failure variant mirrors |
| 427 | /// the Unix `{ ...; } >&2; exit 7` shape by writing every line to the OS |
| 428 | /// stderr handle and exiting 7 after the loop. |
| 429 | #[cfg(windows)] |
| 430 | fn probe_command(sentinel: &str, prefix: &str, success: bool) -> String { |
| 431 | let stdout_line = format!( |
| 432 | "'{prefix}-{{0}}-xxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxx'" |
| 433 | ); |
| 434 | let stderr_line = format!( |
| 435 | "'{prefix}-ERR-{{0}}-xxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxx'" |
| 436 | ); |
| 437 | let stdout_loop = format!("0..2799 | ForEach-Object {{ Write-Output ({stdout_line} -f $_) }}"); |
| 438 | let stderr_loop = format!( |
| 439 | "0..2799 | ForEach-Object {{ if ($_ -eq {SENTINEL_LINE}) {{ [Console]::Error.WriteLine('{sentinel}') }}; [Console]::Error.WriteLine(({stderr_line} -f $_)) }}" |
| 440 | ); |
| 441 | let trailer = if success { "" } else { "; exit 7" }; |
| 442 | format!("{stdout_loop}; {stderr_loop}{trailer}") |
| 443 | } |
| 444 | |
| 445 | fn final_sse() -> String { |
| 446 | [ |
| 447 | chunk(json!({"id":"final","object":"chat.completion.chunk","model":MODEL,"choices":[{"index":0,"delta":{"content":"evidence retained"},"finish_reason":null}]})), |
| 448 | chunk(json!({"id":"final","object":"chat.completion.chunk","model":MODEL,"choices":[{"index":0,"delta":{},"finish_reason":"stop"}],"usage":{"prompt_tokens":20,"completion_tokens":2,"total_tokens":22}})), |
| 449 | "data: [DONE]\n\n".to_string(), |
| 450 | ].join("") |
| 451 | } |
| 452 | |
| 453 | fn chunk(value: Value) -> String { |
| 454 | format!( |
| 455 | "data: {}\n\n", |
| 456 | serde_json::to_string(&value).expect("SSE JSON") |
| 457 | ) |
| 458 | } |
| 459 | |
| 460 | fn sse_response(body: String) -> ResponseTemplate { |
| 461 | ResponseTemplate::new(200) |
| 462 | .insert_header("content-type", "text/event-stream") |
| 463 | .set_body_string(body) |
| 464 | } |
| 465 | |
| 466 | fn json_response(value: Value) -> ResponseTemplate { |
| 467 | ResponseTemplate::new(200).set_body_json(value) |
| 468 | } |
| 469 | |
| 470 | fn preserve_host_env(command: &mut Command) { |
| 471 | command.env_clear(); |
| 472 | for key in [ |
| 473 | "PATH", |
| 474 | "PATHEXT", |
| 475 | "SystemRoot", |
| 476 | "SystemDrive", |
| 477 | "WINDIR", |
| 478 | "COMSPEC", |
| 479 | "TEMP", |
| 480 | "TMP", |
| 481 | "TERM", |
| 482 | "LANG", |
| 483 | "LC_ALL", |
| 484 | ] { |
| 485 | if let Some(value) = std::env::var_os(key) { |
| 486 | command.env(key, value); |
| 487 | } |
| 488 | } |
| 489 | } |
| 490 | |
| 491 | fn run_with_timeout(mut command: Command, timeout: Duration) -> std::process::Output { |
| 492 | let mut child = command.spawn().expect("spawn codewhale exec"); |
| 493 | let stdout = read_in_background(child.stdout.take().expect("stdout")); |
| 494 | let stderr = read_in_background(child.stderr.take().expect("stderr")); |
| 495 | let status = child |
| 496 | .wait_timeout(timeout) |
| 497 | .expect("wait") |
| 498 | .unwrap_or_else(|| { |
| 499 | let _ = child.kill(); |
| 500 | let _ = child.wait(); |
| 501 | panic!("codewhale exec timed out") |
| 502 | }); |
| 503 | std::process::Output { |
| 504 | status, |
| 505 | stdout: stdout.join().expect("stdout thread").expect("read stdout"), |
| 506 | stderr: stderr.join().expect("stderr thread").expect("read stderr"), |
| 507 | } |
| 508 | } |
| 509 | |
| 510 | fn read_in_background<R: Read + Send + 'static>( |
| 511 | mut reader: R, |
| 512 | ) -> std::thread::JoinHandle<std::io::Result<Vec<u8>>> { |
| 513 | std::thread::spawn(move || { |
| 514 | let mut bytes = Vec::new(); |
| 515 | reader.read_to_end(&mut bytes).map(|_| bytes) |
| 516 | }) |
| 517 | } |
| 518 |