返回 CodeWhale
lifecycle_outbox_exec.rs
根目录 / crates / tui / tests / integration / lifecycle_outbox_exec.rs
1 //! End-to-end contract for the `[lifecycle_outbox]` feature on headless
2 //! `codewhale exec`: with a path configured, a run appends one JSONL
3 //! `RuntimeEventEnvelope` line per turn boundary (`turn_start` at message
4 //! dispatch, `turn_end` at the terminal receipt), the per-file `seq` recovers
5 //! across processes, and with no path configured no file is ever created.
6 //!
7 //! A `wiremock` OpenAI-compatible endpoint stands in for the provider, so the
8 //! run is a real `exec` process end to end — same loader, same engine, same
9 //! outbox writer — with no external network.
10
11 #![cfg(unix)]
12
13 use std::io::Read;
14 use std::path::{Path, PathBuf};
15 use std::process::{Command, Stdio};
16 use std::time::Duration;
17
18 use serde_json::{Value, json};
19 use tempfile::TempDir;
20 use wait_timeout::ChildExt;
21 use wiremock::matchers::{method, path};
22 use wiremock::{Mock, MockServer, ResponseTemplate};
23
24 const TEST_MODEL: &str = "lifecycle-outbox-model";
25 const RUN_TIMEOUT: Duration = Duration::from_secs(60);
26
27 /// Placeholder in `outbox_toml` replaced with the isolated home's absolute
28 /// outbox path (so callers can read the file back after the run).
29 const OUTBOX_PATH_TOKEN: &str = "__OUTBOX_PATH__";
30
31 fn sse_chunk(value: Value) -> String {
32 format!(
33 "data: {}\n\n",
34 serde_json::to_string(&value).expect("SSE JSON")
35 )
36 }
37
38 /// Final-answer SSE: one content delta, then a clean stop.
39 fn answer_sse(answer: &str) -> String {
40 [
41 sse_chunk(json!({
42 "id": "chatcmpl-outbox",
43 "object": "chat.completion.chunk",
44 "model": TEST_MODEL,
45 "choices": [{"index": 0, "delta": {"content": answer}, "finish_reason": null}]
46 })),
47 sse_chunk(json!({
48 "id": "chatcmpl-outbox",
49 "object": "chat.completion.chunk",
50 "model": TEST_MODEL,
51 "choices": [{"index": 0, "delta": {}, "finish_reason": "stop"}]
52 })),
53 "data: [DONE]\n\n".to_string(),
54 ]
55 .join("")
56 }
57
58 async fn start_mock_llm() -> MockServer {
59 let server = MockServer::start().await;
60
61 Mock::given(method("GET"))
62 .and(path("/v1/models"))
63 .respond_with(
64 ResponseTemplate::new(200)
65 .insert_header("content-type", "application/json")
66 .set_body_json(json!({
67 "object": "list",
68 "data": [{ "id": TEST_MODEL, "object": "model" }]
69 })),
70 )
71 .mount(&server)
72 .await;
73
74 Mock::given(method("POST"))
75 .and(path("/v1/chat/completions"))
76 .respond_with(
77 ResponseTemplate::new(200)
78 .insert_header("content-type", "text/event-stream")
79 .insert_header("cache-control", "no-cache")
80 .set_body_string(answer_sse("ok")),
81 )
82 .mount(&server)
83 .await;
84
85 server
86 }
87
88 fn preserve_host_env(command: &mut Command) {
89 command.env_clear();
90 for key in [
91 "PATH",
92 "PATHEXT",
93 "SystemRoot",
94 "SystemDrive",
95 "WINDIR",
96 "COMSPEC",
97 "TEMP",
98 "TMP",
99 "TERM",
100 "COLORTERM",
101 "LANG",
102 "LC_ALL",
103 ] {
104 if let Some(value) = std::env::var_os(key) {
105 command.env(key, value);
106 }
107 }
108 }
109
110 /// Run `codewhale exec` against the mock provider with the given
111 /// `[lifecycle_outbox]` config block (already TOML-formatted, may be empty).
112 /// Any `__OUTBOX_PATH__` token in it is replaced with the isolated home's
113 /// absolute outbox path. Returns the isolated home and workspace dirs (the
114 /// latter so callers can assert the outbox `payload.workspace` exactly).
115 fn run_exec_with_outbox_config(
116 server: &MockServer,
117 outbox_toml: &str,
118 expected_exit_code: i32,
119 ) -> (TempDir, TempDir) {
120 let workspace = TempDir::new().expect("workspace tempdir");
121 let home = TempDir::new().expect("home tempdir");
122 let outbox_path = home_outbox_path(&home);
123 let outbox_toml = outbox_toml.replace(OUTBOX_PATH_TOKEN, &outbox_path.display().to_string());
124
125 std::fs::create_dir_all(home.path().join(".codewhale")).expect("create codewhale config dir");
126 std::fs::create_dir_all(home.path().join(".deepseek")).expect("create deepseek config dir");
127 std::fs::write(
128 home.path().join(".codewhale").join("config.toml"),
129 format!("provider = \"deepseek\"\nmodel = \"{TEST_MODEL}\"\n{outbox_toml}"),
130 )
131 .expect("write exec config");
132
133 let mut command = Command::new(crate::binary::codewhale());
134 preserve_host_env(&mut command);
135 command
136 .current_dir(workspace.path())
137 .arg("--workspace")
138 .arg(workspace.path())
139 .arg("--no-project-config")
140 .arg("exec")
141 .arg("--auto")
142 .arg("--model")
143 .arg(TEST_MODEL)
144 .arg("answer briefly")
145 .env("HOME", home.path())
146 .env("USERPROFILE", home.path())
147 .env("XDG_CONFIG_HOME", home.path().join(".config"))
148 .env("XDG_DATA_HOME", home.path().join(".local").join("share"))
149 .env("XDG_CACHE_HOME", home.path().join(".cache"))
150 .env(
151 "CODEWHALE_CONFIG_PATH",
152 home.path().join(".codewhale").join("config.toml"),
153 )
154 .env(
155 "DEEPSEEK_CONFIG_PATH",
156 home.path().join(".deepseek").join("config.toml"),
157 )
158 .env("DEEPSEEK_API_KEY", "ci-test-key-not-real")
159 .env("DEEPSEEK_BASE_URL", server.uri())
160 .env("CODEWHALE_BASE_URL", server.uri())
161 .env("DEEPSEEK_MODEL", TEST_MODEL)
162 .env("CODEWHALE_MODEL", TEST_MODEL)
163 .env("RUST_LOG", "warn")
164 .stdout(Stdio::piped())
165 .stderr(Stdio::piped());
166
167 let mut child = command.spawn().expect("spawn codewhale exec");
168 let stdout_reader = read_pipe_in_background(child.stdout.take().expect("stdout pipe"));
169 let stderr_reader = read_pipe_in_background(child.stderr.take().expect("stderr pipe"));
170
171 let status = match child.wait_timeout(RUN_TIMEOUT).expect("wait for codewhale") {
172 Some(status) => status,
173 None => {
174 let _ = child.kill();
175 let _ = child.wait();
176 let stdout = join_pipe_reader(stdout_reader, "stdout");
177 let stderr = join_pipe_reader(stderr_reader, "stderr");
178 panic!(
179 "codewhale exec timed out after {RUN_TIMEOUT:?}\nstdout:\n{}\nstderr:\n{}",
180 String::from_utf8_lossy(&stdout),
181 String::from_utf8_lossy(&stderr)
182 );
183 }
184 };
185
186 let stdout = join_pipe_reader(stdout_reader, "stdout");
187 let stderr = join_pipe_reader(stderr_reader, "stderr");
188 assert_eq!(
189 status.code(),
190 Some(expected_exit_code),
191 "codewhale exec returned the wrong exit status\nstdout:\n{}\nstderr:\n{}",
192 String::from_utf8_lossy(&stdout),
193 String::from_utf8_lossy(&stderr)
194 );
195
196 (home, workspace)
197 }
198
199 fn read_pipe_in_background<R>(mut reader: R) -> std::thread::JoinHandle<std::io::Result<Vec<u8>>>
200 where
201 R: Read + Send + 'static,
202 {
203 std::thread::spawn(move || {
204 let mut output = Vec::new();
205 reader.read_to_end(&mut output).map(|_| output)
206 })
207 }
208
209 fn join_pipe_reader(
210 handle: std::thread::JoinHandle<std::io::Result<Vec<u8>>>,
211 stream_name: &str,
212 ) -> Vec<u8> {
213 handle
214 .join()
215 .expect("pipe reader join")
216 .unwrap_or_else(|err| panic!("failed to read {stream_name}: {err}"))
217 }
218
219 fn read_outbox_lines(path: &Path) -> Vec<Value> {
220 let text = std::fs::read_to_string(path).expect("read outbox file");
221 text.lines()
222 .map(|line| {
223 serde_json::from_str(line)
224 .unwrap_or_else(|err| panic!("outbox line should parse: {err}\nline: {line}"))
225 })
226 .collect()
227 }
228
229 fn home_outbox_path(home: &TempDir) -> PathBuf {
230 home.path()
231 .join(".codewhale")
232 .join("notifications")
233 .join("outbox.jsonl")
234 }
235
236 #[tokio::test(flavor = "multi_thread")]
237 async fn exec_emits_turn_start_and_turn_end_to_the_configured_outbox() {
238 let server = start_mock_llm().await;
239 let (home, workspace) = run_exec_with_outbox_config(
240 &server,
241 &format!("[lifecycle_outbox]\npath = {}\n", json!(OUTBOX_PATH_TOKEN)),
242 0,
243 );
244
245 let outbox_path = home_outbox_path(&home);
246 assert!(outbox_path.exists(), "outbox file must be created");
247 let lines = read_outbox_lines(&outbox_path);
248 assert_eq!(
249 lines.len(),
250 2,
251 "one turn_start and one turn_end line: {lines:#?}"
252 );
253
254 let start = &lines[0];
255 assert_eq!(start["event"], "turn_start");
256 assert_eq!(start["kind"], "turn.started");
257 assert_eq!(start["schema_version"], 1);
258 assert_eq!(start["seq"], 1);
259 assert!(start["timestamp"].as_str().is_some());
260 // Headless exec has no engine turn id and (for a fresh run) no session
261 // id yet — both are honest absences, never fabricated.
262 assert!(start["turn_id"].is_null());
263 // The model field is bounded and never the raw prompt.
264 assert_eq!(start["payload"]["model"], TEST_MODEL);
265 // Every payload carries the workspace for consumer-side routing; exec
266 // runs with `--workspace <dir>`, so the emitted path must match it.
267 assert_eq!(
268 start["payload"]["workspace"],
269 json!(workspace.path().to_string_lossy().as_ref()),
270 "turn_start must carry the workspace"
271 );
272
273 let end = &lines[1];
274 assert_eq!(end["event"], "turn_end");
275 assert_eq!(end["kind"], "turn.completed");
276 assert_eq!(end["seq"], 2);
277 assert_eq!(end["payload"]["status"], "completed");
278 assert!(end["payload"]["error"].is_null());
279 assert!(end["payload"]["duration_ms"].as_u64().is_some());
280 assert_eq!(
281 end["payload"]["workspace"],
282 json!(workspace.path().to_string_lossy().as_ref()),
283 "turn_end must carry the workspace"
284 );
285 }
286
287 #[tokio::test(flavor = "multi_thread")]
288 async fn exec_without_outbox_config_writes_no_file() {
289 let server = start_mock_llm().await;
290 let (home, _workspace) = run_exec_with_outbox_config(&server, "", 0);
291
292 assert!(
293 !home_outbox_path(&home).exists(),
294 "no outbox file must be created when [lifecycle_outbox] is unset"
295 );
296 }
297
298 #[tokio::test(flavor = "multi_thread")]
299 async fn outbox_seq_recovers_across_processes() {
300 let server = start_mock_llm().await;
301
302 // First run writes seq 1 (turn_start) and 2 (turn_end).
303 let (home, _workspace) = run_exec_with_outbox_config(
304 &server,
305 &format!("[lifecycle_outbox]\npath = {}\n", json!(OUTBOX_PATH_TOKEN)),
306 0,
307 );
308 let shared_outbox = home_outbox_path(&home);
309
310 // Second process, pointing at the SAME file: seq must continue at 3.
311 let (_second_home, _second_workspace) = run_exec_with_outbox_config(
312 &server,
313 &format!(
314 "[lifecycle_outbox]\npath = {}\n",
315 json!(shared_outbox.display().to_string())
316 ),
317 0,
318 );
319
320 let lines = read_outbox_lines(&shared_outbox);
321 assert_eq!(lines.len(), 4, "two runs, four lines: {lines:#?}");
322 let seqs: Vec<u64> = lines
323 .iter()
324 .map(|line| line["seq"].as_u64().expect("seq"))
325 .collect();
326 assert_eq!(
327 seqs,
328 vec![1, 2, 3, 4],
329 "seq must be monotonic across processes"
330 );
331 }
332
333 #[tokio::test(flavor = "multi_thread")]
334 async fn failed_exec_persists_the_terminal_receipt_without_changing_its_exit() {
335 let server = start_mock_llm().await;
336 Mock::given(method("POST"))
337 .and(path("/v1/chat/completions"))
338 .respond_with(ResponseTemplate::new(500).set_body_string("upstream model failure"))
339 .with_priority(1)
340 .mount(&server)
341 .await;
342 let (home, _workspace) = run_exec_with_outbox_config(
343 &server,
344 &format!("[lifecycle_outbox]\npath = {}\n", json!(OUTBOX_PATH_TOKEN)),
345 1,
346 );
347 let lines = read_outbox_lines(&home_outbox_path(&home));
348 assert_eq!(lines.len(), 2, "failed turn retains both boundaries");
349 assert_eq!(lines[0]["event"], "turn_start");
350 assert_eq!(lines[0]["seq"], 1);
351 assert_eq!(lines[1]["event"], "turn_end");
352 assert_eq!(lines[1]["kind"], "turn.failed");
353 assert_eq!(lines[1]["seq"], 2);
354 assert_eq!(lines[1]["payload"]["status"], "failed");
355 assert!(!lines[1]["payload"]["error"].as_str().unwrap().is_empty());
356 }
357
357 lines RUST