返回 CodeWhale
runtime.rs
根目录 / crates / lane / src / runtime.rs
1 //! Runtime backends: tmux durability and inline (#4176). The never-implemented
2 //! `vm` / `ci` kinds are retired and only load from historical records (#6516).
3 //!
4 //! Runtime owns process/session lifecycle and stream-json log capture.
5 //! Fleet modules must not import this module.
6
7 use std::fs::{self, OpenOptions};
8 use std::io::{BufRead, BufReader, Write};
9 use std::path::{Path, PathBuf};
10 use std::process::{Command, ExitStatus, Stdio};
11 use std::thread;
12
13 use anyhow::{Context, Result, bail};
14 use serde::{Deserialize, Serialize};
15
16 use crate::registry::{LaneRecord, LaneRegistry, LaneStatus, TerminalTransition};
17 use crate::worktree::{WorktreeProvision, provision_worktree, remove_worktree_if_expired};
18
19 /// Execution backend for a lane.
20 #[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
21 #[serde(rename_all = "snake_case")]
22 pub enum RuntimeBackendKind {
23 Tmux,
24 Inline,
25 /// Retired (#6516): `vm` and `ci` were advertised but never implemented.
26 /// [`Self::parse`] no longer accepts them; the variants stay only so lane
27 /// records written by earlier releases (always `failed`) still load
28 /// instead of being skipped as corrupt.
29 Vm,
30 Ci,
31 }
32
33 impl RuntimeBackendKind {
34 pub fn as_str(self) -> &'static str {
35 match self {
36 Self::Tmux => "tmux",
37 Self::Inline => "inline",
38 Self::Vm => "vm",
39 Self::Ci => "ci",
40 }
41 }
42
43 pub fn parse(raw: &str) -> Result<Self> {
44 match raw.trim().to_ascii_lowercase().as_str() {
45 "tmux" => Ok(Self::Tmux),
46 "inline" => Ok(Self::Inline),
47 "vm" | "ci" => bail!(
48 "runtime backend `{}` is not available; use tmux or inline",
49 raw.trim().to_ascii_lowercase()
50 ),
51 other => bail!("unknown runtime backend `{other}` (use tmux|inline)"),
52 }
53 }
54 }
55
56 /// Inputs for starting a lane under a runtime backend.
57 #[derive(Clone)]
58 pub struct LaneStartSpec {
59 /// Command argv to run inside the backend (e.g. `codewhale exec …`).
60 pub command: Vec<String>,
61 /// Working directory for the command (defaults to worktree or cwd).
62 pub cwd: Option<PathBuf>,
63 /// Process-local runtime overrides. Values are never written into the
64 /// Lane record or command argv; tmux bridges them through a private 0600
65 /// environment file that the detached shell removes before execution.
66 pub environment: Vec<(String, String)>,
67 /// Executable that exposes Codewhale's hidden `lane-log-proxy` command.
68 /// Required by tmux so arbitrary/binary child output is framed as valid
69 /// NDJSON without trusting a shell pipeline.
70 pub log_proxy: Option<PathBuf>,
71 /// When set, provision an isolated git worktree + branch under this repo.
72 pub worktree: Option<WorktreeProvision>,
73 }
74
75 impl std::fmt::Debug for LaneStartSpec {
76 fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
77 f.debug_struct("LaneStartSpec")
78 .field("command", &self.command)
79 .field("cwd", &self.cwd)
80 .field(
81 "environment_keys",
82 &self
83 .environment
84 .iter()
85 .map(|(key, _)| key)
86 .collect::<Vec<_>>(),
87 )
88 .field("log_proxy", &self.log_proxy)
89 .field("worktree", &self.worktree)
90 .finish()
91 }
92 }
93
94 /// Runtime adapter contract.
95 pub trait RuntimeBackend {
96 fn kind(&self) -> RuntimeBackendKind;
97
98 /// Start the lane process/session; mutates record with attach/log metadata.
99 fn start(
100 &self,
101 registry: &LaneRegistry,
102 record: &mut LaneRecord,
103 spec: &LaneStartSpec,
104 ) -> Result<()>;
105
106 /// Human attach command, if any (tmux).
107 fn attach_command(&self, record: &LaneRecord) -> Option<String>;
108
109 /// Stop the running session/process.
110 ///
111 /// `fence` pins the durable lifecycle generation the caller observed. It is
112 /// evaluated under the registry's per-Lane lock, so a record that moved on
113 /// between the caller's read and this call is refused rather than stopped.
114 /// The return value distinguishes "this call stopped it" from "it was
115 /// already terminal" and "the fence did not match", so a caller can report
116 /// a truthful outcome instead of claiming a transition someone else made.
117 fn stop(
118 &self,
119 registry: &LaneRegistry,
120 record: &mut LaneRecord,
121 fence: Option<u64>,
122 ) -> Result<TerminalTransition>;
123
124 /// Reconcile durable backend state into the Lane record before display.
125 /// Most backends update synchronously and need no refresh; tmux records
126 /// the detached process exit in a private sidecar and folds it in on the
127 /// next read.
128 ///
129 /// Returns `true` when reconciliation *changed durable state*. Reads are
130 /// allowed to fold a finished Runtime exit into the record, but a surface
131 /// that calls this must be able to say so rather than reporting a pure
132 /// observation (#4022).
133 fn reconcile(&self, _registry: &LaneRegistry, _record: &mut LaneRecord) -> Result<bool> {
134 Ok(false)
135 }
136
137 /// Optional worktree TTL cleanup after stop.
138 fn cleanup_worktree(&self, record: &LaneRecord) -> Result<()> {
139 if let Some(path) = record.worktree_path.as_ref() {
140 remove_worktree_if_expired(
141 path,
142 record.worktree_ttl_secs,
143 record.stopped_at.as_deref(),
144 )?;
145 }
146 Ok(())
147 }
148 }
149
150 pub fn resolve_backend(kind: RuntimeBackendKind) -> Box<dyn RuntimeBackend> {
151 match kind {
152 RuntimeBackendKind::Tmux => Box::new(TmuxRuntime),
153 RuntimeBackendKind::Inline => Box::new(InlineRuntime),
154 RuntimeBackendKind::Vm => Box::new(RetiredRuntime {
155 kind: RuntimeBackendKind::Vm,
156 }),
157 RuntimeBackendKind::Ci => Box::new(RetiredRuntime {
158 kind: RuntimeBackendKind::Ci,
159 }),
160 }
161 }
162
163 pub fn backend_for(record: &LaneRecord) -> Box<dyn RuntimeBackend> {
164 resolve_backend(record.runtime)
165 }
166
167 fn append_log_event(log_path: &Path, event: serde_json::Value) -> Result<()> {
168 let mut file = OpenOptions::new()
169 .create(true)
170 .append(true)
171 .open(log_path)
172 .with_context(|| format!("open lane log {}", log_path.display()))?;
173 let mut encoded = serde_json::to_vec(&event).context("serialize lane log event")?;
174 encoded.push(b'\n');
175 file.write_all(&encoded)
176 .with_context(|| format!("write lane log {}", log_path.display()))?;
177 Ok(())
178 }
179
180 const MAX_EXIT_RECEIPT_BYTES: u64 = 4 * 1024;
181 const MAX_ENVIRONMENT_BYTES: u64 = 1024 * 1024;
182 const MAX_CHILD_LOG_FRAME_BYTES: usize = 64 * 1024;
183 const LANE_PROXY_FAILURE_EXIT_CODE: i32 = 125;
184
185 #[derive(Debug, Deserialize, Serialize)]
186 struct LaneExitReceipt {
187 lane_id: String,
188 exit_code: i32,
189 }
190
191 /// Inputs to the hidden Rust log proxy used by detached runtimes.
192 ///
193 /// The proxy is deliberately exposed by the Lane crate so the thin CLI
194 /// facade can invoke it before loading user configuration. Secrets travel in
195 /// the private environment file, never in argv or the Lane record.
196 #[derive(Debug, Clone)]
197 pub struct LaneLogProxySpec {
198 pub command: Vec<String>,
199 pub log_path: PathBuf,
200 pub receipt_path: PathBuf,
201 pub receipt_tmp_path: PathBuf,
202 pub environment_path: Option<PathBuf>,
203 pub lane_id: String,
204 }
205
206 fn lane_exit_receipt_path(log_path: &Path) -> PathBuf {
207 log_path.with_extension("exit.json")
208 }
209
210 fn lane_exit_receipt_tmp_path(log_path: &Path) -> PathBuf {
211 log_path.with_extension("exit.json.tmp")
212 }
213
214 fn lane_environment_path(log_path: &Path) -> PathBuf {
215 log_path.with_extension("env.json")
216 }
217
218 fn lane_environment_tmp_path(path: &Path) -> PathBuf {
219 path.with_extension("json.tmp")
220 }
221
222 fn remove_file_if_present(path: &Path) -> Result<()> {
223 match std::fs::remove_file(path) {
224 Ok(()) => Ok(()),
225 Err(err) if err.kind() == std::io::ErrorKind::NotFound => Ok(()),
226 Err(err) => Err(err).with_context(|| format!("remove {}", path.display())),
227 }
228 }
229
230 fn valid_environment_key(key: &str) -> bool {
231 let mut chars = key.chars();
232 chars
233 .next()
234 .is_some_and(|ch| ch == '_' || ch.is_ascii_alphabetic())
235 && chars.all(|ch| ch == '_' || ch.is_ascii_alphanumeric())
236 }
237
238 fn write_lane_environment(path: &Path, environment: &[(String, String)]) -> Result<()> {
239 for (key, _) in environment {
240 if !valid_environment_key(key) {
241 bail!("invalid lane environment key {key:?}");
242 }
243 }
244 let encoded = serde_json::to_vec(environment).context("serialize private lane environment")?;
245 if encoded.len() as u64 > MAX_ENVIRONMENT_BYTES {
246 bail!(
247 "private lane environment exceeds {} bytes",
248 MAX_ENVIRONMENT_BYTES
249 );
250 }
251
252 let tmp_path = lane_environment_tmp_path(path);
253 remove_file_if_present(path)?;
254 remove_file_if_present(&tmp_path)?;
255 let mut options = OpenOptions::new();
256 options.create_new(true).write(true);
257 #[cfg(unix)]
258 {
259 use std::os::unix::fs::OpenOptionsExt;
260 options.mode(0o600);
261 }
262 let result = (|| {
263 let mut file = options
264 .open(&tmp_path)
265 .with_context(|| format!("create private lane environment {}", tmp_path.display()))?;
266 file.write_all(&encoded)
267 .with_context(|| format!("write private lane environment {}", tmp_path.display()))?;
268 file.sync_all()
269 .with_context(|| format!("sync private lane environment {}", tmp_path.display()))?;
270 fs::rename(&tmp_path, path)
271 .with_context(|| format!("publish private lane environment {}", path.display()))?;
272 Ok(())
273 })();
274 if result.is_err() {
275 let _ = remove_file_if_present(&tmp_path);
276 let _ = remove_file_if_present(path);
277 }
278 result
279 }
280
281 fn append_child_output(log_path: &Path, stream: &str, bytes: &[u8]) -> Result<()> {
282 let mut line = bytes;
283 if let Some(stripped) = line.strip_suffix(b"\n") {
284 line = stripped;
285 }
286 if let Some(stripped) = line.strip_suffix(b"\r") {
287 line = stripped;
288 }
289 if line.is_empty() {
290 return Ok(());
291 }
292 if let Ok(event) = serde_json::from_slice::<serde_json::Value>(line) {
293 append_log_event(log_path, event)
294 } else {
295 append_log_event(
296 log_path,
297 serde_json::json!({
298 "type": "lane_log",
299 "stream": stream,
300 "message": String::from_utf8_lossy(line),
301 }),
302 )
303 }
304 }
305
306 fn stream_child_output(
307 reader: impl std::io::Read,
308 log_path: PathBuf,
309 stream: &'static str,
310 ) -> Result<()> {
311 let mut reader = BufReader::new(reader);
312 let mut line = Vec::new();
313 loop {
314 line.clear();
315 let mut reached_eof = false;
316 while line.len() < MAX_CHILD_LOG_FRAME_BYTES {
317 let buffer = reader
318 .fill_buf()
319 .with_context(|| format!("read lane child {stream}"))?;
320 if buffer.is_empty() {
321 reached_eof = true;
322 break;
323 }
324 let available = buffer
325 .iter()
326 .position(|byte| *byte == b'\n')
327 .map_or(buffer.len(), |position| position + 1);
328 let take = available.min(MAX_CHILD_LOG_FRAME_BYTES - line.len());
329 let ended_line = buffer.get(take.saturating_sub(1)) == Some(&b'\n');
330 line.extend_from_slice(&buffer[..take]);
331 reader.consume(take);
332 if ended_line {
333 break;
334 }
335 }
336 if line.is_empty() && reached_eof {
337 return Ok(());
338 }
339 append_child_output(&log_path, stream, &line)?;
340 if reached_eof {
341 return Ok(());
342 }
343 }
344 }
345
346 fn read_lane_environment(path: &Path) -> Result<Vec<(String, String)>> {
347 let metadata = fs::metadata(path).with_context(|| format!("stat {}", path.display()))?;
348 if metadata.len() > MAX_ENVIRONMENT_BYTES {
349 bail!(
350 "private lane environment {} exceeds {} bytes",
351 path.display(),
352 MAX_ENVIRONMENT_BYTES
353 );
354 }
355 let bytes = fs::read(path).with_context(|| format!("read {}", path.display()))?;
356 if bytes.len() as u64 > MAX_ENVIRONMENT_BYTES {
357 bail!(
358 "private lane environment {} exceeds {} bytes",
359 path.display(),
360 MAX_ENVIRONMENT_BYTES
361 );
362 }
363 let environment: Vec<(String, String)> =
364 serde_json::from_slice(&bytes).with_context(|| format!("parse {}", path.display()))?;
365 for (key, _) in &environment {
366 if !valid_environment_key(key) {
367 bail!("invalid lane environment key {key:?}");
368 }
369 }
370 Ok(environment)
371 }
372
373 fn write_lane_exit_receipt(
374 receipt_path: &Path,
375 receipt_tmp_path: &Path,
376 lane_id: &str,
377 exit_code: i32,
378 ) -> Result<()> {
379 let encoded = serde_json::to_vec(&LaneExitReceipt {
380 lane_id: lane_id.to_string(),
381 exit_code,
382 })
383 .context("serialize lane exit receipt")?;
384 if encoded.len() as u64 > MAX_EXIT_RECEIPT_BYTES {
385 bail!("serialized lane exit receipt exceeds size bound");
386 }
387 remove_file_if_present(receipt_tmp_path)?;
388 let mut options = OpenOptions::new();
389 options.create_new(true).write(true);
390 #[cfg(unix)]
391 {
392 use std::os::unix::fs::OpenOptionsExt;
393 options.mode(0o600);
394 }
395 let result = (|| {
396 let mut file = options
397 .open(receipt_tmp_path)
398 .with_context(|| format!("create {}", receipt_tmp_path.display()))?;
399 file.write_all(&encoded)
400 .with_context(|| format!("write {}", receipt_tmp_path.display()))?;
401 file.sync_all()
402 .with_context(|| format!("sync {}", receipt_tmp_path.display()))?;
403 fs::rename(receipt_tmp_path, receipt_path)
404 .with_context(|| format!("publish {}", receipt_path.display()))?;
405 Ok(())
406 })();
407 if result.is_err() {
408 let _ = remove_file_if_present(receipt_tmp_path);
409 }
410 result
411 }
412
413 fn exit_status_code(status: ExitStatus) -> i32 {
414 if let Some(code) = status.code() {
415 return code;
416 }
417 #[cfg(unix)]
418 {
419 use std::os::unix::process::ExitStatusExt;
420 if let Some(signal) = status.signal() {
421 return 128 + signal;
422 }
423 }
424 LANE_PROXY_FAILURE_EXIT_CODE
425 }
426
427 fn append_proxy_failure(log_path: &Path, lane_id: &str, error: &anyhow::Error) -> Result<()> {
428 append_log_event(
429 log_path,
430 serde_json::json!({
431 "type": "lane_proxy_error",
432 "lane_id": lane_id,
433 "error": format!("{error:#}"),
434 }),
435 )
436 }
437
438 /// Run a child command while framing all output as NDJSON and atomically
439 /// publishing a private exit receipt. Returns the process-style exit code the
440 /// hidden CLI should propagate to tmux.
441 pub fn run_lane_log_proxy(spec: LaneLogProxySpec) -> Result<i32> {
442 let LaneLogProxySpec {
443 command,
444 log_path,
445 receipt_path,
446 receipt_tmp_path,
447 environment_path,
448 lane_id,
449 } = spec;
450
451 let fail = |error: anyhow::Error, exit_code: i32| -> Result<i32> {
452 append_proxy_failure(&log_path, &lane_id, &error)?;
453 write_lane_exit_receipt(&receipt_path, &receipt_tmp_path, &lane_id, exit_code)?;
454 Ok(exit_code)
455 };
456
457 if command.is_empty() {
458 return fail(anyhow::anyhow!("lane log proxy requires a command"), 127);
459 }
460
461 let environment = if let Some(path) = environment_path.as_deref() {
462 let loaded = read_lane_environment(path);
463 let removed = remove_file_if_present(path);
464 match (loaded, removed) {
465 (Ok(environment), Ok(())) => environment,
466 (Err(error), Ok(())) => return fail(error, LANE_PROXY_FAILURE_EXIT_CODE),
467 (Ok(_), Err(error)) | (Err(_), Err(error)) => return Err(error),
468 }
469 } else {
470 Vec::new()
471 };
472
473 let mut child_command = Command::new(&command[0]);
474 child_command.args(&command[1..]);
475 child_command.envs(environment);
476 child_command.stdout(Stdio::piped()).stderr(Stdio::piped());
477 let mut child = match child_command.spawn() {
478 Ok(child) => child,
479 Err(error) => {
480 return fail(
481 anyhow::Error::from(error).context(format!("spawn lane child command {command:?}")),
482 127,
483 );
484 }
485 };
486 let stdout = match child.stdout.take() {
487 Some(stdout) => stdout,
488 None => return fail(anyhow::anyhow!("lane child stdout was not piped"), 125),
489 };
490 let stderr = match child.stderr.take() {
491 Some(stderr) => stderr,
492 None => return fail(anyhow::anyhow!("lane child stderr was not piped"), 125),
493 };
494 let stdout_log = log_path.clone();
495 let stderr_log = log_path.clone();
496 let stdout_thread = thread::spawn(move || stream_child_output(stdout, stdout_log, "stdout"));
497 let stderr_thread = thread::spawn(move || stream_child_output(stderr, stderr_log, "stderr"));
498
499 let status = match child.wait() {
500 Ok(status) => status,
501 Err(error) => {
502 return fail(
503 anyhow::Error::from(error).context(format!("wait for lane child {command:?}")),
504 LANE_PROXY_FAILURE_EXIT_CODE,
505 );
506 }
507 };
508 let stdout_result = stdout_thread
509 .join()
510 .map_err(|_| anyhow::anyhow!("lane proxy stdout logger panicked"))?;
511 let stderr_result = stderr_thread
512 .join()
513 .map_err(|_| anyhow::anyhow!("lane proxy stderr logger panicked"))?;
514 let mut exit_code = exit_status_code(status);
515 if let Some(error) = stdout_result.err().or_else(|| stderr_result.err()) {
516 append_proxy_failure(&log_path, &lane_id, &error)?;
517 exit_code = LANE_PROXY_FAILURE_EXIT_CODE;
518 }
519 write_lane_exit_receipt(&receipt_path, &receipt_tmp_path, &lane_id, exit_code)?;
520 Ok(exit_code)
521 }
522
523 fn read_lane_exit_receipt(log_path: &Path, lane_id: &str) -> Result<Option<LaneExitReceipt>> {
524 let path = lane_exit_receipt_path(log_path);
525 let metadata = match std::fs::metadata(&path) {
526 Ok(metadata) => metadata,
527 Err(err) if err.kind() == std::io::ErrorKind::NotFound => return Ok(None),
528 Err(err) => return Err(err).with_context(|| format!("stat {}", path.display())),
529 };
530 if metadata.len() > MAX_EXIT_RECEIPT_BYTES {
531 bail!(
532 "lane exit receipt {} exceeds {} bytes",
533 path.display(),
534 MAX_EXIT_RECEIPT_BYTES
535 );
536 }
537 let bytes = std::fs::read(&path).with_context(|| format!("read {}", path.display()))?;
538 let receipt: LaneExitReceipt =
539 serde_json::from_slice(&bytes).with_context(|| format!("parse {}", path.display()))?;
540 if receipt.lane_id != lane_id {
541 bail!(
542 "lane exit receipt {} belongs to {}, expected {}",
543 path.display(),
544 receipt.lane_id,
545 lane_id
546 );
547 }
548 Ok(Some(receipt))
549 }
550
551 #[derive(Debug, Clone, Copy, PartialEq, Eq)]
552 enum TmuxSessionState {
553 Present,
554 Absent,
555 }
556
557 fn tmux_command(socket: &Path) -> Command {
558 let mut command = Command::new("tmux");
559 command.arg("-S").arg(socket);
560 command
561 }
562
563 fn ensure_tmux_available() -> Result<()> {
564 let output = Command::new("tmux")
565 .arg("-V")
566 .output()
567 .context("tmux runtime requires the `tmux` executable")?;
568 if !output.status.success() {
569 bail!(
570 "tmux runtime is unavailable: `tmux -V` failed with {}: {}",
571 output.status,
572 String::from_utf8_lossy(&output.stderr).trim()
573 );
574 }
575 Ok(())
576 }
577
578 fn tmux_session_state(socket: &Path, session: &str) -> Result<TmuxSessionState> {
579 let output = tmux_command(socket)
580 .args(["has-session", "-t", session])
581 .stdout(Stdio::null())
582 .stderr(Stdio::piped())
583 .output()
584 .with_context(|| format!("query tmux session {session}"))?;
585 if output.status.success() {
586 return Ok(TmuxSessionState::Present);
587 }
588 let stderr = String::from_utf8_lossy(&output.stderr).to_ascii_lowercase();
589 if stderr.contains("can't find session:")
590 || stderr.contains("no server running on")
591 || (stderr.contains("error connecting to") && stderr.contains("no such file or directory"))
592 {
593 return Ok(TmuxSessionState::Absent);
594 }
595 bail!(
596 "tmux has-session for {session} failed with {}: {}",
597 output.status,
598 stderr.trim()
599 )
600 }
601
602 fn stop_tmux_session(socket: &Path, session: &str) -> Result<()> {
603 let status = tmux_command(socket)
604 .args(["kill-session", "-t", session])
605 .stdout(Stdio::null())
606 .stderr(Stdio::null())
607 .status()
608 .with_context(|| format!("kill tmux session {session}"))?;
609 match tmux_session_state(socket, session).with_context(|| {
610 format!(
611 "confirm tmux session {session} on {} stopped after kill-session ({status})",
612 socket.display()
613 )
614 })? {
615 TmuxSessionState::Absent => {}
616 TmuxSessionState::Present => {
617 bail!("tmux session {session} remains active after kill-session ({status})")
618 }
619 }
620 // A nonzero kill can be benign when the process and session exited in the
621 // same instant. The explicit absence check above is the source of truth.
622 Ok(())
623 }
624
625 fn tmux_log_proxy_command(
626 proxy: &Path,
627 command: &[String],
628 log_path: &Path,
629 receipt_path: &Path,
630 receipt_tmp_path: &Path,
631 environment_path: Option<&Path>,
632 lane_id: &str,
633 ) -> String {
634 let mut argv = vec![
635 proxy.display().to_string(),
636 "lane-log-proxy".to_string(),
637 "--log-path".to_string(),
638 log_path.display().to_string(),
639 "--receipt-path".to_string(),
640 receipt_path.display().to_string(),
641 "--receipt-tmp-path".to_string(),
642 receipt_tmp_path.display().to_string(),
643 "--lane-id".to_string(),
644 lane_id.to_string(),
645 ];
646 if let Some(path) = environment_path {
647 argv.push("--environment-path".to_string());
648 argv.push(path.display().to_string());
649 }
650 argv.push("--".to_string());
651 argv.extend(command.iter().cloned());
652 format!("exec {}", shell_join(&argv))
653 }
654
655 /// The tmux dry-run hook, honoured only in this crate's own unit tests.
656 ///
657 /// It records a Running lane with no process, skips the tmux kill on stop,
658 /// and blinds reconcile, so an inherited `CODEWHALE_LANE_TMUX_DRY_RUN` in a
659 /// real build would fabricate lane state. Release builds ignore it.
660 fn tmux_dry_run() -> bool {
661 cfg!(test) && std::env::var_os("CODEWHALE_LANE_TMUX_DRY_RUN").is_some()
662 }
663
664 fn apply_worktree(record: &mut LaneRecord, spec: &LaneStartSpec) -> Result<Option<PathBuf>> {
665 let Some(wt) = spec.worktree.as_ref() else {
666 return Ok(spec.cwd.clone());
667 };
668 let provisioned = provision_worktree(wt)?;
669 record.worktree_path = Some(provisioned.path.clone());
670 record.branch = Some(provisioned.branch.clone());
671 Ok(Some(provisioned.path))
672 }
673
674 /// Durable local tmux sessions + attach + stream-json log file.
675 #[derive(Debug, Default)]
676 pub struct TmuxRuntime;
677
678 impl RuntimeBackend for TmuxRuntime {
679 fn kind(&self) -> RuntimeBackendKind {
680 RuntimeBackendKind::Tmux
681 }
682
683 fn start(
684 &self,
685 registry: &LaneRegistry,
686 record: &mut LaneRecord,
687 spec: &LaneStartSpec,
688 ) -> Result<()> {
689 if spec.command.is_empty() {
690 bail!("tmux runtime requires a non-empty command");
691 }
692 // Dry-run is an explicit test hook only. A missing/broken tmux binary
693 // must fail closed rather than persisting a fictional Running Lane.
694 let dry_run = tmux_dry_run();
695 if !dry_run && let Err(error) = ensure_tmux_available() {
696 append_log_event(
697 &record.log_path,
698 serde_json::json!({
699 "type": "lane_failed",
700 "lane_id": record.id,
701 "runtime": "tmux",
702 "error": error.to_string(),
703 }),
704 )?;
705 let _ = registry.mark_terminal_if_active(record, LaneStatus::Failed)?;
706 return Err(error);
707 }
708
709 let cwd = apply_worktree(record, spec)?;
710 let session = format!("cw-{}", record.id);
711 let socket = registry.root().join("tmux.sock");
712 record.tmux_session = Some(session.clone());
713 record.tmux_socket = Some(socket.clone());
714 record.attach_target = Some(format!(
715 "tmux -S {} attach -t {}",
716 shell_escape(&socket.display().to_string()),
717 shell_escape(&session)
718 ));
719
720 append_log_event(
721 &record.log_path,
722 serde_json::json!({
723 "type": "lane_started",
724 "lane_id": record.id,
725 "runtime": "tmux",
726 "session": session,
727 "workflow": record.workflow,
728 "fleet": record.fleet,
729 "issue": record.issue,
730 "dry_run": dry_run,
731 }),
732 )?;
733
734 if dry_run {
735 append_log_event(
736 &record.log_path,
737 serde_json::json!({
738 "type": "lane_log",
739 "message": "tmux dry-run: session recorded without spawning process",
740 "command": spec.command,
741 "cwd": cwd.as_ref().map(|p| p.display().to_string()),
742 }),
743 )?;
744 if !registry.mark_running_if_pending(record)? {
745 bail!(
746 "lane `{}` was stopped before tmux dry-run start completed",
747 record.id
748 );
749 }
750 return Ok(());
751 }
752
753 let log_proxy = spec
754 .log_proxy
755 .as_deref()
756 .context("tmux runtime requires a lane log proxy executable")?;
757
758 // Detached session: child output remains an operator journal only.
759 // Terminal control state goes through a separate bounded, atomically
760 // renamed receipt so child stdout cannot forge Lane completion.
761 let receipt_path = lane_exit_receipt_path(&record.log_path);
762 let receipt_tmp_path = lane_exit_receipt_tmp_path(&record.log_path);
763 let environment_path = lane_environment_path(&record.log_path);
764 remove_file_if_present(&receipt_path)?;
765 remove_file_if_present(&receipt_tmp_path)?;
766 remove_file_if_present(&environment_path)?;
767 let environment_path = if spec.environment.is_empty() {
768 None
769 } else {
770 write_lane_environment(&environment_path, &spec.environment)?;
771 Some(environment_path)
772 };
773 let shell_cmd = tmux_log_proxy_command(
774 log_proxy,
775 &spec.command,
776 &record.log_path,
777 &receipt_path,
778 &receipt_tmp_path,
779 environment_path.as_deref(),
780 &record.id,
781 );
782
783 let mut cmd = tmux_command(&socket);
784 cmd.args(["new-session", "-d", "-s", &session]);
785 if let Some(cwd) = cwd.as_ref() {
786 cmd.arg("-c").arg(cwd);
787 }
788 cmd.arg(shell_cmd);
789 let proposed_record = record.clone();
790 let spawned = std::cell::Cell::new(false);
791 let rolled_back = std::cell::Cell::new(false);
792 match registry.mark_running_if_pending_with(
793 record,
794 || {
795 let status = cmd
796 .status()
797 .with_context(|| format!("spawn tmux session {session}"))?;
798 if !status.success() {
799 bail!("tmux new-session failed with {status}");
800 }
801 spawned.set(true);
802 Ok(())
803 },
804 || {
805 stop_tmux_session(&socket, &session)?;
806 rolled_back.set(true);
807 Ok(())
808 },
809 ) {
810 Ok(true) => {}
811 Ok(false) => {
812 if let Some(path) = environment_path.as_deref() {
813 remove_file_if_present(path)?;
814 }
815 let mut stopped_record = proposed_record;
816 stopped_record.stopped_at = record.stopped_at.clone();
817 self.cleanup_worktree(&stopped_record)?;
818 bail!("lane `{}` was stopped before tmux launch", record.id);
819 }
820 Err(error) => {
821 if let Some(path) = environment_path.as_deref() {
822 let _ = remove_file_if_present(path);
823 }
824 if !spawned.get() || rolled_back.get() {
825 let _ = registry.mark_terminal_if_active(record, LaneStatus::Failed)?;
826 let mut failed_record = proposed_record;
827 failed_record.stopped_at = record.stopped_at.clone();
828 self.cleanup_worktree(&failed_record)?;
829 }
830 return Err(error);
831 }
832 }
833 Ok(())
834 }
835
836 fn attach_command(&self, record: &LaneRecord) -> Option<String> {
837 if record.status != LaneStatus::Running {
838 return None;
839 }
840 let socket = record.tmux_socket.as_ref()?;
841 let session = record.tmux_session.as_deref()?;
842 Some(format!(
843 "tmux -S {} attach -t {}",
844 shell_escape(&socket.display().to_string()),
845 shell_escape(session)
846 ))
847 }
848
849 fn stop(
850 &self,
851 registry: &LaneRegistry,
852 record: &mut LaneRecord,
853 fence: Option<u64>,
854 ) -> Result<TerminalTransition> {
855 let dry_run = tmux_dry_run();
856 let transition = registry.mark_terminal_if_active_fenced(
857 record,
858 LaneStatus::Stopped,
859 fence,
860 |current| {
861 if !dry_run {
862 match (
863 current.tmux_socket.as_deref(),
864 current.tmux_session.as_deref(),
865 ) {
866 (Some(socket), Some(session)) => stop_tmux_session(socket, session)?,
867 _ if current.status == LaneStatus::Running => {
868 bail!(
869 "running tmux lane `{}` has incomplete pinned session metadata; refusing unsafe stop",
870 current.id
871 );
872 }
873 _ => {}
874 }
875 }
876 Ok(())
877 },
878 )?;
879 if transition.transitioned() {
880 append_log_event(
881 &record.log_path,
882 serde_json::json!({
883 "type": "lane_stopped",
884 "lane_id": record.id,
885 "session": record.tmux_session,
886 }),
887 )?;
888 remove_file_if_present(&lane_environment_path(&record.log_path))?;
889 self.cleanup_worktree(record)?;
890 }
891 Ok(transition)
892 }
893
894 fn reconcile(&self, registry: &LaneRegistry, record: &mut LaneRecord) -> Result<bool> {
895 // Pending is the pre-launch state. Never reconcile it: a concurrent
896 // status read must not race a fast `tmux new-session` and overwrite
897 // the start path's later Running transition.
898 if record.status != LaneStatus::Running {
899 return Ok(false);
900 }
901
902 let receipt = read_lane_exit_receipt(&record.log_path, &record.id)?;
903 let (lane_status, exit_code, reason) = if let Some(receipt) = receipt {
904 (
905 if receipt.exit_code == 0 {
906 LaneStatus::Completed
907 } else {
908 LaneStatus::Failed
909 },
910 Some(receipt.exit_code),
911 "process_exit_receipt",
912 )
913 } else {
914 if tmux_dry_run() {
915 return Ok(false);
916 }
917 let (Some(socket), Some(session)) = (
918 record.tmux_socket.as_deref(),
919 record.tmux_session.as_deref(),
920 ) else {
921 bail!(
922 "running tmux lane `{}` lacks pinned socket/session metadata; refusing unsafe reconciliation",
923 record.id
924 );
925 };
926 match tmux_session_state(socket, session)? {
927 TmuxSessionState::Present => return Ok(false),
928 TmuxSessionState::Absent => {}
929 }
930 (
931 LaneStatus::Failed,
932 None,
933 "tmux_session_missing_without_exit_receipt",
934 )
935 };
936
937 if registry.mark_terminal_if_active(record, lane_status)? {
938 append_log_event(
939 &record.log_path,
940 serde_json::json!({
941 "type": "lane_reconciled",
942 "lane_id": record.id,
943 "exit_code": exit_code,
944 "status": lane_status.as_str(),
945 "reason": reason,
946 }),
947 )?;
948 remove_file_if_present(&lane_environment_path(&record.log_path))?;
949 self.cleanup_worktree(record)?;
950 return Ok(true);
951 }
952 Ok(false)
953 }
954 }
955
956 /// In-process / local command runtime (no tmux). Used for tests and headless.
957 #[derive(Debug, Default)]
958 pub struct InlineRuntime;
959
960 impl RuntimeBackend for InlineRuntime {
961 fn kind(&self) -> RuntimeBackendKind {
962 RuntimeBackendKind::Inline
963 }
964
965 fn start(
966 &self,
967 registry: &LaneRegistry,
968 record: &mut LaneRecord,
969 spec: &LaneStartSpec,
970 ) -> Result<()> {
971 if spec.command.is_empty() {
972 bail!("inline runtime requires a non-empty command");
973 }
974 let cwd = apply_worktree(record, spec)?;
975 // A concurrent stop that wins `mark_running_if_pending` replaces
976 // `record` with the stored one, which never saw the worktree; keep
977 // this copy so the worktree is still cleaned up, as tmux does.
978 let proposed_record = record.clone();
979 append_log_event(
980 &record.log_path,
981 serde_json::json!({
982 "type": "lane_started",
983 "lane_id": record.id,
984 "runtime": "inline",
985 "command": spec.command,
986 }),
987 )?;
988 if !registry.mark_running_if_pending(record)? {
989 let mut stopped_record = proposed_record;
990 stopped_record.stopped_at = record.stopped_at.clone();
991 self.cleanup_worktree(&stopped_record)?;
992 bail!(
993 "lane `{}` was stopped before inline start completed",
994 record.id
995 );
996 }
997
998 let mut cmd = Command::new(&spec.command[0]);
999 if spec.command.len() > 1 {
1000 cmd.args(&spec.command[1..]);
1001 }
1002 if let Some(cwd) = cwd.as_ref() {
1003 cmd.current_dir(cwd);
1004 }
1005 cmd.envs(spec.environment.iter().map(|(key, value)| (key, value)));
1006 cmd.stdout(Stdio::piped()).stderr(Stdio::piped());
1007 let mut child = match cmd.spawn() {
1008 Ok(child) => child,
1009 Err(err) => {
1010 append_log_event(
1011 &record.log_path,
1012 serde_json::json!({
1013 "type": "lane_failed",
1014 "lane_id": record.id,
1015 "error": err.to_string(),
1016 }),
1017 )?;
1018 let _ = registry.mark_terminal_if_active(record, LaneStatus::Failed)?;
1019 // Terminal now, so no later stop will clean the worktree up.
1020 self.cleanup_worktree(record)?;
1021 return Err(err).with_context(|| format!("run inline command {:?}", spec.command));
1022 }
1023 };
1024 let stdout = child
1025 .stdout
1026 .take()
1027 .context("inline lane child stdout was not piped")?;
1028 let stderr = child
1029 .stderr
1030 .take()
1031 .context("inline lane child stderr was not piped")?;
1032 let stdout_log = record.log_path.clone();
1033 let stderr_log = record.log_path.clone();
1034 let stdout_thread =
1035 thread::spawn(move || stream_child_output(stdout, stdout_log, "stdout"));
1036 let stderr_thread =
1037 thread::spawn(move || stream_child_output(stderr, stderr_log, "stderr"));
1038 let status = child
1039 .wait()
1040 .with_context(|| format!("wait for inline command {:?}", spec.command))?;
1041 let stdout_result = stdout_thread
1042 .join()
1043 .map_err(|_| anyhow::anyhow!("inline stdout logger panicked"))?;
1044 let stderr_result = stderr_thread
1045 .join()
1046 .map_err(|_| anyhow::anyhow!("inline stderr logger panicked"))?;
1047 let logging_error = stdout_result.err().or_else(|| stderr_result.err());
1048
1049 if status.success() && logging_error.is_none() {
1050 append_log_event(
1051 &record.log_path,
1052 serde_json::json!({"type": "lane_completed", "lane_id": record.id}),
1053 )?;
1054 let _ = registry.mark_terminal_if_active(record, LaneStatus::Completed)?;
1055 } else {
1056 append_log_event(
1057 &record.log_path,
1058 serde_json::json!({
1059 "type": "lane_failed",
1060 "lane_id": record.id,
1061 "status": format!("{status}"),
1062 "logging_error": logging_error.as_ref().map(ToString::to_string),
1063 }),
1064 )?;
1065 let _ = registry.mark_terminal_if_active(record, LaneStatus::Failed)?;
1066 }
1067 if let Some(err) = logging_error {
1068 return Err(err).context("stream inline lane output");
1069 }
1070 Ok(())
1071 }
1072
1073 fn attach_command(&self, _record: &LaneRecord) -> Option<String> {
1074 None
1075 }
1076
1077 fn stop(
1078 &self,
1079 registry: &LaneRegistry,
1080 record: &mut LaneRecord,
1081 fence: Option<u64>,
1082 ) -> Result<TerminalTransition> {
1083 let transition = registry.mark_terminal_if_active_fenced(
1084 record,
1085 LaneStatus::Stopped,
1086 fence,
1087 |current| {
1088 if current.status == LaneStatus::Running {
1089 bail!(
1090 "inline lane `{}` cannot be stopped safely from another process",
1091 current.id
1092 );
1093 }
1094 Ok(())
1095 },
1096 )?;
1097 if transition.transitioned() {
1098 self.cleanup_worktree(record)?;
1099 }
1100 Ok(transition)
1101 }
1102 }
1103
1104 /// Backend for a lane record written by an earlier release under the retired
1105 /// `vm` / `ci` runtimes. It never starts anything: [`RuntimeBackendKind::parse`]
1106 /// cannot produce these kinds, so only historical records reach it, and they
1107 /// can still be listed, reconciled, and stopped.
1108 #[derive(Debug)]
1109 struct RetiredRuntime {
1110 kind: RuntimeBackendKind,
1111 }
1112
1113 impl RuntimeBackend for RetiredRuntime {
1114 fn kind(&self) -> RuntimeBackendKind {
1115 self.kind
1116 }
1117
1118 fn start(
1119 &self,
1120 _registry: &LaneRegistry,
1121 _record: &mut LaneRecord,
1122 _spec: &LaneStartSpec,
1123 ) -> Result<()> {
1124 bail!(
1125 "{} runtime is not available; use tmux or inline",
1126 self.kind.as_str()
1127 )
1128 }
1129
1130 fn attach_command(&self, _record: &LaneRecord) -> Option<String> {
1131 None
1132 }
1133
1134 fn stop(
1135 &self,
1136 registry: &LaneRegistry,
1137 record: &mut LaneRecord,
1138 fence: Option<u64>,
1139 ) -> Result<TerminalTransition> {
1140 registry.mark_terminal_if_active_fenced(record, LaneStatus::Stopped, fence, |_| Ok(()))
1141 }
1142 }
1143
1144 fn shell_escape(s: &str) -> String {
1145 format!("'{}'", s.replace('\'', "'\\''"))
1146 }
1147
1148 fn shell_join(args: &[String]) -> String {
1149 args.iter()
1150 .map(|a| shell_escape(a))
1151 .collect::<Vec<_>>()
1152 .join(" ")
1153 }
1154
1155 #[cfg(test)]
1156 mod tests {
1157 use super::*;
1158 use std::ffi::OsString;
1159 use std::sync::{Mutex, MutexGuard, OnceLock};
1160 use tempfile::tempdir;
1161
1162 fn tmux_env_lock() -> MutexGuard<'static, ()> {
1163 static LOCK: OnceLock<Mutex<()>> = OnceLock::new();
1164 LOCK.get_or_init(|| Mutex::new(()))
1165 .lock()
1166 .unwrap_or_else(|poisoned| poisoned.into_inner())
1167 }
1168
1169 struct ScopedEnvVar {
1170 name: &'static str,
1171 previous: Option<OsString>,
1172 }
1173
1174 impl ScopedEnvVar {
1175 fn set(name: &'static str, value: &std::ffi::OsStr) -> Self {
1176 let previous = std::env::var_os(name);
1177 // SAFETY: tmux environment tests hold `tmux_env_lock` and restore
1178 // the process environment in Drop.
1179 unsafe { std::env::set_var(name, value) };
1180 Self { name, previous }
1181 }
1182
1183 #[cfg(unix)]
1184 fn remove(name: &'static str) -> Self {
1185 let previous = std::env::var_os(name);
1186 // SAFETY: tmux environment tests hold `tmux_env_lock` and restore
1187 // the process environment in Drop.
1188 unsafe { std::env::remove_var(name) };
1189 Self { name, previous }
1190 }
1191 }
1192
1193 impl Drop for ScopedEnvVar {
1194 fn drop(&mut self) {
1195 // SAFETY: paired with the serialized mutation above.
1196 unsafe {
1197 if let Some(previous) = self.previous.take() {
1198 std::env::set_var(self.name, previous);
1199 } else {
1200 std::env::remove_var(self.name);
1201 }
1202 }
1203 }
1204 }
1205
1206 #[test]
1207 fn tmux_dry_run_start_attach_stop_roundtrip() {
1208 let _env_guard = tmux_env_lock();
1209 let _dry_run = ScopedEnvVar::set("CODEWHALE_LANE_TMUX_DRY_RUN", std::ffi::OsStr::new("1"));
1210 let dir = tempdir().unwrap();
1211 let reg = LaneRegistry::open(dir.path()).unwrap();
1212 let mut record = reg
1213 .create_pending(
1214 Some("stopship".into()),
1215 Some("stopship".into()),
1216 Some("4375".into()),
1217 None,
1218 RuntimeBackendKind::Tmux,
1219 None,
1220 )
1221 .unwrap();
1222 let backend = TmuxRuntime;
1223 backend
1224 .start(
1225 &reg,
1226 &mut record,
1227 &LaneStartSpec {
1228 command: vec!["echo".into(), "hello-lane".into()],
1229 cwd: None,
1230 environment: Vec::new(),
1231 log_proxy: None,
1232 worktree: None,
1233 },
1234 )
1235 .unwrap();
1236 assert_eq!(record.status, LaneStatus::Running);
1237 assert!(record.tmux_session.is_some());
1238 let expected_socket = reg.root().join("tmux.sock");
1239 assert_eq!(
1240 record.tmux_socket.as_deref(),
1241 Some(expected_socket.as_path())
1242 );
1243 let attach = backend.attach_command(&record).expect("attach");
1244 assert!(attach.contains("tmux -S"));
1245 assert!(attach.contains(&expected_socket.display().to_string()));
1246 assert!(attach.contains("attach -t"));
1247 let log = std::fs::read_to_string(&record.log_path).unwrap();
1248 assert!(log.contains("lane_started"));
1249
1250 backend.stop(&reg, &mut record, None).unwrap();
1251 assert_eq!(record.status, LaneStatus::Stopped);
1252 let reloaded = reg.load(&record.id).unwrap();
1253 assert_eq!(reloaded.status, LaneStatus::Stopped);
1254 }
1255
1256 #[cfg(unix)]
1257 #[test]
1258 fn unavailable_tmux_fails_lane_instead_of_faking_running() {
1259 use std::os::unix::fs::symlink;
1260
1261 let _env_guard = tmux_env_lock();
1262 let dir = tempdir().unwrap();
1263 let bin_dir = dir.path().join("bin");
1264 fs::create_dir(&bin_dir).unwrap();
1265 symlink("/usr/bin/false", bin_dir.join("tmux")).unwrap();
1266 let prior_path = std::env::var_os("PATH").unwrap_or_default();
1267 let combined_path = std::env::join_paths(
1268 std::iter::once(bin_dir).chain(std::env::split_paths(&prior_path)),
1269 )
1270 .unwrap();
1271 let _path = ScopedEnvVar::set("PATH", &combined_path);
1272 let _dry_run = ScopedEnvVar::remove("CODEWHALE_LANE_TMUX_DRY_RUN");
1273
1274 let reg = LaneRegistry::open(dir.path().join("registry")).unwrap();
1275 let mut record = reg
1276 .create_pending(None, None, None, None, RuntimeBackendKind::Tmux, None)
1277 .unwrap();
1278 let error = TmuxRuntime
1279 .start(
1280 &reg,
1281 &mut record,
1282 &LaneStartSpec {
1283 command: vec!["/bin/true".to_string()],
1284 cwd: None,
1285 environment: Vec::new(),
1286 log_proxy: Some(PathBuf::from("/bin/true")),
1287 worktree: None,
1288 },
1289 )
1290 .unwrap_err();
1291 assert!(error.to_string().contains("tmux runtime is unavailable"));
1292 assert_eq!(record.status, LaneStatus::Failed);
1293 assert_eq!(reg.load(&record.id).unwrap().status, LaneStatus::Failed);
1294 assert!(
1295 String::from_utf8_lossy(&fs::read(&record.log_path).unwrap()).contains("lane_failed")
1296 );
1297 }
1298
1299 #[test]
1300 fn retired_remote_runtimes_are_rejected_before_a_lane_is_created() {
1301 for raw in ["vm", "ci", " VM "] {
1302 let error = RuntimeBackendKind::parse(raw).unwrap_err().to_string();
1303 assert!(error.contains("not available"), "{error}");
1304 assert!(error.contains("tmux or inline"), "{error}");
1305 }
1306 assert_eq!(
1307 RuntimeBackendKind::parse("inline").unwrap(),
1308 RuntimeBackendKind::Inline
1309 );
1310 }
1311
1312 #[test]
1313 fn historical_retired_runtime_records_still_load_and_stop() {
1314 for kind in [RuntimeBackendKind::Vm, RuntimeBackendKind::Ci] {
1315 let dir = tempdir().unwrap();
1316 let reg = LaneRegistry::open(dir.path()).unwrap();
1317 let mut record = reg
1318 .create_pending(None, None, None, None, kind, None)
1319 .unwrap();
1320 assert_eq!(reg.list().unwrap().len(), 1, "record must not be skipped");
1321 assert_eq!(reg.load(&record.id).unwrap().runtime, kind);
1322 resolve_backend(kind)
1323 .stop(&reg, &mut record, None)
1324 .expect("a retired-runtime record can still be stopped");
1325 assert!(resolve_backend(kind).attach_command(&record).is_none());
1326 }
1327 }
1328
1329 #[test]
1330 fn inline_runtime_writes_stream_json_log() {
1331 let dir = tempdir().unwrap();
1332 let reg = LaneRegistry::open(dir.path()).unwrap();
1333 let mut record = reg
1334 .create_pending(
1335 Some("demo".into()),
1336 None,
1337 None,
1338 Some("echo".into()),
1339 RuntimeBackendKind::Inline,
1340 None,
1341 )
1342 .unwrap();
1343 InlineRuntime
1344 .start(
1345 &reg,
1346 &mut record,
1347 &LaneStartSpec {
1348 command: vec!["echo".into(), "inline-ok".into()],
1349 cwd: None,
1350 environment: Vec::new(),
1351 log_proxy: None,
1352 worktree: None,
1353 },
1354 )
1355 .unwrap();
1356 assert_eq!(record.status, LaneStatus::Completed);
1357 let log = std::fs::read_to_string(&record.log_path).unwrap();
1358 assert!(log.contains("inline-ok"), "log={log}");
1359 assert!(log.contains("lane_completed"));
1360 }
1361
1362 /// Audit R06-05: an inline command that cannot spawn marks the lane
1363 /// failed and cleans up the worktree it provisioned. The lane is terminal
1364 /// at that point, so no later stop would.
1365 #[test]
1366 fn inline_spawn_failure_cleans_up_its_worktree() {
1367 let dir = tempdir().unwrap();
1368 let repo = dir.path().join("repo");
1369 std::fs::create_dir_all(&repo).unwrap();
1370 crate::worktree::tests::init_repo(&repo);
1371 let wt_path = dir.path().join("wt-inline");
1372 let reg = LaneRegistry::open(dir.path().join("lanes")).unwrap();
1373 let mut record = reg
1374 .create_pending(None, None, None, None, RuntimeBackendKind::Inline, Some(0))
1375 .unwrap();
1376 let result = InlineRuntime.start(
1377 &reg,
1378 &mut record,
1379 &LaneStartSpec {
1380 command: vec!["/nonexistent/codewhale-lane-test-binary".into()],
1381 cwd: None,
1382 environment: Vec::new(),
1383 log_proxy: None,
1384 worktree: Some(WorktreeProvision {
1385 repo_root: repo,
1386 branch: "codex/lane-spawn-fail".into(),
1387 path: wt_path.clone(),
1388 base_ref: Some("main".into()),
1389 }),
1390 },
1391 );
1392 assert!(result.is_err());
1393 assert_eq!(record.status, LaneStatus::Failed);
1394 assert!(
1395 !wt_path.exists(),
1396 "the failed lane's worktree was left behind"
1397 );
1398 }
1399
1400 #[test]
1401 fn lane_start_spec_debug_redacts_environment_values() {
1402 let spec = LaneStartSpec {
1403 command: vec!["codewhale-tui".into()],
1404 cwd: None,
1405 environment: vec![("DEEPSEEK_API_KEY".into(), "secret-value".into())],
1406 log_proxy: None,
1407 worktree: None,
1408 };
1409 let rendered = format!("{spec:?}");
1410 assert!(rendered.contains("DEEPSEEK_API_KEY"));
1411 assert!(!rendered.contains("secret-value"));
1412 }
1413
1414 #[cfg(unix)]
1415 #[test]
1416 fn inline_runtime_persists_typed_receipts_before_process_exit() {
1417 let dir = tempdir().unwrap();
1418 let reg = LaneRegistry::open(dir.path()).unwrap();
1419 let record = reg
1420 .create_pending(
1421 Some("demo".into()),
1422 None,
1423 None,
1424 None,
1425 RuntimeBackendKind::Inline,
1426 None,
1427 )
1428 .unwrap();
1429 let log_path = record.log_path.clone();
1430 let (done_tx, done_rx) = std::sync::mpsc::channel();
1431 let handle = thread::spawn(move || {
1432 let mut record = record;
1433 let result = InlineRuntime.start(
1434 &reg,
1435 &mut record,
1436 &LaneStartSpec {
1437 command: vec![
1438 "sh".into(),
1439 "-c".into(),
1440 "printf '%s\\n' '{\"type\":\"workflow_event\",\"run_id\":\"workflow_live\",\"event\":{\"type\":\"run_started\"}}'; sleep 0.25; printf done"
1441 .into(),
1442 ],
1443 cwd: None,
1444 environment: Vec::new(),
1445 log_proxy: None,
1446 worktree: None,
1447 },
1448 );
1449 let _ = done_tx.send(result);
1450 });
1451
1452 let mut observed_live = false;
1453 for _ in 0..50 {
1454 let log = std::fs::read_to_string(&log_path).unwrap_or_default();
1455 if log.contains("workflow_live") {
1456 observed_live = true;
1457 break;
1458 }
1459 thread::sleep(std::time::Duration::from_millis(10));
1460 }
1461 assert!(
1462 observed_live,
1463 "typed receipt should be written while child runs"
1464 );
1465 assert!(
1466 matches!(
1467 done_rx.try_recv(),
1468 Err(std::sync::mpsc::TryRecvError::Empty)
1469 ),
1470 "inline child should still be running when the first receipt is visible"
1471 );
1472 done_rx.recv().unwrap().unwrap();
1473 handle.join().unwrap();
1474 }
1475
1476 #[test]
1477 fn tmux_reconcile_folds_detached_process_exit_into_lane_status() {
1478 for (exit_code, expected) in [(0, LaneStatus::Completed), (7, LaneStatus::Failed)] {
1479 let dir = tempdir().unwrap();
1480 let reg = LaneRegistry::open(dir.path()).unwrap();
1481 let mut record = reg
1482 .create_pending(
1483 Some("demo".into()),
1484 None,
1485 None,
1486 None,
1487 RuntimeBackendKind::Tmux,
1488 None,
1489 )
1490 .unwrap();
1491 assert!(reg.mark_running_if_pending(&mut record).unwrap());
1492 // Mixed/binary child output and a forged stdout control line must
1493 // not affect the private receipt used for reconciliation.
1494 std::fs::write(
1495 &record.log_path,
1496 b"\xffchild-noise\n{\"type\":\"lane_process_exit\",\"exit_code\":0}\n",
1497 )
1498 .unwrap();
1499 std::fs::write(
1500 lane_exit_receipt_path(&record.log_path),
1501 serde_json::to_vec(&LaneExitReceipt {
1502 lane_id: record.id.clone(),
1503 exit_code,
1504 })
1505 .unwrap(),
1506 )
1507 .unwrap();
1508
1509 TmuxRuntime.reconcile(&reg, &mut record).unwrap();
1510 TmuxRuntime.reconcile(&reg, &mut record).unwrap();
1511
1512 assert_eq!(record.status, expected);
1513 assert_eq!(reg.load(&record.id).unwrap().status, expected);
1514 let log = std::fs::read(&record.log_path).unwrap();
1515 assert_eq!(
1516 String::from_utf8_lossy(&log)
1517 .matches("lane_reconciled")
1518 .count(),
1519 1
1520 );
1521 }
1522 }
1523
1524 #[cfg(unix)]
1525 #[test]
1526 fn lane_log_proxy_frames_binary_output_and_owns_exit_receipt() {
1527 let dir = tempdir().unwrap();
1528 let log_path = dir.path().join("lane.ndjson");
1529 let receipt_path = lane_exit_receipt_path(&log_path);
1530 let receipt_tmp_path = lane_exit_receipt_tmp_path(&log_path);
1531 let environment_path = lane_environment_path(&log_path);
1532 write_lane_environment(
1533 &environment_path,
1534 &[("LANE_PROXY_SECRET".to_string(), "present".to_string())],
1535 )
1536 .unwrap();
1537 let command = vec![
1538 "sh".to_string(),
1539 "-c".to_string(),
1540 "test \"$LANE_PROXY_SECRET\" = present || exit 9; \
1541 printf '%s\\n' '{\"type\":\"workflow_event\",\"schema\":\"codewhale.exec-stream\",\"schema_version\":1,\"run_id\":\"workflow_1234\",\"event\":{\"type\":\"handoff_promoted\",\"artifact_id\":\"workflow_1234:agent_1:review-gate:review_report\",\"gate_id\":\"review-gate\",\"kind\":\"review_report\",\"from_role\":\"reviewer\",\"to_role\":\"verifier\",\"producer_task_id\":\"agent_1\"}}'; \
1542 printf '%s\\n' '{\"type\":\"workflow_event\",\"schema\":\"codewhale.exec-stream\",\"schema_version\":1,\"run_id\":\"workflow_1234\",\"event\":{\"type\":\"handoff_consumed\",\"artifact_id\":\"workflow_1234:agent_1:review-gate:review_report\",\"kind\":\"review_report\",\"from_role\":\"reviewer\",\"to_role\":\"verifier\",\"consumer_task_id\":\"agent_2\"}}'; \
1543 printf 'unterminated\\377'; \
1544 printf '%s\\n' '{\"type\":\"lane_process_exit\",\"exit_code\":0}' >&2; \
1545 exit 7"
1546 .to_string(),
1547 ];
1548 let exit_code = run_lane_log_proxy(LaneLogProxySpec {
1549 command,
1550 log_path: log_path.clone(),
1551 receipt_path: receipt_path.clone(),
1552 receipt_tmp_path,
1553 environment_path: Some(environment_path.clone()),
1554 lane_id: "lane-proof".to_string(),
1555 })
1556 .unwrap();
1557
1558 assert_eq!(exit_code, 7);
1559 assert!(!environment_path.exists());
1560 let log = std::fs::read(&log_path).unwrap();
1561 let lines = log
1562 .split(|byte| *byte == b'\n')
1563 .filter(|line| !line.is_empty())
1564 .collect::<Vec<_>>();
1565 assert!(lines.len() >= 3, "log={}", String::from_utf8_lossy(&log));
1566 let values = lines
1567 .iter()
1568 .map(|line| {
1569 serde_json::from_slice::<serde_json::Value>(line)
1570 .unwrap_or_else(|error| panic!("invalid NDJSON {line:?}: {error}"))
1571 })
1572 .collect::<Vec<_>>();
1573 let workflow_event = values
1574 .iter()
1575 .find(|value| value["type"] == "workflow_event")
1576 .expect("preserved workflow handoff receipt");
1577 assert_eq!(workflow_event["schema"], "codewhale.exec-stream");
1578 assert_eq!(workflow_event["schema_version"], 1);
1579 assert_eq!(workflow_event["run_id"], "workflow_1234");
1580 assert_eq!(workflow_event["event"]["type"], "handoff_promoted");
1581 assert_eq!(
1582 workflow_event["event"]["artifact_id"],
1583 "workflow_1234:agent_1:review-gate:review_report"
1584 );
1585 assert_eq!(workflow_event["event"]["gate_id"], "review-gate");
1586 assert_eq!(workflow_event["event"]["kind"], "review_report");
1587 assert_eq!(workflow_event["event"]["from_role"], "reviewer");
1588 assert_eq!(workflow_event["event"]["to_role"], "verifier");
1589 assert_eq!(workflow_event["event"]["producer_task_id"], "agent_1");
1590 assert!(workflow_event["event"].get("payload").is_none());
1591 let consumed_event = values
1592 .iter()
1593 .find(|value| value["event"]["type"] == "handoff_consumed")
1594 .expect("preserved workflow handoff consumption receipt");
1595 assert_eq!(consumed_event["schema"], "codewhale.exec-stream");
1596 assert_eq!(consumed_event["schema_version"], 1);
1597 assert_eq!(consumed_event["run_id"], "workflow_1234");
1598 assert_eq!(
1599 consumed_event["event"]["artifact_id"],
1600 workflow_event["event"]["artifact_id"]
1601 );
1602 assert_eq!(consumed_event["event"]["kind"], "review_report");
1603 assert_eq!(consumed_event["event"]["from_role"], "reviewer");
1604 assert_eq!(consumed_event["event"]["to_role"], "verifier");
1605 assert_eq!(consumed_event["event"]["consumer_task_id"], "agent_2");
1606 assert!(consumed_event["event"].get("payload").is_none());
1607 let rendered = String::from_utf8_lossy(&log);
1608 assert!(rendered.contains("workflow_event"));
1609 assert!(rendered.contains("lane_process_exit"));
1610 assert!(rendered.contains("lane_log"));
1611 let receipt = read_lane_exit_receipt(&log_path, "lane-proof")
1612 .unwrap()
1613 .expect("private receipt");
1614 assert_eq!(receipt.exit_code, 7);
1615 }
1616
1617 #[test]
1618 fn invalid_environment_key_never_leaves_a_partial_secret_file() {
1619 let dir = tempdir().unwrap();
1620 let path = dir.path().join("lane.env.json");
1621 let result = write_lane_environment(
1622 &path,
1623 &[
1624 ("VALID_KEY".to_string(), "first-secret".to_string()),
1625 ("INVALID-KEY".to_string(), "second-secret".to_string()),
1626 ],
1627 );
1628 assert!(result.is_err());
1629 assert!(!path.exists());
1630 assert!(!lane_environment_tmp_path(&path).exists());
1631 }
1632
1633 #[cfg(unix)]
1634 #[test]
1635 fn private_environment_is_mode_0600_and_malformed_input_is_removed() {
1636 use std::os::unix::fs::PermissionsExt;
1637
1638 let dir = tempdir().unwrap();
1639 let path = dir.path().join("lane.env.json");
1640 write_lane_environment(&path, &[("SECRET".to_string(), "sensitive".to_string())]).unwrap();
1641 assert_eq!(
1642 fs::metadata(&path).unwrap().permissions().mode() & 0o777,
1643 0o600
1644 );
1645
1646 fs::write(&path, b"{malformed").unwrap();
1647 let log_path = dir.path().join("lane.ndjson");
1648 let exit_code = run_lane_log_proxy(LaneLogProxySpec {
1649 command: vec!["/bin/true".to_string()],
1650 receipt_path: lane_exit_receipt_path(&log_path),
1651 receipt_tmp_path: lane_exit_receipt_tmp_path(&log_path),
1652 environment_path: Some(path.clone()),
1653 lane_id: "lane-malformed-env".to_string(),
1654 log_path: log_path.clone(),
1655 })
1656 .unwrap();
1657 assert_eq!(exit_code, LANE_PROXY_FAILURE_EXIT_CODE);
1658 assert!(!path.exists());
1659 assert!(String::from_utf8_lossy(&fs::read(log_path).unwrap()).contains("lane_proxy_error"));
1660 }
1661
1662 #[cfg(unix)]
1663 #[test]
1664 fn tmux_stop_failure_keeps_lane_running_and_preserves_cleanup_targets() {
1665 use std::os::unix::fs::symlink;
1666
1667 let _env_guard = tmux_env_lock();
1668 let dir = tempdir().unwrap();
1669 let bin_dir = dir.path().join("bin");
1670 fs::create_dir(&bin_dir).unwrap();
1671 symlink("/usr/bin/false", bin_dir.join("tmux")).unwrap();
1672 let prior_path = std::env::var_os("PATH").unwrap_or_default();
1673 let combined_path = std::env::join_paths(
1674 std::iter::once(bin_dir.clone()).chain(std::env::split_paths(&prior_path)),
1675 )
1676 .unwrap();
1677 let _path = ScopedEnvVar::set("PATH", &combined_path);
1678 let _dry_run = ScopedEnvVar::remove("CODEWHALE_LANE_TMUX_DRY_RUN");
1679
1680 let reg = LaneRegistry::open(dir.path().join("registry")).unwrap();
1681 let mut record = reg
1682 .create_pending(
1683 Some("demo".into()),
1684 None,
1685 None,
1686 None,
1687 RuntimeBackendKind::Tmux,
1688 Some(0),
1689 )
1690 .unwrap();
1691 record.tmux_session = Some(format!("cw-{}", record.id));
1692 record.tmux_socket = Some(reg.root().join("tmux.sock"));
1693 let worktree = dir.path().join("worktree");
1694 fs::create_dir(&worktree).unwrap();
1695 record.worktree_path = Some(worktree.clone());
1696 let environment_path = lane_environment_path(&record.log_path);
1697 write_lane_environment(
1698 &environment_path,
1699 &[("SECRET".to_string(), "still-private".to_string())],
1700 )
1701 .unwrap();
1702 assert!(reg.mark_running_if_pending(&mut record).unwrap());
1703
1704 let error = TmuxRuntime.stop(&reg, &mut record, None).unwrap_err();
1705 assert!(format!("{error:#}").contains("tmux has-session"));
1706 assert_eq!(record.status, LaneStatus::Running);
1707 assert_eq!(reg.load(&record.id).unwrap().status, LaneStatus::Running);
1708 assert!(environment_path.exists());
1709 assert!(worktree.exists());
1710 }
1711
1712 #[cfg(unix)]
1713 #[test]
1714 fn tmux_reconcile_marks_vanished_session_failed_without_receipt() {
1715 use std::os::unix::fs::PermissionsExt;
1716
1717 let _env_guard = tmux_env_lock();
1718 let dir = tempdir().unwrap();
1719 let bin_dir = dir.path().join("bin");
1720 fs::create_dir(&bin_dir).unwrap();
1721 let tmux = bin_dir.join("tmux");
1722 fs::write(
1723 &tmux,
1724 "#!/bin/sh\nprintf '%s\\n' 'no server running on /tmp/codewhale-test.sock' >&2\nexit 1\n",
1725 )
1726 .unwrap();
1727 let mut permissions = fs::metadata(&tmux).unwrap().permissions();
1728 permissions.set_mode(0o755);
1729 fs::set_permissions(&tmux, permissions).unwrap();
1730 let prior_path = std::env::var_os("PATH").unwrap_or_default();
1731 let combined_path = std::env::join_paths(
1732 std::iter::once(bin_dir).chain(std::env::split_paths(&prior_path)),
1733 )
1734 .unwrap();
1735 let _path = ScopedEnvVar::set("PATH", &combined_path);
1736 let _dry_run = ScopedEnvVar::remove("CODEWHALE_LANE_TMUX_DRY_RUN");
1737
1738 let reg = LaneRegistry::open(dir.path()).unwrap();
1739 let mut record = reg
1740 .create_pending(
1741 Some("demo".into()),
1742 None,
1743 None,
1744 None,
1745 RuntimeBackendKind::Tmux,
1746 None,
1747 )
1748 .unwrap();
1749 record.tmux_session = Some(format!("missing-{}", record.id));
1750 record.tmux_socket = Some(reg.root().join("tmux.sock"));
1751 assert!(reg.mark_running_if_pending(&mut record).unwrap());
1752
1753 TmuxRuntime.reconcile(&reg, &mut record).unwrap();
1754
1755 assert_eq!(record.status, LaneStatus::Failed);
1756 let log = std::fs::read_to_string(&record.log_path).unwrap();
1757 assert!(log.contains("tmux_session_missing_without_exit_receipt"));
1758 }
1759 }
1760
1760 lines RUST