use std::fs;
use std::io::{BufRead, BufReader};
use std::path::{Path, PathBuf};
use std::process::{Child, Command, Stdio};
use std::sync::mpsc;
use std::thread;
use std::time::{Duration, Instant};
fn read_uuid() -> String {
fs::read_to_string("/proc/sys/kernel/random/uuid")
.expect("failed to read UUID from /proc/sys/kernel/random/uuid (not Linux?)")
.trim()
.to_string()
}
fn has_command(name: &str) -> bool {
Command::new("which")
.arg(name)
.output()
.map(|o| o.status.success())
.unwrap_or(false)
}
fn all_pids_by_env(key: &str, value: &str) -> Vec<u32> {
let Ok(proc_dir) = fs::read_dir("/proc") else {
return vec![];
};
let needle = format!("{key}={value}\0");
let mut result = Vec::new();
for entry in proc_dir.flatten() {
let name = match entry.file_name().to_str() {
Some(s) => s.to_string(),
None => continue,
};
let pid: u32 = match name.parse() {
Ok(p) => p,
Err(_) => continue,
};
let environ_path = entry.path().join("environ");
let data = match fs::read(&environ_path) {
Ok(d) => d,
Err(_) => continue,
};
if data.windows(needle.len()).any(|w| w == needle.as_bytes()) {
result.push(pid);
}
}
result
}
fn get_cmdline_args(pid: u32) -> Option<Vec<String>> {
let data = fs::read(format!("/proc/{pid}/cmdline")).ok()?;
Some(
data.split(|&b| b == 0)
.filter_map(|s| std::str::from_utf8(s).ok())
.filter(|s| !s.is_empty())
.map(String::from)
.collect(),
)
}
fn argv_points_to_codex(args: &[String]) -> bool {
let argv0 = match args.first() {
Some(a) => a,
None => return false,
};
let base = Path::new(argv0)
.file_name()
.and_then(|f| f.to_str())
.unwrap_or("");
if matches!(base, "sh" | "bash" | "dash" | "zsh") {
return false;
}
base == "codex"
}
fn format_cmdline(args: &[String]) -> String {
let joined: Vec<&str> = args.iter().map(|s| s.as_str()).collect();
joined.join(" ")
}
fn get_env_var(data: &[u8], key: &str) -> Option<String> {
let prefix = format!("{key}=");
for entry in data.split(|&b| b == 0) {
if let Ok(s) = std::str::from_utf8(entry)
&& let Some(val) = s.strip_prefix(&prefix)
{
return Some(val.to_string());
}
}
None
}
fn parse_ready_session(line: &str) -> Option<String> {
let key = r#""zellij_session":""#;
let start = line.find(key)?;
let rest = &line[start + key.len()..];
let end = rest.find('"')?;
Some(rest[..end].to_string())
}
fn make_init_json(session_id: &str, working_dir: &str) -> String {
let wd = working_dir.replace('\\', "\\\\").replace('"', "\\\"");
format!(
r#"{{"session_id":"{}","title":"test","chat_id":"test","root_message_id":"test","working_dir":"{}","cli_id":"codex","cli_bin":"codex","cli_args":[],"prompt":"","resume":false,"lark_app_id":"local","lark_app_secret":""}}"#,
session_id, wd
)
}
fn spawn_line_reader<R: 'static + std::io::Read + Send>(
pipe: R,
label: &'static str,
) -> mpsc::Receiver<String> {
let (tx, rx) = mpsc::channel();
thread::spawn(move || {
let mut reader = BufReader::new(pipe);
let mut line = String::new();
loop {
line.clear();
match reader.read_line(&mut line) {
Ok(0) => break,
Ok(_) => {
if tx.send(line.clone()).is_err() {
break;
}
}
Err(e) => {
eprintln!("[{label}] read error: {e}");
break;
}
}
}
});
rx
}
struct Guard {
worker: Option<Child>,
zellij_session: Option<String>,
tmp_dir: Option<PathBuf>,
}
impl Drop for Guard {
fn drop(&mut self) {
if let Some(ref mut child) = self.worker {
let _ = child.kill();
let _ = child.wait();
}
if let Some(ref session) = self.zellij_session {
let _ = Command::new("zellij")
.args(["delete-session", session, "-f"])
.output();
}
if let Some(ref dir) = self.tmp_dir {
let _ = fs::remove_dir_all(dir);
}
}
}
#[test]
#[ignore = "live test: requires locally installed and authenticated `codex`, `zellij`, and Linux /proc"]
fn live_codex_term_injected_in_zellij() {
if !has_command("codex") {
eprintln!("skipping live test: `codex` not found in PATH");
return;
}
if !has_command("zellij") {
eprintln!("skipping live test: `zellij` not found in PATH");
return;
}
if !Path::new("/proc").is_dir() {
eprintln!("skipping live test: /proc not available (not Linux?)");
return;
}
let uuid = read_uuid();
let session_id = uuid; let short = &session_id[..session_id.len().min(8)];
let zellij_session_name = format!("beam-{short}");
let tmp = PathBuf::from("/tmp").join(format!("beam-test-{session_id}"));
fs::create_dir_all(&tmp).expect("failed to create temporary directory");
let mut guard = Guard {
worker: None,
zellij_session: Some(zellij_session_name.clone()),
tmp_dir: Some(tmp.clone()),
};
let init_json = make_init_json(&session_id, tmp.to_str().unwrap());
let init_path = tmp.join("init.json");
fs::write(&init_path, &init_json).expect("failed to write init config");
let worker_bin = std::env!("CARGO_BIN_EXE_beam-worker");
let mut worker = Command::new(worker_bin)
.arg("--init-path")
.arg(&init_path)
.env("TERM", "dumb") .env("BEAM_HOME", &tmp)
.stdout(Stdio::piped())
.stderr(Stdio::piped())
.stdin(Stdio::piped()) .spawn()
.unwrap_or_else(|e| panic!("failed to spawn beam-worker: {e}"));
let stderr = worker.stderr.take().expect("worker stderr not captured");
let _stderr_rx = spawn_line_reader(stderr, "worker-stderr");
let stdout = worker.stdout.take().expect("worker stdout not captured");
let stdout_rx = spawn_line_reader(stdout, "worker-stdout");
guard.worker = Some(worker);
let deadline = Instant::now() + Duration::from_secs(30);
loop {
if Instant::now() > deadline {
panic!(
"worker did not send Ready within 30 s timeout \
(session_id={session_id}, zellij_session={zellij_session_name})"
);
}
match stdout_rx.recv_timeout(Duration::from_millis(500)) {
Ok(line) => {
let trimmed = line.trim().to_string();
if trimmed.contains(r#""type":"ready""#) {
let parsed = parse_ready_session(&trimmed).unwrap_or_else(|| {
panic!("could not extract zellij_session from Ready line: {trimmed}")
});
assert_eq!(
parsed, zellij_session_name,
"Ready reports zellij_session={parsed} \
but we expected {zellij_session_name}"
);
break;
}
}
Err(mpsc::RecvTimeoutError::Timeout) => continue,
Err(mpsc::RecvTimeoutError::Disconnected) => {
panic!("worker stdout reader exited unexpectedly");
}
}
}
eprintln!("worker Ready received, zellij session={zellij_session_name}");
let deadline = Instant::now() + Duration::from_secs(15);
let codex_pid = loop {
if Instant::now() > deadline {
let all = all_pids_by_env("BEAM_SESSION_ID", &session_id);
eprintln!("[diagnostic] all PIDs with BEAM_SESSION_ID={session_id}:");
for pid in &all {
let args = get_cmdline_args(*pid).unwrap_or_default();
let verdict = if argv_points_to_codex(&args) {
"ACCEPT"
} else {
"REJECT (not codex — probably wrapper shell)"
};
eprintln!(
" pid={} cmdline=[{}] {verdict}",
pid,
format_cmdline(&args),
);
}
panic!(
"could not find a codex process with BEAM_SESSION_ID={} \
within 15 s timeout ({} candidate(s) found, none accepted)",
session_id,
all.len(),
);
}
let mut accepted = Vec::new();
for pid in all_pids_by_env("BEAM_SESSION_ID", &session_id) {
if let Some(args) = get_cmdline_args(pid)
&& argv_points_to_codex(&args)
{
accepted.push((pid, args));
}
}
if accepted.len() == 1 {
break accepted.into_iter().next().unwrap().0;
}
if accepted.len() > 1 {
let desc: Vec<String> = accepted
.iter()
.map(|(pid, args)| format!("pid={} cmdline=[{}]", pid, format_cmdline(args)))
.collect();
panic!(
"multiple ({}) codex candidates for session {}: {}",
accepted.len(),
session_id,
desc.join("; "),
);
}
std::thread::sleep(Duration::from_millis(200));
};
eprintln!("found codex pid={codex_pid}");
let environ_path = format!("/proc/{codex_pid}/environ");
let environ_data = fs::read(&environ_path)
.unwrap_or_else(|e| panic!("failed to read {environ_path}: {e} (codex exited?)"));
let term_val = get_env_var(&environ_data, "TERM")
.unwrap_or_else(|| panic!("TERM not found in environ of codex pid={codex_pid}"));
assert_eq!(
term_val, "xterm-256color",
"codex (pid={codex_pid}) should have TERM=xterm-256color \
despite worker TERM=dumb; got TERM={term_val}"
);
eprintln!("PASS: codex pid={codex_pid} TERM={term_val}");
}