use super::*;
use crate::testutil::{env_here, here, read_pid, sh_env, temp, wait_until};
fn no_waker() -> Waker {
Arc::new(Mutex::new(None))
}
fn spawn(id: u64, command: &str) -> Task {
Task::spawn(
id,
command,
command,
&here(),
24,
80,
2000,
&env_here(),
no_waker(),
)
.unwrap()
}
fn wait_finished(t: &mut Task) {
assert!(
wait_until(Duration::from_secs(5), || {
t.poll_exit().unwrap();
t.finished.is_some()
}),
"task never finished"
);
}
#[test]
fn spawn_reads_output_and_exits_zero() {
let mut t = spawn(1, "printf 'alpha\\nomega\\n'");
let mut preview = String::new();
wait_until(Duration::from_secs(5), || {
t.poll_exit().unwrap();
preview = t.resolve_preview(Instant::now()).text;
t.finished.is_some() && preview.contains("omega")
});
assert_eq!(t.exit_code, Some(0));
assert!(preview.contains("omega"), "preview was {preview:?}");
t.terminate();
}
#[test]
fn nonzero_exit_is_recorded() {
let mut t = spawn(2, "exit 3");
wait_finished(&mut t);
assert_eq!(t.exit_code, Some(3));
assert_eq!(
t.lifecycle(Instant::now(), Duration::from_secs(10)),
Lifecycle::Failed
);
t.terminate();
}
#[test]
fn lifecycle_and_parked_agree_across_the_window_edge() {
let mut t = spawn(5, "sleep 5");
let quiet_since = *t.last_activity.lock().unwrap();
let window = Duration::from_secs(10);
let inside = quiet_since + Duration::from_secs(9);
assert_eq!(t.lifecycle(inside, window), Lifecycle::Active);
assert!(!t.parked(inside, window));
let past = quiet_since + Duration::from_secs(11);
assert_eq!(t.lifecycle(past, window), Lifecycle::Idle);
assert!(t.parked(past, window));
t.terminate();
}
#[test]
fn sub_window_quiet_gaps_never_read_as_idle() {
let mut t = spawn(7, "sleep 5");
let window = Duration::from_secs(10);
let start = *t.last_activity.lock().unwrap();
for gaps in 1..=4u32 {
let probe = start + Duration::from_secs(9) * gaps;
assert_eq!(t.lifecycle(probe, window), Lifecycle::Active);
assert!(!t.parked(probe, window));
*t.last_activity.lock().unwrap() = probe;
}
t.terminate();
}
#[test]
fn finished_tasks_are_never_parked() {
let mut t = spawn(6, "exit 0");
wait_finished(&mut t);
let now = *t.last_activity.lock().unwrap() + Duration::from_secs(11);
assert!(!t.parked(now, Duration::from_secs(10)));
t.terminate();
}
#[test]
fn resize_is_reflected_in_the_grid() {
let mut t = Task::spawn(
3,
"sleep 5",
"sleep 5",
&here(),
24,
80,
2000,
&env_here(),
no_waker(),
)
.unwrap();
t.resize(30, 100).unwrap();
assert_eq!(t.parser.lock().size(), (30, 100));
t.terminate();
}
#[test]
fn exited_leader_stays_a_zombie_until_drop() {
use nix::sys::signal::kill;
let mut t = spawn(4, "exit 7");
wait_finished(&mut t);
assert_eq!(t.exit_code, Some(7));
let pid = Pid::from_raw(t.pid.expect("spawn always yields a pid") as i32);
assert!(
kill(pid, None).is_ok(),
"leader was reaped by the latch; the pgid reservation is gone"
);
drop(t);
assert!(kill(pid, None).is_err(), "Drop did not collect the zombie");
}
#[test]
fn terminate_reaches_stragglers_after_leader_exit() {
use nix::sys::signal::kill;
let dir = temp("task_straggler");
let spid = dir.join("spid");
let cmd = format!("trap '' HUP; sleep 300 & echo $! > {}", spid.display());
let mut t = Task::spawn(5, &cmd, &cmd, &here(), 24, 80, 2000, &sh_env(), no_waker()).unwrap();
wait_finished(&mut t); let straggler = read_pid(&spid);
assert!(kill(straggler, None).is_ok(), "straggler should be alive");
t.terminate(); assert!(
wait_until(Duration::from_secs(5), || kill(straggler, None).is_err()),
"TERM after leader exit never reached the straggler"
);
let _ = std::fs::remove_dir_all(&dir);
}
#[test]
fn killed_leader_latches_137_via_collect() {
let mut t = spawn(8, "sleep 300");
t.force_kill(); assert!(
wait_until(Duration::from_secs(5), || t.try_collect()),
"KILLed leader was never collected"
);
assert_eq!(t.exit_code, Some(137));
}
#[test]
fn group_gone_reaps_then_probes_past_the_zombie() {
use nix::errno::Errno;
let mut t = spawn(30, "exit 0");
wait_finished(&mut t);
let pgid = Pid::from_raw(t.pid.expect("spawn always yields a pid") as i32);
assert_ne!(
killpg(pgid, None::<Signal>),
Err(Errno::ESRCH),
"an unreaped zombie must keep the group id resolvable"
);
assert!(
t.group_gone(),
"a zombie-only group must probe gone in one reap+probe pass"
);
assert_eq!(killpg(pgid, None::<Signal>), Err(Errno::ESRCH));
}
#[test]
fn group_gone_holds_while_a_member_survives() {
use nix::sys::signal::kill;
let dir = temp("task_gone");
let spid = dir.join("spid");
let cmd = format!("trap '' HUP; sleep 300 & echo $! > {}", spid.display());
let mut t = Task::spawn(31, &cmd, &cmd, &here(), 24, 80, 2000, &sh_env(), no_waker()).unwrap();
wait_finished(&mut t);
let straggler = read_pid(&spid);
assert!(!t.group_gone(), "a surviving member must hold the probe");
assert!(t.reaped, "the probe reaps the exited leader to see past it");
let _ = kill(straggler, Signal::SIGKILL);
assert!(
wait_until(Duration::from_secs(5), || t.group_gone()),
"the group must probe gone once its last member dies"
);
let _ = std::fs::remove_dir_all(&dir);
}
#[test]
fn group_gone_never_reaps_a_live_leader() {
let mut t = spawn(32, "sleep 300");
assert!(!t.group_gone(), "a live leader is a live group");
assert!(!t.reaped, "the probe must not reap a running leader");
t.terminate();
}
#[test]
fn viewport_scrolls_and_snaps_live_on_input() {
let mut t = spawn(9, "cat");
for i in 0..50 {
grid(&t.parser).process(format!("line{i}\r\n").as_bytes());
}
assert_eq!(t.scroll_offset(), 0);
t.scroll_view(ScrollAction::Up(10));
assert_eq!(t.scroll_offset(), 10);
t.scroll_view(ScrollAction::Down(4));
assert_eq!(t.scroll_offset(), 6);
t.scroll_view(ScrollAction::Top);
let top = t.scroll_offset();
assert!(top > 0);
assert!(
t.screen_lines()[0].starts_with("line0"),
"Top must show the oldest stored row, got {:?}",
t.screen_lines()[0]
);
t.scroll_view(ScrollAction::Live);
t.scroll_view(ScrollAction::Up(10_000));
assert_eq!(t.scroll_offset(), top);
t.send_input(b"x").unwrap();
assert_eq!(t.scroll_offset(), 0);
t.terminate();
}
#[test]
fn screen_lines_yields_one_entry_per_grid_row() {
let mut t = spawn(60, "sleep 300");
grid(&t.parser).process(b"top");
let lines = t.screen_lines();
assert_eq!(lines.len(), 24, "one entry per grid row");
assert_eq!(lines[0], "top");
assert_eq!(lines[23], "", "the blank bottom row keeps its slot");
t.terminate();
}
#[test]
fn queued_writes_reach_the_child_in_order() {
let mut t = spawn(10, "cat");
t.send_input(b"zqfirstqz\n").unwrap();
t.send_input(b"zqsecondqz\n").unwrap();
let mut contents = String::new();
wait_until(Duration::from_secs(5), || {
contents = grid(&t.parser).contents();
contents.contains("zqsecondqz")
});
let first = contents
.find("zqfirstqz")
.expect("first message never echoed");
let second = contents
.find("zqsecondqz")
.expect("second message never echoed");
assert!(first < second, "queued writes reordered: {contents:?}");
t.terminate();
}
struct FailingWriter {
limit: usize,
written: usize,
}
impl Write for FailingWriter {
fn write(&mut self, buf: &[u8]) -> io::Result<usize> {
if self.written + buf.len() > self.limit {
return Err(io::Error::other("slave side closed"));
}
self.written += buf.len();
Ok(buf.len())
}
fn flush(&mut self) -> io::Result<()> {
Ok(())
}
}
#[test]
fn write_error_keeps_draining_the_pending_counter() {
let (tx, rx) = channel::<Vec<u8>>();
let pending = AtomicUsize::new(0);
let msgs: [&[u8]; 3] = [b"fits", b"fails", b"queued-behind"];
for msg in msgs {
admit_write(&tx, &pending, msg.to_vec()).unwrap();
}
let total: usize = msgs.iter().map(|m| m.len()).sum();
assert_eq!(pending.load(Ordering::Acquire), total);
drop(tx);
let mut w = FailingWriter {
limit: msgs[0].len(),
written: 0,
};
drain_writes(rx, &mut w, &pending);
assert_eq!(
pending.load(Ordering::Acquire),
0,
"accounting must survive a dead writer"
);
assert_eq!(
w.written,
msgs[0].len(),
"post-error messages must be discarded, not written"
);
}
#[test]
fn input_hints_track_child_modes() {
let mut t = spawn(8, "sleep 5");
assert_eq!(t.input_hints(), (false, false, false));
grid(&t.parser).process(b"\x1b[?1000h");
assert_eq!(t.input_hints(), (true, false, false));
grid(&t.parser).process(b"\x1b[?1000l\x1b[?1049h");
assert_eq!(t.input_hints(), (false, true, true));
grid(&t.parser).process(b"\x1b[?1007l");
assert_eq!(t.input_hints(), (false, true, false));
t.terminate();
}
#[test]
fn scrape_exit_hint_waits_for_reader_eof() {
const ID: &str = "c8c4a5cc-0b32-4ba0-a6b4-6ed08c218e0d";
let cmd = format!("printf 'Resume this session with:\\nclaude --resume {ID}\\n'");
let mut t = Task::spawn(20, &cmd, &cmd, &here(), 24, 80, 2000, &sh_env(), no_waker()).unwrap();
t.harness = Some(&crate::harness::Claude);
assert!(
wait_until(Duration::from_secs(60), || {
t.poll_exit().unwrap();
t.output_complete()
}),
"child never exited"
);
let (release, parked) = channel::<()>();
t.handle = Some(thread::spawn(move || {
let _ = parked.recv();
}));
t.scrape_exit_hint();
assert_eq!(t.scraped_id, None, "the scrape must wait for reader EOF");
drop(release);
assert!(
wait_until(Duration::from_secs(60), || t.reader_done()),
"the stand-in reader never stopped"
);
t.scrape_exit_hint();
assert_eq!(t.scraped_id.as_deref(), Some(ID));
}
#[test]
fn scrape_exit_hint_lands_an_open_sync_frame() {
const ID: &str = "7f3b9c1e-5a2d-4e8f-9b6a-0c4d2e8f1a3b";
let cmd = format!("printf '\\033[?2026hResume this session with:\\nclaude --resume {ID}\\n'");
let mut t = Task::spawn(21, &cmd, &cmd, &here(), 24, 80, 2000, &sh_env(), no_waker()).unwrap();
t.harness = Some(&crate::harness::Claude);
assert!(
wait_until(Duration::from_secs(60), || {
t.poll_exit().unwrap();
t.finished.is_some() && t.reader_done()
}),
"child never exited"
);
assert!(
!grid(&t.parser).text_with_history().contains(ID),
"premise: the unclosed frame still buffers the hint at scrape time"
);
t.scrape_exit_hint();
assert_eq!(t.scraped_id.as_deref(), Some(ID));
}
#[test]
fn finalize_preview_freezes_the_final_primary_line() {
use crate::protocol::PreviewSource;
let dir = temp("task_final_primary");
let flag = dir.join("flag");
let cmd = format!(
"until [ -e '{}' ]; do sleep 0.05; done; printf 'test result: ok\\n'",
flag.display()
);
let mut t = Task::spawn(40, &cmd, &cmd, &here(), 24, 80, 2000, &sh_env(), no_waker()).unwrap();
let early = t.resolve_preview(Instant::now());
assert!(!early.frozen);
std::fs::write(&flag, b"").unwrap();
assert!(
wait_until(Duration::from_secs(60), || {
t.poll_exit().unwrap();
t.output_complete()
}),
"child never completed"
);
t.finalize_preview();
let p = t.resolve_preview(Instant::now());
assert_eq!(
(p.text.as_str(), p.source, p.frozen),
("test result: ok", PreviewSource::Floor, true)
);
let _ = std::fs::remove_dir_all(&dir);
}
#[test]
fn finalize_preview_keeps_the_last_render_across_alt_teardown() {
use crate::protocol::PreviewSource;
let dir = temp("task_final_alt");
let teardown = dir.join("teardown");
let exit = dir.join("exit");
let cmd = format!(
"printf 'prelaunch junk\\n'; \
printf '\\033[?1049h\\033]0;working\\007app body'; \
until [ -e '{td}' ]; do sleep 0.05; done; printf '\\033[?1049l'; \
until [ -e '{ex}' ]; do sleep 0.05; done",
td = teardown.display(),
ex = exit.display()
);
let mut t = Task::spawn(41, &cmd, &cmd, &here(), 24, 80, 2000, &sh_env(), no_waker()).unwrap();
assert!(
wait_until(Duration::from_secs(5), || {
t.resolve_preview(Instant::now()).source == PreviewSource::Title
}),
"title never rendered"
);
std::fs::write(&teardown, b"").unwrap();
assert!(
wait_until(Duration::from_secs(60), || {
!grid(&t.parser).alternate_screen()
}),
"teardown never reached the grid"
);
assert_eq!(
t.resolve_preview(Instant::now()).source,
PreviewSource::Title,
"premise: the demotion hold keeps the title rendered"
);
std::fs::write(&exit, b"").unwrap();
assert!(
wait_until(Duration::from_secs(60), || {
t.poll_exit().unwrap();
t.output_complete()
}),
"child never completed"
);
t.finalize_preview();
assert_eq!(
grid(&t.parser).live_floor(),
"prelaunch junk",
"premise: 1049l restored the pre-launch primary screen"
);
let p = t.resolve_preview(Instant::now());
assert_eq!(
(p.text.as_str(), p.source, p.frozen),
("working", PreviewSource::Title, true)
);
let _ = std::fs::remove_dir_all(&dir);
}
#[test]
fn finalize_preview_freezes_primary_output_after_alt_teardown() {
use crate::protocol::PreviewSource;
let dir = temp("task_final_alt_output");
let flag = dir.join("flag");
let cmd = format!(
"printf 'prelaunch junk\\n'; \
printf '\\033[?1049h\\033]0;working\\007app body'; \
until [ -e '{}' ]; do sleep 0.05; done; \
printf '\\033[?1049ldone\\n'",
flag.display()
);
let mut t = Task::spawn(43, &cmd, &cmd, &here(), 24, 80, 2000, &sh_env(), no_waker()).unwrap();
assert!(
wait_until(Duration::from_secs(5), || {
t.resolve_preview(Instant::now()).source == PreviewSource::Title
}),
"title never rendered"
);
std::fs::write(&flag, b"").unwrap();
assert!(
wait_until(Duration::from_secs(60), || {
t.poll_exit().unwrap();
t.output_complete()
}),
"child never completed"
);
t.finalize_preview();
let p = t.resolve_preview(Instant::now());
assert_eq!(
(p.text.as_str(), p.source, p.frozen),
("done", PreviewSource::Floor, true),
"the post-teardown line must win over the stale title"
);
let _ = std::fs::remove_dir_all(&dir);
}
#[test]
fn summary_adapter_anchors_live_and_freezes_completion_at_exit() {
use crate::protocol::PreviewSource;
let dir = temp("task_anchor_e2e");
let flag = dir.join("flag");
let cmd = format!(
"printf '• Working (3s • esc to interrupt)\\n\\n› \\n synth-model high · 1 in · 2 out'; \
until [ -e '{f}' ]; do sleep 0.05; done; \
printf '\\033[H\\033[2J• Ran echo ok\\n\\n› \\n synth-model high · 2 in · 3 out'",
f = flag.display()
);
let mut t = Task::spawn(42, &cmd, &cmd, &here(), 24, 80, 2000, &sh_env(), no_waker()).unwrap();
assert!(t.summary_adapter.is_none(), "printf selects nothing");
t.summary_adapter = crate::harness::summary::select("codex");
assert!(t.summary_adapter.is_some());
let mut live = t.resolve_preview(Instant::now());
assert!(
wait_until(Duration::from_secs(5), || {
live = t.resolve_preview(Instant::now());
live.source == PreviewSource::Anchor
}),
"anchor never resolved, last preview {live:?}"
);
assert_eq!(
(live.text.as_str(), live.rule, live.frozen),
("synth-model high · Working", Some("codex:working"), false)
);
std::fs::write(&flag, b"").unwrap();
assert!(
wait_until(Duration::from_secs(60), || {
t.poll_exit().unwrap();
t.output_complete()
}),
"child never completed"
);
t.finalize_preview();
let p = t.resolve_preview(Instant::now());
assert_eq!(
(p.text.as_str(), p.source, p.rule, p.frozen),
(
"synth-model high · Ran echo ok",
PreviewSource::Anchor,
Some("codex:ran"),
true
)
);
let _ = std::fs::remove_dir_all(&dir);
}
#[test]
fn probe_replies_reach_the_child_through_the_allowlist() {
let dir = temp("task_probe");
let out = dir.join("out");
let cmd = format!(
"stty -icanon -echo min 1 time 0; printf '\\033[>c\\033[c\\033[6n'; \
head -c 11 > {}",
out.display()
);
let mut t = Task::spawn(11, &cmd, &cmd, &here(), 24, 80, 2000, &sh_env(), no_waker()).unwrap();
let mut got = Vec::new();
wait_until(Duration::from_secs(5), || {
got = std::fs::read(&out).unwrap_or_default();
got.len() >= 11
});
assert!(
got.starts_with(b"\x1b[?6c\x1b["),
"child must read the primary DA reply first (no secondary-DA \
leak); got {got:?}"
);
assert!(
got.ends_with(b"R"),
"CPR reply must follow the DA reply; got {got:?}"
);
t.terminate();
let _ = std::fs::remove_dir_all(&dir);
}