#![cfg(unix)]
use std::io::{BufRead, BufReader, Write};
use std::os::unix::net::UnixStream;
use std::path::PathBuf;
use std::process::{Child, Command, Stdio};
use std::sync::mpsc;
use std::time::{Duration, Instant};
use autofork_core::protocol::{
encode, Event, EventKind, Request, RequestBody, Response, ResponseBody,
};
use autofork_core::PROTO_VERSION;
struct Harness {
_tmp: tempfile::TempDir,
home: PathBuf,
socket: PathBuf,
project: PathBuf,
daemon: Option<Child>,
poll_grace_ms: Option<u64>,
wake_grace_secs: Option<u64>,
gate_grace_secs: Option<u64>,
chain_grace_secs: Option<u64>,
liveness_sweep_secs: Option<u64>,
session_sweep_secs: Option<u64>,
final_runner_bin: Option<PathBuf>,
daemon_env: Vec<(String, String)>,
}
impl Harness {
fn new(idle_deadline: &str, wake_debounce: &str) -> Self {
let tmp = tempfile::tempdir().unwrap();
let base = tmp.path().to_path_buf();
let home = base.join("fsan");
let project = base.join("proj");
std::fs::create_dir_all(&home).unwrap();
std::fs::create_dir_all(project.join(".autofork/forks")).unwrap();
std::fs::write(
home.join("config.toml"),
format!(
"default_idle_deadline = \"{idle_deadline}\"\nquiet_period = \"1h\"\nwake_debounce = \"{wake_debounce}\"\n",
),
)
.unwrap();
Self {
socket: base.join("d.sock"),
_tmp: tmp,
home,
project,
daemon: None,
poll_grace_ms: None,
wake_grace_secs: None,
gate_grace_secs: None,
chain_grace_secs: None,
liveness_sweep_secs: None,
session_sweep_secs: None,
final_runner_bin: None,
daemon_env: Vec::new(),
}
}
fn daemon_env(mut self, key: &str, value: &str) -> Self {
self.daemon_env.push((key.to_string(), value.to_string()));
self
}
fn poll_grace_ms(mut self, ms: u64) -> Self {
self.poll_grace_ms = Some(ms);
self
}
fn wake_grace_secs(mut self, secs: u64) -> Self {
self.wake_grace_secs = Some(secs);
self
}
fn gate_grace_secs(mut self, secs: u64) -> Self {
self.gate_grace_secs = Some(secs);
self
}
fn chain_grace_secs(mut self, secs: u64) -> Self {
self.chain_grace_secs = Some(secs);
self
}
fn liveness_sweep_secs(mut self, secs: u64) -> Self {
self.liveness_sweep_secs = Some(secs);
self
}
fn session_sweep_secs(mut self, secs: u64) -> Self {
self.session_sweep_secs = Some(secs);
self
}
fn recording_final_runner(&mut self) -> PathBuf {
let record = self.project.join("final-run.argv");
let script = self.project.join("fake-final-run.sh");
std::fs::write(
&script,
format!(
"#!/bin/sh\nprintf '%s\\n' \"$*\" >> \"{}\"\n",
record.display()
),
)
.unwrap();
#[cfg(unix)]
{
use std::os::unix::fs::PermissionsExt;
std::fs::set_permissions(&script, std::fs::Permissions::from_mode(0o755)).unwrap();
}
self.final_runner_bin = Some(script);
record
}
fn env_recording_final_runner(&mut self) -> PathBuf {
let record = self.project.join("final-run.env");
let script = self.project.join("fake-final-run-env.sh");
std::fs::write(
&script,
format!(
"#!/bin/sh\nprintf '%s|%s\\n' \"${{CLAUDE_CODE_OAUTH_TOKEN:--}}\" \"${{ANTHROPIC_API_KEY:--}}\" >> \"{}\"\n",
record.display()
),
)
.unwrap();
#[cfg(unix)]
{
use std::os::unix::fs::PermissionsExt;
std::fs::set_permissions(&script, std::fs::Permissions::from_mode(0o755)).unwrap();
}
self.final_runner_bin = Some(script);
record
}
fn append_config(&self, line: &str) {
use std::io::Write as _;
let mut f = std::fs::OpenOptions::new()
.append(true)
.open(self.home.join("config.toml"))
.unwrap();
writeln!(f, "{line}").unwrap();
}
fn write_fork(&self, rel: &str, content: &str) {
let path = self.project.join(".autofork/forks").join(rel);
std::fs::create_dir_all(path.parent().unwrap()).unwrap();
std::fs::write(path, content).unwrap();
}
fn write_logging_hook(&self, rel: &str, on: &str) -> PathBuf {
let log = self.project.join(format!("{}.log", rel.replace('/', "_")));
let path = self.project.join(".autofork/hooks").join(rel);
std::fs::create_dir_all(path.parent().unwrap()).unwrap();
std::fs::write(
path,
format!(
"---\nhook: true\non: {on}\n\
command: printf '%s\\n' \"$AUTOFORK_EVENT|${{AUTOFORK_SOURCE:-}}|${{AUTOFORK_END_REASON:-}}|${{AUTOFORK_IDLE_SECS:-}}|$AUTOFORK_SESSION_ID\" >> \"{}\"\n\
---\nlease-keeper documentation\n",
log.display()
),
)
.unwrap();
log
}
fn wait_for_hook_lines(&self, log: &PathBuf, n: usize, timeout: Duration) -> Vec<String> {
let start = Instant::now();
loop {
let lines: Vec<String> = std::fs::read_to_string(log)
.unwrap_or_default()
.lines()
.map(|l| l.to_string())
.collect();
if lines.len() >= n {
return lines;
}
assert!(
start.elapsed() < timeout,
"expected {n} hook lines, have {lines:?}"
);
std::thread::sleep(Duration::from_millis(50));
}
}
fn write_feed(&self, rel: &str, on: &str, deliver: &str, body: &str) {
let path = self.project.join(".autofork/hooks").join(rel);
std::fs::create_dir_all(path.parent().unwrap()).unwrap();
std::fs::write(
path,
format!(
"---\nhook: true\non: {on}\ndeliver: {deliver}\n\
command: |-\n printf '%s' \"{body}\"\n---\nfeed documentation\n"
),
)
.unwrap();
}
fn wait_for_reports(&self, session: &str, timeout: Duration) -> Vec<String> {
let start = Instant::now();
loop {
if let ResponseBody::Reports { blocks } = self.request(RequestBody::TakeReports {
session_id: session.to_string(),
wait_ms: None,
}) {
if !blocks.is_empty() {
return blocks;
}
}
assert!(start.elapsed() < timeout, "no feed block was spooled");
std::thread::sleep(Duration::from_millis(50));
}
}
fn watched_dir(&self) -> PathBuf {
let dir = self.project.join("watched");
std::fs::create_dir_all(&dir).unwrap();
dir
}
fn write_transcript(&self, tokens: u64) -> PathBuf {
let path = self.project.join("transcript.jsonl");
std::fs::write(
&path,
format!(
"{{\"type\":\"assistant\",\"message\":{{\"model\":\"m\",\"usage\":{{\"input_tokens\":{tokens},\"cache_read_input_tokens\":0,\"cache_creation_input_tokens\":0}}}}}}\n"
),
)
.unwrap();
path
}
fn append_transcript(&self, tokens: u64) {
use std::io::Write as _;
let mut f = std::fs::OpenOptions::new()
.append(true)
.open(self.project.join("transcript.jsonl"))
.unwrap();
writeln!(
f,
"{{\"type\":\"assistant\",\"message\":{{\"model\":\"m\",\"usage\":{{\"input_tokens\":{tokens},\"cache_read_input_tokens\":0,\"cache_creation_input_tokens\":0}}}}}}"
)
.unwrap();
}
fn start_daemon(&mut self) {
let mut cmd = Command::new(env!("CARGO_BIN_EXE_autofork-daemon"));
cmd.env("AUTOFORK_HOME", &self.home)
.env("AUTOFORK_SOCKET", &self.socket)
.env("AUTOFORK_CLAUDE_DIR", self.home.join("claude"))
.env("AUTOFORK_AGENTS_DIR", self.home.join("agents"))
.env("RUST_LOG", "debug")
.stdout(Stdio::null())
.stderr(Stdio::null());
if let Some(ms) = self.poll_grace_ms {
cmd.env("AUTOFORK_POLL_LOSS_GRACE_MS", ms.to_string());
}
if let Some(secs) = self.wake_grace_secs {
cmd.env("AUTOFORK_WAKE_GRACE_SECS", secs.to_string());
}
if let Some(secs) = self.gate_grace_secs {
cmd.env("AUTOFORK_GATE_GRACE_SECS", secs.to_string());
}
if let Some(secs) = self.chain_grace_secs {
cmd.env("AUTOFORK_CHAIN_GRACE_SECS", secs.to_string());
}
if let Some(secs) = self.liveness_sweep_secs {
cmd.env("AUTOFORK_LIVENESS_SWEEP_SECS", secs.to_string());
}
if let Some(secs) = self.session_sweep_secs {
cmd.env("AUTOFORK_SESSION_SWEEP_SECS", secs.to_string());
}
if let Some(bin) = &self.final_runner_bin {
cmd.env("AUTOFORK_FINAL_RUNNER_BIN", bin);
}
for (k, v) in &self.daemon_env {
cmd.env(k, v);
}
let child = cmd.spawn().unwrap();
self.daemon = Some(child);
let start = Instant::now();
loop {
if UnixStream::connect(&self.socket).is_ok() {
return;
}
assert!(
start.elapsed() < Duration::from_secs(10),
"daemon never came up"
);
std::thread::sleep(Duration::from_millis(25));
}
}
fn kill_daemon(&mut self) {
if let Some(mut child) = self.daemon.take() {
let _ = child.kill();
let _ = child.wait();
}
}
fn event(&self, kind: EventKind, session: &str) -> Event {
Event {
event: kind,
session_id: session.to_string(),
transcript_path: None,
cwd: self.project.clone(),
project_root: self.project.clone(),
source: None,
reason: None,
model: None,
enable_tags: None,
disable_tags: None,
waking: None,
notif_tool_use_id: None,
notif_task_id: None,
notif_status: None,
notif_continue: None,
context_tokens: None,
context_window: None,
client: None,
busy: None,
harness: None,
env: None,
}
}
fn prompt_submit(&self, session: &str, waking: bool) -> Event {
let mut ev = self.event(EventKind::PromptSubmit, session);
ev.waking = Some(waking);
ev
}
fn prompt_submit_notif(&self, session: &str, tool_use_id: &str, status: &str) -> Event {
let mut ev = self.event_t(EventKind::PromptSubmit, session);
ev.waking = Some(false);
ev.notif_tool_use_id = Some(tool_use_id.to_string());
ev.notif_status = Some(status.to_string());
ev.notif_continue = Some(false);
ev
}
fn prompt_submit_notif_cont(&self, session: &str, tool_use_id: &str) -> Event {
let mut ev = self.prompt_submit_notif(session, tool_use_id, "completed");
ev.notif_continue = Some(true);
ev
}
fn event_t(&self, kind: EventKind, session: &str) -> Event {
let mut ev = self.event(kind, session);
ev.transcript_path = Some(self.project.join("transcript.jsonl"));
ev
}
fn append_transcript_line(&self, line: &str) {
use std::io::Write as _;
let mut f = std::fs::OpenOptions::new()
.create(true)
.append(true)
.open(self.project.join("transcript.jsonl"))
.unwrap();
writeln!(f, "{line}").unwrap();
}
fn append_fork_spawn(&self, tool_use_id: &str, fork: &str) {
let prompt = format!(
"Read the file /x/{fork}.md and follow the instructions in its body. \
Context for this run: fork '{fork}', trigger 'idle', parent session s, \
conversation c, project root /p."
);
let line = serde_json::json!({
"type": "assistant",
"message": { "content": [
{ "type": "tool_use", "id": tool_use_id, "name": "Agent",
"input": { "subagent_type": "fork", "prompt": prompt } },
] }
});
self.append_transcript_line(&line.to_string());
}
fn append_background_launch(&self, tool_use_id: &str, task_id: &str) {
let line = serde_json::json!({
"type": "user",
"message": { "content": [
{ "type": "tool_result", "tool_use_id": tool_use_id, "content": [
{ "type": "text", "text": format!(
"Command running in background with ID: {task_id}. Output is being \
written to: /tmp/{task_id}.output. You will be notified when it \
completes.") },
] },
] }
});
self.append_transcript_line(&line.to_string());
}
fn append_monitor_launch(&self, tool_use_id: &str, task_id: &str) {
let line = serde_json::json!({
"type": "user",
"message": { "content": [
{ "type": "tool_result", "tool_use_id": tool_use_id, "content": [
{ "type": "text", "text": format!(
"Monitor started (task {task_id}, persistent — runs until TaskStop \
or session end). You will be notified on each event.") },
] },
] }
});
self.append_transcript_line(&line.to_string());
}
fn append_task_stop(&self, task_id: &str) {
let line = serde_json::json!({
"type": "assistant",
"message": { "content": [
{ "type": "tool_use", "id": format!("toolu_stop_{task_id}"), "name": "TaskStop",
"input": { "task_id": task_id } },
] }
});
self.append_transcript_line(&line.to_string());
}
fn append_completion_notification(&self, tool_use_id: &str, status: &str) {
self.append_completion_notification_result(tool_use_id, status, "report");
}
fn append_completion_notification_result(&self, tool_use_id: &str, status: &str, result: &str) {
let content = format!(
"<task-notification>\n<task-id>t-{tool_use_id}</task-id>\n\
<tool-use-id>{tool_use_id}</tool-use-id>\n<status>{status}</status>\n\
<summary>Agent \"x\" finished</summary>\n<result>{result}</result>"
);
let line = serde_json::json!({
"type": "user",
"message": { "content": content }
});
self.append_transcript_line(&line.to_string());
}
fn request(&self, body: RequestBody) -> ResponseBody {
let stream = UnixStream::connect(&self.socket).unwrap();
stream
.set_read_timeout(Some(Duration::from_secs(30)))
.unwrap();
let mut writer = stream.try_clone().unwrap();
let req = Request {
proto: PROTO_VERSION,
id: 1,
body,
};
writer.write_all(encode(&req).unwrap().as_bytes()).unwrap();
let mut line = String::new();
BufReader::new(stream).read_line(&mut line).unwrap();
serde_json::from_str::<Response>(line.trim()).unwrap().body
}
fn send_event(&self, ev: Event) -> ResponseBody {
self.request(RequestBody::Event(ev))
}
fn park_stop_wait(&self, ev: Event) -> mpsc::Receiver<ResponseBody> {
let socket = self.socket.clone();
let (tx, rx) = mpsc::channel();
std::thread::spawn(move || {
let stream = UnixStream::connect(&socket).unwrap();
stream
.set_read_timeout(Some(Duration::from_secs(60)))
.unwrap();
let mut writer = stream.try_clone().unwrap();
let req = Request {
proto: PROTO_VERSION,
id: 1,
body: RequestBody::StopWait(ev),
};
writer.write_all(encode(&req).unwrap().as_bytes()).unwrap();
let mut line = String::new();
let body = match BufReader::new(stream).read_line(&mut line) {
Ok(n) if n > 0 => serde_json::from_str::<Response>(line.trim()).unwrap().body,
_ => ResponseBody::Waited,
};
let _ = tx.send(body);
});
rx
}
fn status_recent_runs(&self) -> usize {
match self.request(RequestBody::Status) {
ResponseBody::StatusInfo(info) => info.recent_runs.len(),
other => panic!("unexpected: {other:?}"),
}
}
fn open_sessions(&self) -> Vec<autofork_core::protocol::SessionInfo> {
match self.request(RequestBody::Status) {
ResponseBody::StatusInfo(info) => info.sessions,
other => panic!("unexpected: {other:?}"),
}
}
fn has_open_session(&self, session: &str) -> bool {
self.open_sessions().iter().any(|s| s.session_id == session)
}
fn drop_stop_wait(&self, ev: Event) {
let stream = UnixStream::connect(&self.socket).unwrap();
let mut writer = stream.try_clone().unwrap();
let req = Request {
proto: PROTO_VERSION,
id: 1,
body: RequestBody::StopWait(ev),
};
writer.write_all(encode(&req).unwrap().as_bytes()).unwrap();
std::thread::sleep(Duration::from_millis(150));
drop(writer);
drop(stream);
}
}
impl Drop for Harness {
fn drop(&mut self) {
self.kill_daemon();
}
}
fn assert_ack(body: ResponseBody) {
assert!(
matches!(body, ResponseBody::Ack),
"expected ack, got {body:?}"
);
}
fn wake_payload(body: ResponseBody) -> String {
match body {
ResponseBody::Wake { payload, .. } => payload,
other => panic!("expected a Wake, got {other:?}"),
}
}
fn wake_forks(body: ResponseBody) -> Vec<autofork_core::protocol::WakeFork> {
match body {
ResponseBody::Wake { forks, .. } => forks.expect("wake carries structured forks"),
other => panic!("expected a Wake, got {other:?}"),
}
}
#[test]
fn idle_wake_names_the_fork() {
let mut h = Harness::new("1s", "0");
h.write_fork(
"journal.md",
"---\nfork: true\nrun_on: [idle]\n---\nwrite the journal now",
);
h.start_daemon();
assert_ack(h.send_event(h.event(EventKind::SessionStart, "s1")));
let rx = h.park_stop_wait(h.event(EventKind::Stop, "s1"));
let payload = wake_payload(rx.recv_timeout(Duration::from_secs(10)).unwrap());
assert!(payload.contains("source: autofork"));
assert!(payload.contains("due: journal (trigger: idle)"));
assert!(payload.contains("subagent_type \"fork\""));
assert!(payload.contains("journal.md"));
assert!(payload.contains("parent session s1"));
assert!(payload.contains(&format!("project root {}", h.project.display())));
assert!(payload.contains("Do not read that file yourself"));
assert!(payload.contains("skip spawning it"));
assert_eq!(h.status_recent_runs(), 1);
}
#[test]
fn wake_debounce_zero_is_immediate() {
let mut h = Harness::new("1s", "0");
h.write_fork("j.md", "---\nfork: true\nrun_on: [idle]\n---\nbody");
h.start_daemon();
assert_ack(h.send_event(h.event(EventKind::SessionStart, "s1")));
let rx = h.park_stop_wait(h.event(EventKind::Stop, "s1"));
let start = Instant::now();
let payload = wake_payload(rx.recv_timeout(Duration::from_secs(5)).unwrap());
assert!(payload.contains("due: j"));
assert!(
start.elapsed() < Duration::from_secs(4),
"too slow: {:?}",
start.elapsed()
);
}
#[test]
fn prompt_submit_cancels_parked_wait_without_stamping() {
let mut h = Harness::new("1s", "3");
h.write_fork("j.md", "---\nfork: true\nrun_on: [idle]\n---\nbody");
h.start_daemon();
assert_ack(h.send_event(h.event(EventKind::SessionStart, "s1")));
let rx = h.park_stop_wait(h.event(EventKind::Stop, "s1"));
std::thread::sleep(Duration::from_millis(1500));
assert_ack(h.send_event(h.event(EventKind::PromptSubmit, "s1")));
let body = rx.recv_timeout(Duration::from_secs(5)).unwrap();
assert!(
matches!(body, ResponseBody::Waited),
"expected Waited, got {body:?}"
);
assert_eq!(h.status_recent_runs(), 0, "throttle stamped despite cancel");
}
#[test]
fn shutdown_resolves_parked_wait() {
let mut h = Harness::new("1h", "0");
h.write_fork("j.md", "---\nfork: true\nrun_on: [idle]\n---\nbody");
h.start_daemon();
assert_ack(h.send_event(h.event(EventKind::SessionStart, "s1")));
let rx = h.park_stop_wait(h.event(EventKind::Stop, "s1"));
std::thread::sleep(Duration::from_millis(300));
assert_ack(h.request(RequestBody::Shutdown { drain: false }));
let body = rx.recv_timeout(Duration::from_secs(5)).unwrap();
assert!(
matches!(body, ResponseBody::Waited),
"expected Waited, got {body:?}"
);
}
#[test]
fn disable_tag_filters_fork_but_untagged_wakes() {
let mut h = Harness::new("1s", "0");
h.write_fork(
"tagged.md",
"---\nfork: true\nrun_on: [idle]\ntags: [ci]\n---\nTAGGED",
);
h.write_fork("plain.md", "---\nfork: true\nrun_on: [idle]\n---\nPLAIN");
h.start_daemon();
let mut start_ev = h.event(EventKind::SessionStart, "s1");
start_ev.disable_tags = Some(vec!["ci".into()]);
assert_ack(h.send_event(start_ev));
let mut stop = h.event(EventKind::Stop, "s1");
stop.disable_tags = Some(vec!["ci".into()]);
let rx = h.park_stop_wait(stop);
let payload = wake_payload(rx.recv_timeout(Duration::from_secs(10)).unwrap());
assert!(payload.contains("due: plain"));
assert!(
!payload.contains("due: tagged"),
"disabled fork leaked: {payload}"
);
}
#[test]
fn enable_list_excludes_untagged_fork() {
let mut h = Harness::new("1s", "0");
h.write_fork(
"tagged.md",
"---\nfork: true\nrun_on: [idle]\ntags: [ci]\n---\nTAGGED",
);
h.write_fork("plain.md", "---\nfork: true\nrun_on: [idle]\n---\nPLAIN");
h.start_daemon();
let mut start_ev = h.event(EventKind::SessionStart, "s1");
start_ev.enable_tags = Some(vec!["ci".into()]);
assert_ack(h.send_event(start_ev));
let mut stop = h.event(EventKind::Stop, "s1");
stop.enable_tags = Some(vec!["ci".into()]);
let rx = h.park_stop_wait(stop);
let payload = wake_payload(rx.recv_timeout(Duration::from_secs(10)).unwrap());
assert!(payload.contains("due: tagged"));
assert!(
!payload.contains("due: plain"),
"untagged fork ran despite whitelist: {payload}"
);
}
#[test]
fn throttle_suppresses_second_wake() {
let mut h = Harness::new("1s", "0");
h.write_fork(
"j.md",
"---\nfork: true\nrun_on: [idle]\nthrottle: 1h\n---\nbody",
);
h.start_daemon();
assert_ack(h.send_event(h.event(EventKind::SessionStart, "s1")));
let rx = h.park_stop_wait(h.event(EventKind::Stop, "s1"));
let _ = wake_payload(rx.recv_timeout(Duration::from_secs(10)).unwrap());
let rx2 = h.park_stop_wait(h.event(EventKind::Stop, "s1"));
std::thread::sleep(Duration::from_millis(1500));
assert_ack(h.send_event(h.event(EventKind::PromptSubmit, "s1")));
let body = rx2.recv_timeout(Duration::from_secs(5)).unwrap();
assert!(
matches!(body, ResponseBody::Waited),
"throttled fork woke again: {body:?}"
);
}
#[test]
fn tag_throttle_suppresses_group_but_other_tag_wakes() {
let mut h = Harness::new("1s", "0");
std::fs::write(
h.home.join("config.toml"),
"default_idle_deadline = \"1s\"\nquiet_period = \"1h\"\nwake_debounce = \"0\"\n[tag_throttles]\nci = \"1h\"\n",
)
.unwrap();
h.write_fork(
"a.md",
"---\nfork: true\nrun_on: [idle]\ntags: [ci]\n---\nA",
);
h.write_fork(
"b.md",
"---\nfork: true\nrun_on: [idle]\ntags: [ci]\n---\nB",
);
h.write_fork(
"c.md",
"---\nfork: true\nrun_on: [idle]\ntags: [docs]\n---\nC",
);
h.start_daemon();
assert_ack(h.send_event(h.event(EventKind::SessionStart, "s1")));
let rx = h.park_stop_wait(h.event(EventKind::Stop, "s1"));
let payload = wake_payload(rx.recv_timeout(Duration::from_secs(10)).unwrap());
assert!(payload.contains("due: a") && payload.contains("due: b") && payload.contains("due: c"));
assert_ack(h.send_event(h.prompt_submit("s1", true)));
let rx2 = h.park_stop_wait(h.event(EventKind::Stop, "s1"));
let payload = wake_payload(rx2.recv_timeout(Duration::from_secs(10)).unwrap());
assert!(
payload.contains("due: c"),
"docs fork should still wake: {payload}"
);
assert!(
!payload.contains("due: a"),
"ci fork a not throttled: {payload}"
);
assert!(
!payload.contains("due: b"),
"ci fork b not throttled: {payload}"
);
}
#[test]
fn after_dependent_held_until_predecessor_completes() {
let mut h = Harness::new("1s", "0");
h.write_fork("alpha.md", "---\nfork: true\nrun_on: [idle]\n---\nALPHA");
h.write_fork(
"beta.md",
"---\nfork: true\nrun_on: [idle]\nafter: alpha\n---\nBETA",
);
h.start_daemon();
h.write_transcript(100);
assert_ack(h.send_event(h.event_t(EventKind::SessionStart, "s1")));
let rx = h.park_stop_wait(h.event_t(EventKind::Stop, "s1"));
let payload = wake_payload(rx.recv_timeout(Duration::from_secs(10)).unwrap());
assert!(payload.contains("due: alpha"), "{payload}");
assert!(payload.contains("held back by autofork"), "{payload}");
assert!(payload.contains("'beta' (after 'alpha')"), "{payload}");
assert!(!payload.contains("due: beta"), "{payload}");
h.append_fork_spawn("toolu_alpha", "alpha");
let rx2 = h.park_stop_wait(h.event_t(EventKind::Stop, "s1"));
std::thread::sleep(Duration::from_millis(400));
h.append_completion_notification("toolu_alpha", "completed");
assert_ack(h.send_event(h.prompt_submit("s1", false)));
assert!(matches!(
rx2.recv_timeout(Duration::from_secs(5)).unwrap(),
ResponseBody::Waited
));
let rx3 = h.park_stop_wait(h.event_t(EventKind::Stop, "s1"));
let release = wake_payload(rx3.recv_timeout(Duration::from_secs(10)).unwrap());
assert!(
release.contains("due: beta (trigger: idle) — released, 'alpha' finished"),
"{release}"
);
assert!(release.contains("Read the file"), "{release}");
assert!(release.contains("beta.md"), "{release}");
assert!(
release.contains("append the report(s) 'alpha' returned"),
"{release}"
);
let rx4 = h.park_stop_wait(h.event_t(EventKind::Stop, "s1"));
assert!(
rx4.recv_timeout(Duration::from_millis(2500)).is_err(),
"release fired twice"
);
}
#[test]
fn priority_layers_forks_into_waves() {
let mut h = Harness::new("1s", "0");
h.write_fork(
"aaa-last.md",
"---\nfork: true\nrun_on: [idle]\npriority: 10\n---\nLAST",
);
h.write_fork(
"zzz-first.md",
"---\nfork: true\nrun_on: [idle]\n---\nFIRST",
);
h.start_daemon();
h.write_transcript(100);
assert_ack(h.send_event(h.event_t(EventKind::SessionStart, "s1")));
let rx = h.park_stop_wait(h.event_t(EventKind::Stop, "s1"));
let body = rx.recv_timeout(Duration::from_secs(10)).unwrap();
let payload = wake_payload(body.clone());
assert!(payload.contains("due: zzz-first"), "{payload}");
assert!(!payload.contains("due: aaa-last"), "{payload}");
assert!(payload.contains("held back by autofork"), "{payload}");
assert!(
payload.contains("'aaa-last' (after 'zzz-first')"),
"{payload}"
);
let forks = wake_forks(body);
assert_eq!(forks.len(), 1);
assert_eq!(forks[0].name, "zzz-first");
h.append_fork_spawn("toolu_first", "zzz-first");
let rx2 = h.park_stop_wait(h.event_t(EventKind::Stop, "s1"));
std::thread::sleep(Duration::from_millis(400));
h.append_completion_notification("toolu_first", "completed");
assert_ack(h.send_event(h.prompt_submit("s1", false)));
assert!(matches!(
rx2.recv_timeout(Duration::from_secs(5)).unwrap(),
ResponseBody::Waited
));
let rx3 = h.park_stop_wait(h.event_t(EventKind::Stop, "s1"));
let body = rx3.recv_timeout(Duration::from_secs(10)).unwrap();
let release = wake_payload(body.clone());
assert!(
release.contains("due: aaa-last (trigger: idle) — released, earlier forks finished"),
"{release}"
);
assert!(
release.contains("The forks ordered before this one have finished."),
"{release}"
);
assert!(!release.contains("append the report(s)"), "{release}");
let forks = wake_forks(body);
assert_eq!(forks.len(), 1);
assert_eq!(forks[0].name, "aaa-last");
assert!(forks[0].after.is_empty());
}
#[test]
fn after_wins_over_priority_and_reports_still_pipe() {
let mut h = Harness::new("1s", "0");
h.write_fork("alpha.md", "---\nfork: true\nrun_on: [idle]\n---\nA");
h.write_fork(
"beta.md",
"---\nfork: true\nrun_on: [idle]\nafter: alpha\npriority: -5\n---\nB",
);
h.start_daemon();
h.write_transcript(100);
assert_ack(h.send_event(h.event_t(EventKind::SessionStart, "s1")));
let rx = h.park_stop_wait(h.event_t(EventKind::Stop, "s1"));
let payload = wake_payload(rx.recv_timeout(Duration::from_secs(10)).unwrap());
assert!(payload.contains("due: alpha"), "{payload}");
assert!(!payload.contains("due: beta"), "{payload}");
h.append_fork_spawn("toolu_alpha", "alpha");
let rx2 = h.park_stop_wait(h.event_t(EventKind::Stop, "s1"));
std::thread::sleep(Duration::from_millis(400));
h.append_completion_notification("toolu_alpha", "completed");
assert_ack(h.send_event(h.prompt_submit("s1", false)));
assert!(matches!(
rx2.recv_timeout(Duration::from_secs(5)).unwrap(),
ResponseBody::Waited
));
let rx3 = h.park_stop_wait(h.event_t(EventKind::Stop, "s1"));
let release = wake_payload(rx3.recv_timeout(Duration::from_secs(10)).unwrap());
assert!(
release.contains("append the report(s) 'alpha' returned"),
"{release}"
);
}
#[test]
fn skill_attached_fork_wake_tells_the_fork_to_load_the_skill() {
let mut h = Harness::new("1s", "0");
let skill_dir = h.project.join(".claude/skills/feedback");
std::fs::create_dir_all(&skill_dir).unwrap();
std::fs::write(
skill_dir.join("SKILL.md"),
"---\nname: feedback\ndescription: d\n---\nskill body",
)
.unwrap();
std::fs::write(
skill_dir.join("FORK.md"),
"---\nfork: true\nrun_on: [idle]\n---\napply the skill",
)
.unwrap();
h.start_daemon();
h.write_transcript(100);
assert_ack(h.send_event(h.event_t(EventKind::SessionStart, "s1")));
let rx = h.park_stop_wait(h.event_t(EventKind::Stop, "s1"));
let payload = wake_payload(rx.recv_timeout(Duration::from_secs(10)).unwrap());
assert!(payload.contains("due: feedback"), "{payload}");
assert!(payload.contains("belongs to the skill at"), "{payload}");
assert!(payload.contains("SKILL.md"), "{payload}");
assert!(payload.contains("not already in your context"), "{payload}");
}
#[test]
fn foreign_task_completion_starts_a_new_pause() {
let mut h = Harness::new("1s", "0").wake_grace_secs(0);
h.write_fork("journal.md", "---\nfork: true\nrun_on: [idle]\n---\nJ");
h.start_daemon();
h.write_transcript(100);
assert_ack(h.send_event(h.event_t(EventKind::SessionStart, "s1")));
let rx = h.park_stop_wait(h.event_t(EventKind::Stop, "s1"));
wake_payload(rx.recv_timeout(Duration::from_secs(10)).unwrap());
let rx2 = h.park_stop_wait(h.event_t(EventKind::Stop, "s1"));
std::thread::sleep(Duration::from_millis(400));
assert_ack(h.send_event(h.prompt_submit_notif("s1", "toolu_users_build", "completed")));
assert!(matches!(
rx2.recv_timeout(Duration::from_secs(5)).unwrap(),
ResponseBody::Waited
));
let rx3 = h.park_stop_wait(h.event_t(EventKind::Stop, "s1"));
let payload = wake_payload(rx3.recv_timeout(Duration::from_secs(10)).unwrap());
assert!(payload.contains("due: journal"), "{payload}");
}
#[test]
fn own_fork_completion_matches_even_without_an_intervening_stop() {
let mut h = Harness::new("1s", "0").wake_grace_secs(0);
h.write_fork("journal.md", "---\nfork: true\nrun_on: [idle]\n---\nJ");
h.start_daemon();
h.write_transcript(100);
assert_ack(h.send_event(h.event_t(EventKind::SessionStart, "s1")));
let rx = h.park_stop_wait(h.event_t(EventKind::Stop, "s1"));
wake_payload(rx.recv_timeout(Duration::from_secs(10)).unwrap());
h.append_fork_spawn("toolu_j", "journal");
assert_ack(h.send_event(h.prompt_submit_notif("s1", "toolu_j", "completed")));
let rx2 = h.park_stop_wait(h.event_t(EventKind::Stop, "s1"));
assert!(
rx2.recv_timeout(Duration::from_millis(2500)).is_err(),
"own fork completion re-fired the idle fork without an intervening Stop"
);
}
#[test]
fn own_fork_completion_does_not_restart_the_pause() {
let mut h = Harness::new("1s", "0").wake_grace_secs(0);
h.write_fork("journal.md", "---\nfork: true\nrun_on: [idle]\n---\nJ");
h.start_daemon();
h.write_transcript(100);
assert_ack(h.send_event(h.event_t(EventKind::SessionStart, "s1")));
let rx = h.park_stop_wait(h.event_t(EventKind::Stop, "s1"));
wake_payload(rx.recv_timeout(Duration::from_secs(10)).unwrap());
h.append_fork_spawn("toolu_j", "journal");
let rx2 = h.park_stop_wait(h.event_t(EventKind::Stop, "s1"));
std::thread::sleep(Duration::from_millis(400));
assert_ack(h.send_event(h.prompt_submit_notif("s1", "toolu_j", "completed")));
assert!(matches!(
rx2.recv_timeout(Duration::from_secs(5)).unwrap(),
ResponseBody::Waited
));
let rx3 = h.park_stop_wait(h.event_t(EventKind::Stop, "s1"));
assert!(
rx3.recv_timeout(Duration::from_millis(2500)).is_err(),
"own fork completion re-fired the idle fork"
);
}
#[test]
fn context_threshold_wakes_and_latches_once() {
let mut h = Harness::new("1h", "0"); h.write_fork(
"ctx.md",
"---\nfork: true\nrun_on:\n - context_tokens: 1000\n---\ncontext filling",
);
h.start_daemon();
let transcript = h.write_transcript(2000);
let mut start = h.event(EventKind::SessionStart, "s1");
start.transcript_path = Some(transcript.clone());
assert_ack(h.send_event(start));
let mut stop = h.event(EventKind::Stop, "s1");
stop.transcript_path = Some(transcript.clone());
let rx = h.park_stop_wait(stop);
let payload = wake_payload(rx.recv_timeout(Duration::from_secs(10)).unwrap());
assert!(
payload.contains("due: ctx (trigger: context_tokens:1000)"),
"{payload}"
);
let mut stop2 = h.event(EventKind::Stop, "s1");
stop2.transcript_path = Some(transcript);
let rx2 = h.park_stop_wait(stop2);
std::thread::sleep(Duration::from_millis(400));
assert_ack(h.send_event(h.event(EventKind::PromptSubmit, "s1")));
let body = rx2.recv_timeout(Duration::from_secs(5)).unwrap();
assert!(
matches!(body, ResponseBody::Waited),
"context re-fired: {body:?}"
);
}
#[test]
fn context_used_respects_1m_model_window() {
let mut h = Harness::new("1h", "0"); h.write_fork(
"ctx75.md",
"---\nfork: true\nrun_on:\n - context_used: 75%\n---\nnearly full",
);
h.start_daemon();
let transcript = h.write_transcript(300_000);
let mut start = h.event(EventKind::SessionStart, "s1");
start.transcript_path = Some(transcript.clone());
start.model = Some("claude-opus-4-8[1m]".to_string());
assert_ack(h.send_event(start));
let mut stop = h.event(EventKind::Stop, "s1");
stop.transcript_path = Some(transcript.clone());
let rx = h.park_stop_wait(stop);
std::thread::sleep(Duration::from_millis(400));
assert_ack(h.send_event(h.prompt_submit("s1", true)));
let body = rx.recv_timeout(Duration::from_secs(5)).unwrap();
assert!(
matches!(body, ResponseBody::Waited),
"context fired at 30% of a 1M window: {body:?}"
);
h.append_transcript(800_000);
let mut stop2 = h.event(EventKind::Stop, "s1");
stop2.transcript_path = Some(transcript);
let rx2 = h.park_stop_wait(stop2);
let payload = wake_payload(rx2.recv_timeout(Duration::from_secs(10)).unwrap());
assert!(
payload.contains("due: ctx75 (trigger: context_used:75%)"),
"{payload}"
);
}
#[test]
fn oversized_gauge_bumps_unmarked_window() {
let mut h = Harness::new("1h", "0");
h.write_fork(
"ctx75.md",
"---\nfork: true\nrun_on:\n - context_used: 75%\n---\nnearly full",
);
h.start_daemon();
let transcript = h.write_transcript(300_000);
let mut start = h.event(EventKind::SessionStart, "s1");
start.transcript_path = Some(transcript.clone());
assert_ack(h.send_event(start));
let mut stop = h.event(EventKind::Stop, "s1");
stop.transcript_path = Some(transcript);
let rx = h.park_stop_wait(stop);
std::thread::sleep(Duration::from_millis(400));
assert_ack(h.send_event(h.prompt_submit("s1", true)));
let body = rx.recv_timeout(Duration::from_secs(5)).unwrap();
assert!(
matches!(body, ResponseBody::Waited),
"context fired despite oversized-gauge bump: {body:?}"
);
}
#[test]
fn debounce_batches_forks_across_the_window() {
let mut h = Harness::new("1s", "2");
h.write_fork("a.md", "---\nfork: true\nrun_on:\n - idle: 1\n---\nA");
h.write_fork("b.md", "---\nfork: true\nrun_on:\n - idle: 2\n---\nB");
h.start_daemon();
assert_ack(h.send_event(h.event(EventKind::SessionStart, "s1")));
let rx = h.park_stop_wait(h.event(EventKind::Stop, "s1"));
let payload = wake_payload(rx.recv_timeout(Duration::from_secs(10)).unwrap());
assert!(payload.contains("due: a (trigger: idle:1)"), "{payload}");
assert!(payload.contains("due: b (trigger: idle:2)"), "{payload}");
assert_eq!(payload.matches("After spawning all forks above").count(), 1);
assert_eq!(h.status_recent_runs(), 2);
}
#[test]
fn idle_fork_fires_at_most_once_per_pause() {
let mut h = Harness::new("1s", "0");
h.write_fork("j.md", "---\nfork: true\nrun_on: [idle]\n---\nbody");
h.start_daemon();
assert_ack(h.send_event(h.event(EventKind::SessionStart, "s1")));
let rx = h.park_stop_wait(h.event(EventKind::Stop, "s1"));
let payload = wake_payload(rx.recv_timeout(Duration::from_secs(10)).unwrap());
assert!(payload.contains("due: j"));
assert_eq!(h.status_recent_runs(), 1);
assert_ack(h.send_event(h.prompt_submit("s1", false)));
let rx2 = h.park_stop_wait(h.event(EventKind::Stop, "s1"));
std::thread::sleep(Duration::from_millis(1500));
assert!(
rx2.try_recv().is_err(),
"fork re-fired within the same pause"
);
assert_ack(h.send_event(h.prompt_submit("s1", false)));
assert!(matches!(
rx2.recv_timeout(Duration::from_secs(5)).unwrap(),
ResponseBody::Waited
));
assert_eq!(h.status_recent_runs(), 1);
assert_ack(h.send_event(h.prompt_submit("s1", true)));
let rx3 = h.park_stop_wait(h.event(EventKind::Stop, "s1"));
let payload = wake_payload(rx3.recv_timeout(Duration::from_secs(10)).unwrap());
assert!(
payload.contains("due: j"),
"new pause did not re-arm the fork"
);
assert_eq!(h.status_recent_runs(), 2);
}
#[test]
fn ambiguous_prompt_within_grace_is_treated_as_continuation() {
let mut h = Harness::new("1s", "0");
h.write_fork("j.md", "---\nfork: true\nrun_on: [idle]\n---\nbody");
h.start_daemon();
assert_ack(h.send_event(h.event(EventKind::SessionStart, "s1")));
let rx = h.park_stop_wait(h.event(EventKind::Stop, "s1"));
let _ = wake_payload(rx.recv_timeout(Duration::from_secs(10)).unwrap());
assert_ack(h.send_event(h.event(EventKind::PromptSubmit, "s1")));
let rx2 = h.park_stop_wait(h.event(EventKind::Stop, "s1"));
std::thread::sleep(Duration::from_millis(1500));
assert!(
rx2.try_recv().is_err(),
"belt failed: ambiguous prompt advanced the pause"
);
assert_ack(h.send_event(h.prompt_submit("s1", false)));
assert!(matches!(
rx2.recv_timeout(Duration::from_secs(5)).unwrap(),
ResponseBody::Waited
));
}
#[test]
fn throttle_holds_across_pauses() {
let mut h = Harness::new("1s", "0");
h.write_fork(
"j.md",
"---\nfork: true\nrun_on: [idle]\nthrottle: 1h\n---\nbody",
);
h.start_daemon();
assert_ack(h.send_event(h.event(EventKind::SessionStart, "s1")));
let rx = h.park_stop_wait(h.event(EventKind::Stop, "s1"));
let _ = wake_payload(rx.recv_timeout(Duration::from_secs(10)).unwrap());
assert_ack(h.send_event(h.prompt_submit("s1", true)));
let rx2 = h.park_stop_wait(h.event(EventKind::Stop, "s1"));
std::thread::sleep(Duration::from_millis(1500));
assert!(
rx2.try_recv().is_err(),
"throttle didn't hold across pauses"
);
assert_ack(h.send_event(h.prompt_submit("s1", false)));
assert!(matches!(
rx2.recv_timeout(Duration::from_secs(5)).unwrap(),
ResponseBody::Waited
));
}
#[test]
fn lost_poll_closes_session_after_grace() {
let mut h = Harness::new("1h", "0").poll_grace_ms(400);
h.write_fork("j.md", "---\nfork: true\nrun_on: [idle]\n---\nbody");
h.start_daemon();
assert_ack(h.send_event(h.event(EventKind::SessionStart, "s1")));
assert!(h.has_open_session("s1"));
h.drop_stop_wait(h.event(EventKind::Stop, "s1"));
std::thread::sleep(Duration::from_millis(150));
assert!(h.has_open_session("s1"), "closed before the grace elapsed");
std::thread::sleep(Duration::from_millis(500));
assert!(
!h.has_open_session("s1"),
"lost poll did not close the session"
);
assert_ack(h.send_event(h.event(EventKind::SessionStart, "s1")));
assert!(
h.has_open_session("s1"),
"a later event did not re-open the session"
);
}
#[test]
fn event_within_grace_keeps_session_open() {
let mut h = Harness::new("1h", "0").poll_grace_ms(700);
h.write_fork("j.md", "---\nfork: true\nrun_on: [idle]\n---\nbody");
h.start_daemon();
assert_ack(h.send_event(h.event(EventKind::SessionStart, "s1")));
h.drop_stop_wait(h.event(EventKind::Stop, "s1"));
std::thread::sleep(Duration::from_millis(200));
assert_ack(h.send_event(h.event(EventKind::SessionStart, "s1")));
std::thread::sleep(Duration::from_millis(800));
assert!(
h.has_open_session("s1"),
"grace-close fired despite a fresh event"
);
}
#[test]
fn answered_poll_never_triggers_grace_close() {
let mut h = Harness::new("1s", "0").poll_grace_ms(400);
h.write_fork("j.md", "---\nfork: true\nrun_on: [idle]\n---\nbody");
h.start_daemon();
assert_ack(h.send_event(h.event(EventKind::SessionStart, "s1")));
let rx = h.park_stop_wait(h.event(EventKind::Stop, "s1"));
let _ = wake_payload(rx.recv_timeout(Duration::from_secs(10)).unwrap());
std::thread::sleep(Duration::from_millis(600)); assert!(
h.has_open_session("s1"),
"an answered poll wrongly closed the session"
);
}
#[test]
fn stale_annotation_for_idle_open_session_without_poll() {
let mut h = Harness::new("1s", "0"); h.start_daemon();
assert_ack(h.send_event(h.event(EventKind::SessionStart, "s1")));
std::thread::sleep(Duration::from_millis(3300));
let stale = h
.open_sessions()
.into_iter()
.find(|s| s.session_id == "s1")
.map(|s| s.stale)
.unwrap_or(false);
assert!(
stale,
"an old open session with no poll should be flagged stale"
);
}
#[test]
fn list_forks_marks_only_marked_files_and_status_and_shutdown() {
let mut h = Harness::new("1h", "0");
h.write_fork(
"info/FORK.md",
"---\nfork: true\ndescription: nested fork\nrun_on: [idle]\nthrottle: 5m\n---\nbody",
);
h.write_fork("oops.md", "---\nrun_on: [idle]\n---\nnope");
h.start_daemon();
assert_ack(h.send_event(h.event(EventKind::SessionStart, "s1")));
match h.request(RequestBody::Status) {
ResponseBody::StatusInfo(info) => {
assert_eq!(info.daemon_proto, PROTO_VERSION);
assert_eq!(info.sessions.len(), 1);
}
other => panic!("unexpected: {other:?}"),
}
match h.request(RequestBody::ListForks {
project_root: h.project.clone(),
cwd: h.project.clone(),
}) {
ResponseBody::ForkList { items } => {
assert_eq!(items.len(), 1);
assert_eq!(items[0].name, "info");
assert_eq!(items[0].throttle_secs, Some(300));
let has_warn = items
.iter()
.any(|f| f.warnings.iter().any(|w| w.contains("no `fork: true`")));
assert!(has_warn, "missing fork-like warning: {items:?}");
}
other => panic!("unexpected: {other:?}"),
}
assert_ack(h.request(RequestBody::Shutdown { drain: true }));
let start = Instant::now();
loop {
if UnixStream::connect(&h.socket).is_err() {
break;
}
assert!(
start.elapsed() < Duration::from_secs(10),
"daemon didn't exit"
);
std::thread::sleep(Duration::from_millis(50));
}
if let Some(mut child) = h.daemon.take() {
let _ = child.wait();
}
}
#[test]
fn prune_closes_stale_sessions_only() {
let mut h = Harness::new("1s", "0");
h.start_daemon();
assert_ack(h.send_event(h.event(EventKind::SessionStart, "s1")));
assert_ack(h.send_event(h.event(EventKind::SessionStart, "s3")));
let parked = h.park_stop_wait(h.event(EventKind::Stop, "s3"));
std::thread::sleep(Duration::from_millis(3200));
assert_ack(h.send_event(h.event(EventKind::SessionStart, "s2")));
let stale: Vec<String> = h
.open_sessions()
.into_iter()
.filter(|s| s.stale)
.map(|s| s.session_id)
.collect();
assert_eq!(stale, vec!["s1".to_string()], "status stale annotation");
match h.request(RequestBody::Prune) {
ResponseBody::Pruned { sessions } => {
assert_eq!(sessions.len(), 1, "pruned: {sessions:?}");
assert_eq!(sessions[0].session_id, "s1");
assert_eq!(sessions[0].status, "closed");
}
other => panic!("unexpected: {other:?}"),
}
assert!(!h.has_open_session("s1"), "stale session still open");
assert!(h.has_open_session("s2"), "active session was pruned");
assert!(h.has_open_session("s3"), "parked session was pruned");
match h.request(RequestBody::Prune) {
ResponseBody::Pruned { sessions } => assert!(sessions.is_empty(), "{sessions:?}"),
other => panic!("unexpected: {other:?}"),
}
assert_ack(h.send_event(h.event(EventKind::SessionStart, "s1")));
assert!(h.has_open_session("s1"), "event did not re-open");
assert_ack(h.send_event(h.prompt_submit("s3", true)));
let _ = parked.recv_timeout(Duration::from_secs(5));
}
fn oc_event(h: &Harness, kind: EventKind, session: &str) -> Event {
let mut ev = h.event(kind, session);
ev.client = Some("opencode".to_string());
ev
}
#[test]
fn opencode_wake_carries_structured_forks() {
let mut h = Harness::new("1s", "0");
h.write_fork(
"journal.md",
"---\nfork: true\nrun_on: [idle]\n---\nwrite the journal now",
);
h.start_daemon();
assert_ack(h.send_event(oc_event(&h, EventKind::SessionStart, "oc1")));
let rx = h.park_stop_wait(oc_event(&h, EventKind::Stop, "oc1"));
let forks = wake_forks(rx.recv_timeout(Duration::from_secs(10)).unwrap());
assert_eq!(forks.len(), 1);
let f = &forks[0];
assert_eq!(f.name, "journal");
assert_eq!(f.trigger, "idle");
assert!(!f.overlap);
assert!(f.after.is_empty());
assert!(f.path.ends_with("journal.md"), "{}", f.path);
assert!(f.prompt.contains("Read the file"), "{}", f.prompt);
assert!(f.prompt.contains(&f.path), "{}", f.prompt);
assert!(f.prompt.contains("parent session oc1"), "{}", f.prompt);
assert!(
f.prompt.contains("Your final message is your report"),
"{}",
f.prompt
);
assert_eq!(h.status_recent_runs(), 1);
}
#[test]
fn opencode_context_gauge_rides_on_the_event() {
let mut h = Harness::new("1h", "0"); h.write_fork(
"distill.md",
"---\nfork: true\nrun_on:\n - context_used: 50%\n---\nDISTILL",
);
h.start_daemon();
assert_ack(h.send_event(oc_event(&h, EventKind::SessionStart, "oc1")));
let mut low = oc_event(&h, EventKind::Stop, "oc1");
low.context_tokens = Some(10_000);
let rx = h.park_stop_wait(low);
std::thread::sleep(Duration::from_millis(400));
assert_ack(h.send_event({
let mut ev = oc_event(&h, EventKind::PromptSubmit, "oc1");
ev.waking = Some(true);
ev
}));
assert!(matches!(
rx.recv_timeout(Duration::from_secs(5)).unwrap(),
ResponseBody::Waited
));
let mut high = oc_event(&h, EventKind::Stop, "oc1");
high.context_tokens = Some(150_000);
let rx = h.park_stop_wait(high);
let forks = wake_forks(rx.recv_timeout(Duration::from_secs(10)).unwrap());
assert_eq!(forks.len(), 1);
assert_eq!(forks[0].name, "distill");
assert_eq!(forks[0].trigger, "context_used:50%");
}
#[test]
fn opencode_reported_window_governs_context_thresholds() {
let mut h = Harness::new("1h", "0"); h.write_fork(
"distill.md",
"---\nfork: true\nrun_on:\n - context_used: 75%\n---\nDISTILL",
);
h.start_daemon();
assert_ack(h.send_event(oc_event(&h, EventKind::SessionStart, "oc1")));
let mut low = oc_event(&h, EventKind::Stop, "oc1");
low.model = Some("claude-sonnet-4-5".to_string());
low.context_tokens = Some(150_000);
low.context_window = Some(1_000_000);
let rx = h.park_stop_wait(low);
std::thread::sleep(Duration::from_millis(400));
assert_ack(h.send_event({
let mut ev = oc_event(&h, EventKind::PromptSubmit, "oc1");
ev.waking = Some(true);
ev
}));
assert!(matches!(
rx.recv_timeout(Duration::from_secs(5)).unwrap(),
ResponseBody::Waited
));
let mut high = oc_event(&h, EventKind::Stop, "oc1");
high.model = Some("claude-sonnet-4-5".to_string());
high.context_tokens = Some(800_000);
high.context_window = Some(1_000_000);
let rx = h.park_stop_wait(high);
let forks = wake_forks(rx.recv_timeout(Duration::from_secs(10)).unwrap());
assert_eq!(forks.len(), 1);
assert_eq!(forks[0].name, "distill");
assert_eq!(forks[0].trigger, "context_used:75%");
}
#[test]
fn opencode_fork_completion_releases_after_dependent() {
let mut h = Harness::new("1s", "0");
h.write_fork("alpha.md", "---\nfork: true\nrun_on: [idle]\n---\nALPHA");
h.write_fork(
"beta.md",
"---\nfork: true\nrun_on: [idle]\nafter: alpha\n---\nBETA",
);
h.start_daemon();
assert_ack(h.send_event(oc_event(&h, EventKind::SessionStart, "oc1")));
let rx = h.park_stop_wait(oc_event(&h, EventKind::Stop, "oc1"));
let forks = wake_forks(rx.recv_timeout(Duration::from_secs(10)).unwrap());
assert_eq!(forks.len(), 1);
assert_eq!(forks[0].name, "alpha");
assert_ack(h.request(RequestBody::ForkSpawned {
session_id: "oc1".into(),
fork: "alpha".into(),
run_ref: "ses_fork_alpha".into(),
}));
let parked = h.park_stop_wait(oc_event(&h, EventKind::Stop, "oc1"));
std::thread::sleep(Duration::from_millis(400));
assert_ack(h.request(RequestBody::ForkCompleted {
session_id: "oc1".into(),
fork: "alpha".into(),
run_ref: "ses_fork_alpha".into(),
status: "completed".into(),
cont: None,
}));
assert!(matches!(
parked.recv_timeout(Duration::from_secs(5)).unwrap(),
ResponseBody::Waited
));
let rx = h.park_stop_wait(oc_event(&h, EventKind::Stop, "oc1"));
let forks = wake_forks(rx.recv_timeout(Duration::from_secs(10)).unwrap());
assert_eq!(forks.len(), 1);
assert_eq!(forks[0].name, "beta");
assert_eq!(forks[0].after, vec!["alpha".to_string()]);
assert!(forks[0].prompt.contains("beta.md"));
}
#[test]
fn opencode_fork_run_sessions_are_never_scheduled() {
let mut h = Harness::new("1s", "0");
h.write_fork(
"journal.md",
"---\nfork: true\nrun_on: [idle]\n---\nJOURNAL",
);
h.start_daemon();
assert_ack(h.send_event(oc_event(&h, EventKind::SessionStart, "oc1")));
let rx = h.park_stop_wait(oc_event(&h, EventKind::Stop, "oc1"));
let forks = wake_forks(rx.recv_timeout(Duration::from_secs(10)).unwrap());
assert_eq!(forks[0].name, "journal");
assert_ack(h.request(RequestBody::ForkSpawned {
session_id: "oc1".into(),
fork: "journal".into(),
run_ref: "ses_fork_run".into(),
}));
assert_ack(h.send_event(oc_event(&h, EventKind::SessionStart, "ses_fork_run")));
assert!(
!h.has_open_session("ses_fork_run"),
"fork run was registered"
);
let rx = h.park_stop_wait(oc_event(&h, EventKind::Stop, "ses_fork_run"));
assert!(matches!(
rx.recv_timeout(Duration::from_secs(5)).unwrap(),
ResponseBody::Waited
));
assert!(
!h.has_open_session("ses_fork_run"),
"fork run was registered"
);
assert_eq!(h.status_recent_runs(), 1);
assert_ack(h.request(RequestBody::ForkCompleted {
session_id: "oc1".into(),
fork: "journal".into(),
run_ref: "ses_fork_run".into(),
status: "completed".into(),
cont: None,
}));
let rx = h.park_stop_wait(oc_event(&h, EventKind::Stop, "ses_fork_run"));
assert!(matches!(
rx.recv_timeout(Duration::from_secs(5)).unwrap(),
ResponseBody::Waited
));
assert!(h.has_open_session("oc1"), "parent must stay registered");
}
#[test]
fn every_fires_at_turn_boundary_without_a_long_idle() {
let mut h = Harness::new("1h", "0");
h.write_fork(
"periodic.md",
"---\nfork: true\nrun_on:\n - every: 1s\n---\nPERIODIC",
);
h.start_daemon();
assert_ack(h.send_event(h.event(EventKind::SessionStart, "s1")));
let rx = h.park_stop_wait(h.event(EventKind::Stop, "s1"));
let payload = wake_payload(rx.recv_timeout(Duration::from_secs(10)).unwrap());
assert!(
payload.contains("due: periodic (trigger: every:1)"),
"{payload}"
);
let rx = h.park_stop_wait(h.event(EventKind::Stop, "s1"));
assert!(
rx.recv_timeout(Duration::from_millis(2500)).is_err(),
"every re-fired during a quiet pause"
);
assert_ack(h.send_event(h.prompt_submit("s1", true)));
let _ = rx.recv_timeout(Duration::from_secs(5)); let rx = h.park_stop_wait(h.event(EventKind::Stop, "s1"));
let payload = wake_payload(rx.recv_timeout(Duration::from_secs(10)).unwrap());
assert!(
payload.contains("due: periodic (trigger: every:1)"),
"{payload}"
);
}
#[test]
fn busy_poll_fires_every_but_never_idle() {
let mut h = Harness::new("1s", "0");
h.write_fork(
"idler.md",
"---\nfork: true\nrun_on:\n - idle: 1s\n---\nIDLER",
);
h.write_fork(
"periodic.md",
"---\nfork: true\nrun_on:\n - every: 1s\n---\nPERIODIC",
);
h.start_daemon();
assert_ack(h.send_event(oc_event(&h, EventKind::SessionStart, "oc1")));
let mut ev = oc_event(&h, EventKind::Stop, "oc1");
ev.busy = Some(true);
let rx = h.park_stop_wait(ev);
let forks = wake_forks(rx.recv_timeout(Duration::from_secs(10)).unwrap());
assert_eq!(
forks.len(),
1,
"only the every fork may fire on a busy poll"
);
assert_eq!(forks[0].name, "periodic");
assert_eq!(forks[0].trigger, "every:1");
let rx = h.park_stop_wait(oc_event(&h, EventKind::Stop, "oc1"));
let forks = wake_forks(rx.recv_timeout(Duration::from_secs(10)).unwrap());
assert!(forks.iter().any(|f| f.name == "idler"), "{forks:?}");
}
#[test]
fn chain_continue_rearms_the_fork_within_the_pause() {
let mut h = Harness::new("1s", "0").wake_grace_secs(0);
h.write_fork(
"goal.md",
"---\nfork: true\nrun_on: [idle]\nchain: true\n---\nGOAL",
);
h.start_daemon();
h.write_transcript(100);
assert_ack(h.send_event(h.event_t(EventKind::SessionStart, "s1")));
let rx = h.park_stop_wait(h.event_t(EventKind::Stop, "s1"));
let payload = wake_payload(rx.recv_timeout(Duration::from_secs(10)).unwrap());
assert!(payload.contains("due: goal"), "{payload}");
assert!(payload.contains("<<autofork:continue>>"), "{payload}");
h.append_fork_spawn("toolu_g1", "goal");
h.append_completion_notification_result(
"toolu_g1",
"completed",
"goal not met, queued more work\n<<autofork:continue>>",
);
assert_ack(h.send_event(h.prompt_submit_notif_cont("s1", "toolu_g1")));
let rx2 = h.park_stop_wait(h.event_t(EventKind::Stop, "s1"));
let payload = wake_payload(rx2.recv_timeout(Duration::from_secs(10)).unwrap());
assert!(payload.contains("due: goal"), "{payload}");
h.append_fork_spawn("toolu_g2", "goal");
h.append_completion_notification("toolu_g2", "completed");
assert_ack(h.send_event(h.prompt_submit_notif("s1", "toolu_g2", "completed")));
let rx3 = h.park_stop_wait(h.event_t(EventKind::Stop, "s1"));
assert!(
rx3.recv_timeout(Duration::from_millis(2500)).is_err(),
"chain re-fired without the sentinel"
);
}
#[test]
fn chain_continue_via_prompt_ids_without_transcript_notification() {
let mut h = Harness::new("1s", "0").wake_grace_secs(0);
h.write_fork(
"goal.md",
"---\nfork: true\nrun_on: [idle]\nchain: true\n---\nGOAL",
);
h.start_daemon();
h.write_transcript(100);
assert_ack(h.send_event(h.event_t(EventKind::SessionStart, "s1")));
let rx = h.park_stop_wait(h.event_t(EventKind::Stop, "s1"));
wake_payload(rx.recv_timeout(Duration::from_secs(10)).unwrap());
h.append_fork_spawn("toolu_g1", "goal");
assert_ack(h.send_event(h.prompt_submit_notif_cont("s1", "toolu_g1")));
let rx2 = h.park_stop_wait(h.event_t(EventKind::Stop, "s1"));
let payload = wake_payload(rx2.recv_timeout(Duration::from_secs(10)).unwrap());
assert!(payload.contains("due: goal"), "{payload}");
}
#[test]
fn chain_sentinel_ignored_without_optin() {
let mut h = Harness::new("1s", "0").wake_grace_secs(0);
h.write_fork("journal.md", "---\nfork: true\nrun_on: [idle]\n---\nJ");
h.start_daemon();
h.write_transcript(100);
assert_ack(h.send_event(h.event_t(EventKind::SessionStart, "s1")));
let rx = h.park_stop_wait(h.event_t(EventKind::Stop, "s1"));
let payload = wake_payload(rx.recv_timeout(Duration::from_secs(10)).unwrap());
assert!(!payload.contains("<<autofork:continue>>"), "{payload}");
h.append_fork_spawn("toolu_j1", "journal");
h.append_completion_notification_result(
"toolu_j1",
"completed",
"quoting the docs\n<<autofork:continue>>",
);
assert_ack(h.send_event(h.prompt_submit_notif_cont("s1", "toolu_j1")));
let rx2 = h.park_stop_wait(h.event_t(EventKind::Stop, "s1"));
assert!(
rx2.recv_timeout(Duration::from_millis(2500)).is_err(),
"sentinel re-armed a fork without chain: true"
);
}
#[test]
fn chain_limit_caps_refires() {
let mut h = Harness::new("1s", "0").wake_grace_secs(0);
h.write_fork(
"goal.md",
"---\nfork: true\nrun_on: [idle]\nchain: true\nchain_limit: 2\n---\nGOAL",
);
h.start_daemon();
h.write_transcript(100);
assert_ack(h.send_event(h.event_t(EventKind::SessionStart, "s1")));
let rx = h.park_stop_wait(h.event_t(EventKind::Stop, "s1"));
wake_payload(rx.recv_timeout(Duration::from_secs(10)).unwrap());
h.append_fork_spawn("toolu_g1", "goal");
h.append_completion_notification_result("toolu_g1", "completed", "more\n<<autofork:continue>>");
assert_ack(h.send_event(h.prompt_submit_notif_cont("s1", "toolu_g1")));
let rx2 = h.park_stop_wait(h.event_t(EventKind::Stop, "s1"));
wake_payload(rx2.recv_timeout(Duration::from_secs(10)).unwrap());
h.append_fork_spawn("toolu_g2", "goal");
h.append_completion_notification_result("toolu_g2", "completed", "more\n<<autofork:continue>>");
assert_ack(h.send_event(h.prompt_submit_notif_cont("s1", "toolu_g2")));
let rx3 = h.park_stop_wait(h.event_t(EventKind::Stop, "s1"));
assert!(
rx3.recv_timeout(Duration::from_millis(2500)).is_err(),
"chain exceeded its chain_limit"
);
}
#[test]
fn opencode_continue_field_rearms_and_settles() {
let mut h = Harness::new("1s", "0").wake_grace_secs(0);
h.write_fork(
"goal.md",
"---\nfork: true\nrun_on: [idle]\nchain: true\n---\nGOAL",
);
h.start_daemon();
assert_ack(h.send_event(oc_event(&h, EventKind::SessionStart, "oc1")));
let rx = h.park_stop_wait(oc_event(&h, EventKind::Stop, "oc1"));
let forks = wake_forks(rx.recv_timeout(Duration::from_secs(10)).unwrap());
assert_eq!(forks[0].name, "goal");
assert!(forks[0].chain, "structured spec must carry chain");
assert!(forks[0].prompt.contains("<<autofork:continue>>"));
assert_ack(h.request(RequestBody::ForkSpawned {
session_id: "oc1".into(),
fork: "goal".into(),
run_ref: "ses_g1".into(),
}));
let parked = h.park_stop_wait(oc_event(&h, EventKind::Stop, "oc1"));
std::thread::sleep(Duration::from_millis(400));
assert_ack(h.request(RequestBody::ForkCompleted {
session_id: "oc1".into(),
fork: "goal".into(),
run_ref: "ses_g1".into(),
status: "completed".into(),
cont: Some(true),
}));
assert!(matches!(
parked.recv_timeout(Duration::from_secs(5)).unwrap(),
ResponseBody::Waited
));
let rx2 = h.park_stop_wait(oc_event(&h, EventKind::Stop, "oc1"));
let forks = wake_forks(rx2.recv_timeout(Duration::from_secs(10)).unwrap());
assert_eq!(forks[0].name, "goal");
assert_ack(h.request(RequestBody::ForkSpawned {
session_id: "oc1".into(),
fork: "goal".into(),
run_ref: "ses_g2".into(),
}));
assert_ack(h.request(RequestBody::ForkCompleted {
session_id: "oc1".into(),
fork: "goal".into(),
run_ref: "ses_g2".into(),
status: "completed".into(),
cont: None,
}));
let rx3 = h.park_stop_wait(oc_event(&h, EventKind::Stop, "oc1"));
assert!(
rx3.recv_timeout(Duration::from_millis(2500)).is_err(),
"chain re-fired after settling"
);
}
#[test]
fn runaway_breaker_stops_an_epoch_pumped_chain() {
let mut h = Harness::new("1s", "0")
.wake_grace_secs(0)
.chain_grace_secs(0);
h.append_config("runaway_limit = 2");
h.write_fork(
"goal.md",
"---\nfork: true\nrun_on: [idle]\nchain: true\n---\nGOAL",
);
h.start_daemon();
assert_ack(h.send_event(oc_event(&h, EventKind::SessionStart, "oc1")));
let pump = |h: &Harness| {
let mut ev = oc_event(h, EventKind::PromptSubmit, "oc1");
ev.waking = Some(true);
assert_ack(h.send_event(ev));
};
let rx = h.park_stop_wait(oc_event(&h, EventKind::Stop, "oc1"));
let forks = wake_forks(rx.recv_timeout(Duration::from_secs(10)).unwrap());
assert_eq!(forks[0].name, "goal");
assert_ack(h.request(RequestBody::ForkSpawned {
session_id: "oc1".into(),
fork: "goal".into(),
run_ref: "ses_g1".into(),
}));
assert_ack(h.request(RequestBody::ForkCompleted {
session_id: "oc1".into(),
fork: "goal".into(),
run_ref: "ses_g1".into(),
status: "completed".into(),
cont: Some(true),
}));
pump(&h);
let rx2 = h.park_stop_wait(oc_event(&h, EventKind::Stop, "oc1"));
let forks = wake_forks(rx2.recv_timeout(Duration::from_secs(10)).unwrap());
assert_eq!(forks[0].name, "goal");
assert_ack(h.request(RequestBody::ForkSpawned {
session_id: "oc1".into(),
fork: "goal".into(),
run_ref: "ses_g2".into(),
}));
assert_ack(h.request(RequestBody::ForkCompleted {
session_id: "oc1".into(),
fork: "goal".into(),
run_ref: "ses_g2".into(),
status: "completed".into(),
cont: Some(true),
}));
pump(&h);
let rx3 = h.park_stop_wait(oc_event(&h, EventKind::Stop, "oc1"));
assert!(
rx3.recv_timeout(Duration::from_millis(2500)).is_err(),
"runaway breaker failed: the epoch-pumped chain re-fired past the cap"
);
}
#[test]
fn chain_grace_downgrades_duplicate_activity_reports() {
let mut h = Harness::new("1s", "0").wake_grace_secs(0);
h.write_fork(
"goal.md",
"---\nfork: true\nrun_on: [idle]\nchain: true\nchain_limit: 2\n---\nGOAL",
);
h.start_daemon();
assert_ack(h.send_event(oc_event(&h, EventKind::SessionStart, "oc1")));
let duplicate_report = |h: &Harness| {
let mut ev = oc_event(h, EventKind::PromptSubmit, "oc1");
ev.waking = Some(true);
assert_ack(h.send_event(ev));
};
let rx = h.park_stop_wait(oc_event(&h, EventKind::Stop, "oc1"));
let forks = wake_forks(rx.recv_timeout(Duration::from_secs(10)).unwrap());
assert_eq!(forks[0].name, "goal");
assert_ack(h.request(RequestBody::ForkSpawned {
session_id: "oc1".into(),
fork: "goal".into(),
run_ref: "ses_g1".into(),
}));
assert_ack(h.request(RequestBody::ForkCompleted {
session_id: "oc1".into(),
fork: "goal".into(),
run_ref: "ses_g1".into(),
status: "completed".into(),
cont: Some(true),
}));
duplicate_report(&h);
let rx2 = h.park_stop_wait(oc_event(&h, EventKind::Stop, "oc1"));
let forks = wake_forks(rx2.recv_timeout(Duration::from_secs(10)).unwrap());
assert_eq!(forks[0].name, "goal");
assert_ack(h.request(RequestBody::ForkSpawned {
session_id: "oc1".into(),
fork: "goal".into(),
run_ref: "ses_g2".into(),
}));
assert_ack(h.request(RequestBody::ForkCompleted {
session_id: "oc1".into(),
fork: "goal".into(),
run_ref: "ses_g2".into(),
status: "completed".into(),
cont: Some(true),
}));
duplicate_report(&h);
let rx3 = h.park_stop_wait(oc_event(&h, EventKind::Stop, "oc1"));
assert!(
rx3.recv_timeout(Duration::from_millis(2500)).is_err(),
"duplicate activity report reset the pause and defeated chain_limit"
);
}
#[test]
fn overlap_false_holds_across_pause_resets() {
let mut h = Harness::new("1s", "0").wake_grace_secs(0);
h.write_fork("journal.md", "---\nfork: true\nrun_on: [idle]\n---\nJ");
h.start_daemon();
assert_ack(h.send_event(oc_event(&h, EventKind::SessionStart, "oc1")));
let rx = h.park_stop_wait(oc_event(&h, EventKind::Stop, "oc1"));
let forks = wake_forks(rx.recv_timeout(Duration::from_secs(10)).unwrap());
assert_eq!(forks[0].name, "journal");
assert_ack(h.request(RequestBody::ForkSpawned {
session_id: "oc1".into(),
fork: "journal".into(),
run_ref: "ses_j1".into(),
}));
let mut ev = oc_event(&h, EventKind::PromptSubmit, "oc1");
ev.waking = Some(true);
assert_ack(h.send_event(ev));
let rx2 = h.park_stop_wait(oc_event(&h, EventKind::Stop, "oc1"));
assert!(
rx2.recv_timeout(Duration::from_millis(2500)).is_err(),
"overlap: false fork re-fired while its run was still in flight"
);
assert_ack(h.request(RequestBody::ForkCompleted {
session_id: "oc1".into(),
fork: "journal".into(),
run_ref: "ses_j1".into(),
status: "completed".into(),
cont: None,
}));
let mut ev = oc_event(&h, EventKind::PromptSubmit, "oc1");
ev.waking = Some(true);
assert_ack(h.send_event(ev));
let rx3 = h.park_stop_wait(oc_event(&h, EventKind::Stop, "oc1"));
let forks = wake_forks(rx3.recv_timeout(Duration::from_secs(10)).unwrap());
assert_eq!(forks[0].name, "journal");
}
#[test]
fn gate_holds_idle_forks_until_chain_settles() {
let mut h = Harness::new("1s", "0").wake_grace_secs(0);
h.write_fork(
"goal.md",
"---\nfork: true\nrun_on: [idle: 0s]\nchain: true\ngate: true\n---\nGOAL",
);
h.write_fork("handover.md", "---\nfork: true\nrun_on: [idle: 1s]\n---\nH");
h.start_daemon();
h.write_transcript(100);
assert_ack(h.send_event(h.event_t(EventKind::SessionStart, "s1")));
let rx = h.park_stop_wait(h.event_t(EventKind::Stop, "s1"));
let payload = wake_payload(rx.recv_timeout(Duration::from_secs(10)).unwrap());
assert!(payload.contains("due: goal"), "{payload}");
assert!(!payload.contains("due: handover"), "{payload}");
h.append_fork_spawn("toolu_g1", "goal");
let rx2 = h.park_stop_wait(h.event_t(EventKind::Stop, "s1"));
assert!(
rx2.recv_timeout(Duration::from_millis(2500)).is_err(),
"gate failed to hold the idle fork mid-chain"
);
h.append_completion_notification_result("toolu_g1", "completed", "more\n<<autofork:continue>>");
assert_ack(h.send_event(h.prompt_submit_notif_cont("s1", "toolu_g1")));
let rx3 = h.park_stop_wait(h.event_t(EventKind::Stop, "s1"));
let payload = wake_payload(rx3.recv_timeout(Duration::from_secs(10)).unwrap());
assert!(payload.contains("due: goal"), "{payload}");
assert!(!payload.contains("due: handover"), "{payload}");
h.append_fork_spawn("toolu_g2", "goal");
h.append_completion_notification("toolu_g2", "completed");
assert_ack(h.send_event(h.prompt_submit_notif("s1", "toolu_g2", "completed")));
let rx4 = h.park_stop_wait(h.event_t(EventKind::Stop, "s1"));
let payload = wake_payload(rx4.recv_timeout(Duration::from_secs(10)).unwrap());
assert!(payload.contains("due: handover"), "{payload}");
assert!(!payload.contains("due: goal"), "{payload}");
}
#[test]
fn gate_belt_lifts_a_fumbled_wake() {
let mut h = Harness::new("1s", "0")
.wake_grace_secs(0)
.gate_grace_secs(1);
h.write_fork(
"goal.md",
"---\nfork: true\nrun_on: [idle: 0s]\nchain: true\ngate: true\n---\nGOAL",
);
h.write_fork("handover.md", "---\nfork: true\nrun_on: [idle: 1s]\n---\nH");
h.start_daemon();
h.write_transcript(100);
assert_ack(h.send_event(h.event_t(EventKind::SessionStart, "s1")));
let rx = h.park_stop_wait(h.event_t(EventKind::Stop, "s1"));
let payload = wake_payload(rx.recv_timeout(Duration::from_secs(10)).unwrap());
assert!(payload.contains("due: goal"), "{payload}");
let rx2 = h.park_stop_wait(h.event_t(EventKind::Stop, "s1"));
let payload = wake_payload(rx2.recv_timeout(Duration::from_secs(10)).unwrap());
assert!(payload.contains("due: handover"), "{payload}");
assert!(!payload.contains("due: goal"), "{payload}");
}
#[test]
fn lifecycle_hooks_fire_across_the_session_life() {
let mut h = Harness::new("30m", "0");
let log = h.write_logging_hook(
"lease.md",
"[session_start, activity, \"idle: 1s\", session_end]",
);
h.start_daemon();
let mut start = h.event(EventKind::SessionStart, "s1");
start.source = Some("startup".into());
assert_ack(h.send_event(start));
let lines = h.wait_for_hook_lines(&log, 1, Duration::from_secs(5));
assert_eq!(lines[0], "session_start|startup|||s1");
let mut again = h.event(EventKind::SessionStart, "s1");
again.source = Some("compact".into());
assert_ack(h.send_event(again));
assert_ack(h.send_event(h.prompt_submit("s1", true)));
let lines = h.wait_for_hook_lines(&log, 2, Duration::from_secs(5));
assert_eq!(lines[1], "activity||||s1");
let rx = h.park_stop_wait(h.event(EventKind::Stop, "s1"));
let lines = h.wait_for_hook_lines(&log, 3, Duration::from_secs(10));
assert_eq!(lines[2], "idle|||1|s1");
assert!(
rx.try_recv().is_err(),
"idle hook firing must not resolve the parked poll"
);
let mut end = h.event(EventKind::SessionEnd, "s1");
end.reason = Some("logout".into());
assert_ack(h.send_event(end));
let lines = h.wait_for_hook_lines(&log, 4, Duration::from_secs(5));
assert_eq!(lines[3], "session_end||logout||s1");
assert_ack(h.send_event(h.event(EventKind::SessionEnd, "s1")));
std::thread::sleep(Duration::from_millis(600));
assert_eq!(
h.wait_for_hook_lines(&log, 4, Duration::from_secs(1)).len(),
4
);
}
#[test]
fn idle_hook_fires_once_per_pause_and_rearms_on_activity() {
let mut h = Harness::new("30m", "0");
let log = h.write_logging_hook("park.md", "[\"idle: 1s\"]");
h.start_daemon();
assert_ack(h.send_event(h.event(EventKind::SessionStart, "s1")));
let rx = h.park_stop_wait(h.event(EventKind::Stop, "s1"));
let lines = h.wait_for_hook_lines(&log, 1, Duration::from_secs(10));
assert_eq!(lines.len(), 1);
assert_ack(h.send_event(h.prompt_submit("s1", false)));
let _ = rx.recv_timeout(Duration::from_secs(5)).unwrap();
let _rx2 = h.park_stop_wait(h.event(EventKind::Stop, "s1"));
std::thread::sleep(Duration::from_millis(1800));
assert_eq!(
h.wait_for_hook_lines(&log, 1, Duration::from_secs(1)).len(),
1
);
assert_ack(h.send_event(h.prompt_submit("s1", true)));
let _rx3 = h.park_stop_wait(h.event(EventKind::Stop, "s1"));
let lines = h.wait_for_hook_lines(&log, 2, Duration::from_secs(10));
assert_eq!(lines.len(), 2);
}
#[test]
fn poll_loss_close_fires_session_end_with_reason_lost() {
let mut h = Harness::new("1s", "0").poll_grace_ms(300);
let log = h.write_logging_hook("cleanup.md", "[session_end]");
h.start_daemon();
assert_ack(h.send_event(h.event(EventKind::SessionStart, "s1")));
h.drop_stop_wait(h.event(EventKind::Stop, "s1"));
let lines = h.wait_for_hook_lines(&log, 1, Duration::from_secs(10));
assert_eq!(lines[0], "session_end||lost||s1");
}
#[test]
fn resume_hook_fires_only_for_the_resume_source() {
let mut h = Harness::new("30m", "0");
let log = h.write_logging_hook("rejoin.md", "[resume]");
h.start_daemon();
let mut startup = h.event(EventKind::SessionStart, "s1");
startup.source = Some("startup".into());
assert_ack(h.send_event(startup));
let mut resume = h.event(EventKind::SessionStart, "s2");
resume.source = Some("resume".into());
assert_ack(h.send_event(resume));
let lines = h.wait_for_hook_lines(&log, 1, Duration::from_secs(5));
assert_eq!(lines, vec!["resume|resume|||s2".to_string()]);
std::thread::sleep(Duration::from_millis(400));
assert_eq!(
h.wait_for_hook_lines(&log, 1, Duration::from_secs(1)).len(),
1
);
}
fn cx_event(h: &Harness, kind: EventKind, session: &str) -> Event {
let mut ev = h.event(kind, session);
ev.client = Some("codex".to_string());
ev
}
#[test]
fn codex_wake_carries_structured_forks_and_gauge() {
let mut h = Harness::new("1s", "0");
h.write_fork(
"journal.md",
"---\nfork: true\nrun_on: [idle]\n---\nwrite the journal now",
);
h.write_fork(
"distill.md",
"---\nfork: true\nrun_on:\n - context_used: 50%\n---\nDISTILL",
);
h.start_daemon();
let sid = "01a01f24-3113-76c3-a00a-74ac3948e630";
assert_ack(h.send_event(cx_event(&h, EventKind::SessionStart, sid)));
let mut ev = cx_event(&h, EventKind::Stop, sid);
ev.context_tokens = Some(160_000);
ev.context_window = Some(258_400);
let rx = h.park_stop_wait(ev.clone());
let forks = wake_forks(rx.recv_timeout(Duration::from_secs(10)).unwrap());
assert_eq!(forks.len(), 1);
assert_eq!(forks[0].name, "distill");
assert_eq!(forks[0].trigger, "context_used:50%");
let rx = h.park_stop_wait(ev);
let forks = wake_forks(rx.recv_timeout(Duration::from_secs(10)).unwrap());
assert_eq!(forks.len(), 1);
assert_eq!(forks[0].name, "journal");
assert!(forks[0].prompt.contains(&format!("parent session {sid}")));
assert_eq!(h.status_recent_runs(), 2);
}
#[test]
fn codex_chain_grace_downgrades_duplicate_activity() {
let mut h = Harness::new("1s", "0").wake_grace_secs(0);
h.write_fork(
"goal.md",
"---\nfork: true\nrun_on: [idle]\nchain: true\nchain_limit: 2\n---\nGOAL",
);
h.start_daemon();
let sid = "01a01f24-aaaa-76c3-a00a-74ac3948e630";
assert_ack(h.send_event(cx_event(&h, EventKind::SessionStart, sid)));
let duplicate_report = |h: &Harness| {
let mut ev = cx_event(h, EventKind::PromptSubmit, sid);
ev.waking = Some(true);
assert_ack(h.send_event(ev));
};
for run in [
"01a01f30-0001-7000-8000-000000000001",
"01a01f30-0002-7000-8000-000000000002",
] {
let rx = h.park_stop_wait(cx_event(&h, EventKind::Stop, sid));
let forks = wake_forks(rx.recv_timeout(Duration::from_secs(10)).unwrap());
assert_eq!(forks[0].name, "goal");
assert_ack(h.request(RequestBody::ForkSpawned {
session_id: sid.into(),
fork: "goal".into(),
run_ref: run.into(),
}));
assert_ack(h.request(RequestBody::ForkCompleted {
session_id: sid.into(),
fork: "goal".into(),
run_ref: run.into(),
status: "completed".into(),
cont: Some(true),
}));
duplicate_report(&h);
}
let rx = h.park_stop_wait(cx_event(&h, EventKind::Stop, sid));
assert!(
rx.recv_timeout(Duration::from_millis(2500)).is_err(),
"duplicate activity report reset the pause and defeated chain_limit"
);
}
#[test]
fn codex_fork_run_sessions_are_never_scheduled() {
let mut h = Harness::new("1s", "0");
h.write_fork("journal.md", "---\nfork: true\nrun_on: [idle]\n---\nJ");
h.start_daemon();
let sid = "01a01f24-bbbb-76c3-a00a-74ac3948e630";
let run = "01a01f30-cccc-7000-8000-000000000001";
assert_ack(h.send_event(cx_event(&h, EventKind::SessionStart, sid)));
let rx = h.park_stop_wait(cx_event(&h, EventKind::Stop, sid));
let forks = wake_forks(rx.recv_timeout(Duration::from_secs(10)).unwrap());
assert_eq!(forks[0].name, "journal");
assert_ack(h.request(RequestBody::ForkSpawned {
session_id: sid.into(),
fork: "journal".into(),
run_ref: run.into(),
}));
assert_ack(h.send_event(cx_event(&h, EventKind::SessionStart, run)));
assert!(!h.has_open_session(run), "fork run was registered");
let rx = h.park_stop_wait(cx_event(&h, EventKind::Stop, run));
assert!(matches!(
rx.recv_timeout(Duration::from_secs(5)).unwrap(),
ResponseBody::Waited
));
assert_eq!(h.status_recent_runs(), 1);
}
#[test]
fn wake_forks_carry_resolved_model_and_mode() {
let mut h = Harness::new("1s", "0");
h.append_config("[fork_models]");
h.append_config("codex = \"gpt-5.1-codex-mini\"");
h.append_config("\"claude-code\" = \"haiku\"");
h.write_fork(
"journal.md",
"---\nfork: true\nrun_on: [idle]\nmodel:\n opencode: anthropic/claude-haiku-4-5\nmode:\n codex: workspace-write\n---\nJ",
);
h.start_daemon();
assert_ack(h.send_event(oc_event(&h, EventKind::SessionStart, "oc-m")));
let rx = h.park_stop_wait(oc_event(&h, EventKind::Stop, "oc-m"));
let forks = wake_forks(rx.recv_timeout(Duration::from_secs(10)).unwrap());
assert_eq!(
forks[0].model.as_deref(),
Some("anthropic/claude-haiku-4-5")
);
assert_eq!(forks[0].mode, None);
let sid = "01a01f24-dddd-76c3-a00a-74ac3948e630";
assert_ack(h.send_event(cx_event(&h, EventKind::SessionStart, sid)));
let rx = h.park_stop_wait(cx_event(&h, EventKind::Stop, sid));
let forks = wake_forks(rx.recv_timeout(Duration::from_secs(10)).unwrap());
assert_eq!(forks[0].model.as_deref(), Some("gpt-5.1-codex-mini"));
assert_eq!(forks[0].mode.as_deref(), Some("workspace-write"));
assert_ack(h.send_event(h.event(EventKind::SessionStart, "cc1")));
let rx = h.park_stop_wait(h.event(EventKind::Stop, "cc1"));
let forks = wake_forks(rx.recv_timeout(Duration::from_secs(10)).unwrap());
assert_eq!(forks[0].model.as_deref(), Some("haiku"));
}
#[test]
fn codex_peek_due_runs_the_goal_fast_path() {
let mut h = Harness::new("1h", "0"); h.write_fork(
"goal.md",
"---\nfork: true\nrun_on:\n - idle: 0s\nchain: true\n---\nGOAL",
);
h.write_fork(
"note.md",
"---\nfork: true\nrun_on:\n - idle: 0s\n---\nNOTE",
);
h.start_daemon();
let sid = "01a01f24-eeee-76c3-a00a-74ac3948e630";
assert_ack(h.send_event(cx_event(&h, EventKind::SessionStart, sid)));
let ResponseBody::Due { forks } = h.request(RequestBody::PeekDue {
session_id: sid.into(),
}) else {
panic!("expected Due");
};
assert_eq!(forks.len(), 1, "only the chain fork rides the fast path");
assert_eq!(forks[0].name, "goal");
assert!(forks[0].chain);
assert_ack(h.request(RequestBody::ForkSpawned {
session_id: sid.into(),
fork: "goal".into(),
run_ref: "01a01f30-aaaa-7000-8000-000000000001".into(),
}));
assert_ack(h.request(RequestBody::ForkCompleted {
session_id: sid.into(),
fork: "goal".into(),
run_ref: "01a01f30-aaaa-7000-8000-000000000001".into(),
status: "completed".into(),
cont: Some(true),
}));
let ResponseBody::Due { forks } = h.request(RequestBody::PeekDue {
session_id: sid.into(),
}) else {
panic!("expected Due");
};
assert_eq!(forks.len(), 1);
assert_ack(h.request(RequestBody::ForkSpawned {
session_id: sid.into(),
fork: "goal".into(),
run_ref: "01a01f30-aaaa-7000-8000-000000000002".into(),
}));
assert_ack(h.request(RequestBody::ForkCompleted {
session_id: sid.into(),
fork: "goal".into(),
run_ref: "01a01f30-aaaa-7000-8000-000000000002".into(),
status: "completed".into(),
cont: None, }));
let ResponseBody::Due { forks } = h.request(RequestBody::PeekDue {
session_id: sid.into(),
}) else {
panic!("expected Due");
};
assert!(forks.is_empty(), "settled chain must not re-fire");
let rx = h.park_stop_wait(cx_event(&h, EventKind::Stop, sid));
let forks = wake_forks(rx.recv_timeout(Duration::from_secs(10)).unwrap());
assert_eq!(forks.len(), 1);
assert_eq!(forks[0].name, "note");
}
#[test]
fn spooled_reports_deliver_once_in_order() {
let mut h = Harness::new("1h", "0");
h.start_daemon();
assert_ack(h.send_event(h.event(EventKind::SessionStart, "cc-spool")));
assert_ack(h.request(RequestBody::SpoolReport {
session_id: "cc-spool".into(),
fork: "journal".into(),
text: "first".into(),
}));
assert_ack(h.request(RequestBody::SpoolReport {
session_id: "cc-spool".into(),
fork: "notes".into(),
text: "second".into(),
}));
let ResponseBody::Reports { blocks } = h.request(RequestBody::TakeReports {
session_id: "cc-spool".into(),
wait_ms: None,
}) else {
panic!("expected Reports");
};
assert_eq!(blocks, vec!["first".to_string(), "second".to_string()]);
let ResponseBody::Reports { blocks } = h.request(RequestBody::TakeReports {
session_id: "cc-spool".into(),
wait_ms: None,
}) else {
panic!("expected Reports");
};
assert!(blocks.is_empty(), "taking clears the spool");
}
#[test]
fn codex_poll_reserves_idle_zero_chain_forks_for_peek_due() {
let mut h = Harness::new("1h", "0");
h.write_fork(
"goal.md",
"---\nfork: true\nrun_on:\n - idle: 0s\nchain: true\n---\nGOAL",
);
h.start_daemon();
let sid = "01a01f24-ffff-76c3-a00a-74ac3948e630";
assert_ack(h.send_event(cx_event(&h, EventKind::SessionStart, sid)));
let rx = h.park_stop_wait(cx_event(&h, EventKind::Stop, sid));
assert!(
rx.recv_timeout(Duration::from_millis(2500)).is_err(),
"the poll grabbed a fast-path fork"
);
let ResponseBody::Due { forks } = h.request(RequestBody::PeekDue {
session_id: sid.into(),
}) else {
panic!("expected Due");
};
assert_eq!(forks.len(), 1);
assert_eq!(forks[0].name, "goal");
assert_ack(h.send_event(oc_event(&h, EventKind::SessionStart, "oc-goal")));
let rx = h.park_stop_wait(oc_event(&h, EventKind::Stop, "oc-goal"));
let forks = wake_forks(rx.recv_timeout(Duration::from_secs(10)).unwrap());
assert_eq!(forks[0].name, "goal");
}
#[test]
fn take_final_runs_flushes_only_unrun_idle_forks_in_order() {
let mut h = Harness::new("1s", "0");
h.write_fork("early.md", "---\nfork: true\nrun_on:\n - idle: 1s\n---\nE");
h.write_fork(
"journal.md",
"---\nfork: true\nrun_on:\n - idle: 30m\n---\nJ",
);
h.write_fork(
"handover.md",
"---\nfork: true\nrun_on:\n - idle: 30m\nafter: [journal]\n---\nH",
);
h.start_daemon();
assert_ack(h.send_event(h.event(EventKind::SessionStart, "cc-flush")));
let rx = h.park_stop_wait(h.event(EventKind::Stop, "cc-flush"));
let forks = wake_forks(rx.recv_timeout(Duration::from_secs(10)).unwrap());
assert_eq!(forks.len(), 1);
assert_eq!(forks[0].name, "early");
let ResponseBody::Due { forks } = h.request(RequestBody::TakeFinalRuns {
session_id: "cc-flush".into(),
}) else {
panic!("expected Due");
};
assert_eq!(
forks.iter().map(|f| f.name.as_str()).collect::<Vec<_>>(),
vec!["journal", "handover"]
);
assert!(forks[0].after.is_empty());
assert_eq!(forks[1].after, vec!["journal".to_string()]);
assert!(
forks[0].trigger.contains("at close"),
"{}",
forks[0].trigger
);
let ResponseBody::Due { forks } = h.request(RequestBody::TakeFinalRuns {
session_id: "cc-flush".into(),
}) else {
panic!("expected Due");
};
assert!(forks.is_empty(), "final runs must stamp what they hand out");
}
#[test]
fn a_gate_fork_leads_the_close_batch_instead_of_swallowing_it() {
let mut h = Harness::new("1s", "0");
h.write_fork(
"goal-supervisor.md",
"---\nfork: true\nrun_on: [idle: 0s]\ngate: true\n---\nSUPERVISE",
);
h.write_fork(
"context-curator.md",
"---\nfork: true\nrun_on:\n - idle: 30m\n---\nCURATE",
);
h.write_fork(
"handover.md",
"---\nfork: true\nrun_on:\n - idle: 30m\n---\nHAND OVER",
);
h.start_daemon();
assert_ack(h.send_event(h.event(EventKind::SessionStart, "s-gated")));
let ResponseBody::Due { forks } = h.request(RequestBody::TakeFinalRuns {
session_id: "s-gated".into(),
}) else {
panic!("expected Due");
};
let names: Vec<&str> = forks.iter().map(|f| f.name.as_str()).collect();
assert_eq!(names.first(), Some(&"goal-supervisor"), "{names:?}");
assert!(names.contains(&"handover"), "{names:?}");
assert!(names.contains(&"context-curator"), "{names:?}");
for spec in &forks {
assert!(
!spec.after.contains(&"goal-supervisor".to_string()),
"{:?} should not depend on the gate: {:?}",
spec.name,
spec.after
);
}
}
#[test]
fn wake_forks_carry_model_fallback_lists() {
let mut h = Harness::new("1s", "0");
h.append_config("[fork_models]");
h.append_config("codex = [\"gpt-5.6-luna\", \"gpt-5.5\"]");
h.write_fork(
"journal.md",
"---\nfork: true\nrun_on: [idle]\nmodel:\n opencode: [github-copilot/gemini-3.7-flash, anthropic/claude-haiku-4-5]\n---\nJ",
);
h.start_daemon();
assert_ack(h.send_event(oc_event(&h, EventKind::SessionStart, "oc-fb")));
let rx = h.park_stop_wait(oc_event(&h, EventKind::Stop, "oc-fb"));
let forks = wake_forks(rx.recv_timeout(Duration::from_secs(10)).unwrap());
assert_eq!(
forks[0].model.as_deref(),
Some("github-copilot/gemini-3.7-flash")
);
assert_eq!(
forks[0].model_fallbacks,
vec!["anthropic/claude-haiku-4-5".to_string()]
);
let sid = "01a01f24-abcd-76c3-a00a-74ac3948e630";
assert_ack(h.send_event(cx_event(&h, EventKind::SessionStart, sid)));
let rx = h.park_stop_wait(cx_event(&h, EventKind::Stop, sid));
let forks = wake_forks(rx.recv_timeout(Duration::from_secs(10)).unwrap());
assert_eq!(forks[0].model.as_deref(), Some("gpt-5.6-luna"));
assert_eq!(forks[0].model_fallbacks, vec!["gpt-5.5".to_string()]);
}
#[test]
fn fork_models_resolve_by_parent_model_with_default_catchall() {
let mut h = Harness::new("1s", "0");
h.append_config("[fork_models.\"claude-code\"]");
h.append_config("opus = \"sonnet\"");
h.append_config("fable = [\"sonnet\", \"haiku\"]");
h.append_config("default = \"haiku\"");
h.write_fork("journal.md", "---\nfork: true\nrun_on: [idle]\n---\nJ");
h.start_daemon();
let mut start = h.event(EventKind::SessionStart, "cc-opus");
start.model = Some("opus".to_string());
assert_ack(h.send_event(start));
let mut stop = h.event(EventKind::Stop, "cc-opus");
stop.model = Some("opus".to_string());
let rx = h.park_stop_wait(stop);
let forks = wake_forks(rx.recv_timeout(Duration::from_secs(10)).unwrap());
assert_eq!(forks[0].model.as_deref(), Some("sonnet"));
assert!(forks[0].model_fallbacks.is_empty());
let mut start = h.event(EventKind::SessionStart, "cc-fable");
start.model = Some("fable".to_string());
assert_ack(h.send_event(start));
let mut stop = h.event(EventKind::Stop, "cc-fable");
stop.model = Some("fable".to_string());
let rx = h.park_stop_wait(stop);
let forks = wake_forks(rx.recv_timeout(Duration::from_secs(10)).unwrap());
assert_eq!(forks[0].model.as_deref(), Some("sonnet"));
assert_eq!(forks[0].model_fallbacks, vec!["haiku".to_string()]);
let mut start = h.event(EventKind::SessionStart, "cc-other");
start.model = Some("some-future-model".to_string());
assert_ack(h.send_event(start));
let mut stop = h.event(EventKind::Stop, "cc-other");
stop.model = Some("some-future-model".to_string());
let rx = h.park_stop_wait(stop);
let forks = wake_forks(rx.recv_timeout(Duration::from_secs(10)).unwrap());
assert_eq!(forks[0].model.as_deref(), Some("haiku"));
}
#[test]
fn background_work_holds_the_idle_clock_until_it_finishes() {
let mut h = Harness::new("1s", "0").wake_grace_secs(0);
h.write_fork(
"goal.md",
"---\nfork: true\nrun_on:\n - idle: 0s\n---\nGOAL",
);
h.start_daemon();
h.write_transcript(100);
assert_ack(h.send_event(h.event_t(EventKind::SessionStart, "s1")));
h.append_background_launch("toolu_bg", "bg1");
let rx = h.park_stop_wait(h.event_t(EventKind::Stop, "s1"));
assert!(
rx.recv_timeout(Duration::from_millis(2500)).is_err(),
"an idle:0s fork fired while background work was still running"
);
h.append_completion_notification("toolu_bg", "completed");
assert_ack(h.send_event(h.prompt_submit("s1", false)));
let rx = h.park_stop_wait(h.event_t(EventKind::Stop, "s1"));
let payload = wake_payload(rx.recv_timeout(Duration::from_secs(10)).unwrap());
assert!(payload.contains("due: goal"), "{payload}");
}
#[test]
fn background_hold_expires_so_unfinished_work_cannot_silence_forks() {
let mut h = Harness::new("1s", "0").wake_grace_secs(0);
h.append_config("background_hold_timeout = \"2s\"");
h.write_fork(
"journal.md",
"---\nfork: true\nrun_on:\n - idle: 0s\n---\nJOURNAL",
);
h.start_daemon();
h.write_transcript(100);
assert_ack(h.send_event(h.event_t(EventKind::SessionStart, "s1")));
h.append_background_launch("toolu_server", "srv1");
let rx = h.park_stop_wait(h.event_t(EventKind::Stop, "s1"));
assert!(rx.recv_timeout(Duration::from_millis(1200)).is_err());
let payload = wake_payload(rx.recv_timeout(Duration::from_secs(10)).unwrap());
assert!(payload.contains("due: journal"), "{payload}");
}
#[test]
fn a_fork_can_opt_out_of_the_background_hold() {
let mut h = Harness::new("1s", "0").wake_grace_secs(0);
h.write_fork(
"goal.md",
"---\nfork: true\nrun_on:\n - idle: 0s\nbackground_hold: false\n---\nGOAL",
);
h.write_fork(
"journal.md",
"---\nfork: true\nrun_on:\n - idle: 0s\n---\nJOURNAL",
);
h.start_daemon();
h.write_transcript(100);
assert_ack(h.send_event(h.event_t(EventKind::SessionStart, "s1")));
h.append_monitor_launch("toolu_mon", "mon1");
let rx = h.park_stop_wait(h.event_t(EventKind::Stop, "s1"));
let payload = wake_payload(rx.recv_timeout(Duration::from_secs(10)).unwrap());
assert!(payload.contains("due: goal"), "{payload}");
assert!(
!payload.contains("due: journal"),
"a holding fork fired while a Monitor was still running: {payload}"
);
h.append_task_stop("mon1");
let rx = h.park_stop_wait(h.event_t(EventKind::Stop, "s1"));
let payload = wake_payload(rx.recv_timeout(Duration::from_secs(10)).unwrap());
assert!(payload.contains("due: journal"), "{payload}");
assert!(!payload.contains("due: goal"), "{payload}");
}
#[test]
fn a_fork_can_opt_in_to_the_hold_when_the_config_default_is_off() {
let mut h = Harness::new("1s", "0").wake_grace_secs(0);
h.append_config("background_hold = false");
h.write_fork(
"journal.md",
"---\nfork: true\nrun_on:\n - idle: 0s\nbackground_hold: true\n---\nJOURNAL",
);
h.start_daemon();
h.write_transcript(100);
assert_ack(h.send_event(h.event_t(EventKind::SessionStart, "s1")));
h.append_background_launch("toolu_bg", "bg1");
let rx = h.park_stop_wait(h.event_t(EventKind::Stop, "s1"));
assert!(
rx.recv_timeout(Duration::from_millis(2500)).is_err(),
"a fork with background_hold: true fired while background work was running"
);
h.append_completion_notification("toolu_bg", "completed");
assert_ack(h.send_event(h.prompt_submit("s1", false)));
let rx = h.park_stop_wait(h.event_t(EventKind::Stop, "s1"));
let payload = wake_payload(rx.recv_timeout(Duration::from_secs(10)).unwrap());
assert!(payload.contains("due: journal"), "{payload}");
}
#[test]
fn background_hold_can_be_switched_off() {
let mut h = Harness::new("1s", "0").wake_grace_secs(0);
h.append_config("background_hold = false");
h.write_fork(
"journal.md",
"---\nfork: true\nrun_on:\n - idle: 0s\n---\nJOURNAL",
);
h.start_daemon();
h.write_transcript(100);
assert_ack(h.send_event(h.event_t(EventKind::SessionStart, "s1")));
h.append_background_launch("toolu_bg", "bg1");
let rx = h.park_stop_wait(h.event_t(EventKind::Stop, "s1"));
let payload = wake_payload(rx.recv_timeout(Duration::from_secs(10)).unwrap());
assert!(payload.contains("due: journal"), "{payload}");
}
#[test]
fn own_fork_spawns_never_count_as_background_work() {
let mut h = Harness::new("1s", "0").wake_grace_secs(0);
h.write_fork(
"second.md",
"---\nfork: true\nrun_on:\n - idle: 0s\n---\nSECOND",
);
h.start_daemon();
h.write_transcript(100);
assert_ack(h.send_event(h.event_t(EventKind::SessionStart, "s1")));
h.append_fork_spawn("toolu_f1", "first");
h.append_transcript_line(
&serde_json::json!({
"type": "user",
"message": { "content": [
{ "type": "tool_result", "tool_use_id": "toolu_f1", "content": [
{ "type": "text", "text": "Async agent launched successfully.\nagentId: a1" },
] },
] }
})
.to_string(),
);
let rx = h.park_stop_wait(h.event_t(EventKind::Stop, "s1"));
let payload = wake_payload(rx.recv_timeout(Duration::from_secs(10)).unwrap());
assert!(payload.contains("due: second"), "{payload}");
}
fn dead_harness() -> autofork_core::harness::Harness {
let mut child = Command::new("true").spawn().unwrap();
let pid = child.id();
child.wait().unwrap();
autofork_core::harness::Harness {
pid,
start: None,
bin: None,
}
}
struct FakeClient(Child);
impl FakeClient {
fn spawn() -> Self {
FakeClient(
Command::new("sleep")
.arg("120")
.stdout(Stdio::null())
.stderr(Stdio::null())
.spawn()
.unwrap(),
)
}
fn harness(&self) -> autofork_core::harness::Harness {
autofork_core::harness::of_pid(self.0.id()).unwrap()
}
fn kill(&mut self) {
let _ = self.0.kill();
let _ = self.0.wait();
}
}
#[test]
fn a_dead_harness_closes_the_session_without_any_session_end() {
let mut h = Harness::new("30m", "0").liveness_sweep_secs(1);
let log = h.write_logging_hook("cleanup.md", "[session_end]");
h.start_daemon();
let mut ev = h.event(EventKind::SessionStart, "s-gone");
ev.harness = Some(dead_harness());
assert_ack(h.send_event(ev));
let lines = h.wait_for_hook_lines(&log, 1, Duration::from_secs(15));
assert_eq!(lines[0], "session_end||gone||s-gone");
}
#[test]
fn a_live_harness_survives_the_session_timeout_reaper() {
let mut h = Harness::new("30m", "0")
.liveness_sweep_secs(1)
.session_sweep_secs(1);
h.append_config("session_timeout = \"1s\"");
let log = h.write_logging_hook("cleanup.md", "[session_end]");
h.start_daemon();
let mut client = FakeClient::spawn();
let mut ev = h.event(EventKind::SessionStart, "s-live");
ev.harness = Some(client.harness());
assert_ack(h.send_event(ev));
std::thread::sleep(Duration::from_secs(4));
assert!(
std::fs::read_to_string(&log).unwrap_or_default().is_empty(),
"a live client's session must not be reaped"
);
client.kill();
}
#[test]
fn a_close_the_client_never_reported_still_flushes_its_idle_forks() {
let mut h = Harness::new("30m", "0").liveness_sweep_secs(1);
let record = h.recording_final_runner();
h.write_fork(
"handover.md",
"---\nfork: true\nrun_on:\n - idle: 30m\n---\nHAND OVER",
);
h.start_daemon();
let mut client = FakeClient::spawn();
let mut ev = h.event(EventKind::SessionStart, "s-flush");
ev.harness = Some(client.harness());
assert_ack(h.send_event(ev));
let _rx = h.park_stop_wait({
let mut ev = h.event(EventKind::Stop, "s-flush");
ev.harness = Some(client.harness());
ev
});
std::thread::sleep(Duration::from_millis(300));
client.kill();
let start = Instant::now();
let argv = loop {
let argv = std::fs::read_to_string(&record).unwrap_or_default();
if !argv.is_empty() {
break argv;
}
assert!(
start.elapsed() < Duration::from_secs(15),
"the end-runner never ran"
);
std::thread::sleep(Duration::from_millis(100));
};
assert!(argv.contains("final-run"), "{argv}");
assert!(argv.contains("--session s-flush"), "{argv}");
let specs_path = argv
.split_whitespace()
.skip_while(|a| *a != "--specs")
.nth(1)
.expect("the runner is handed a specs file");
let specs = std::fs::read_to_string(specs_path).unwrap();
assert!(specs.contains("handover"), "{specs}");
}
#[test]
fn the_end_runner_authenticates_as_the_session_not_as_the_daemon() {
let mut h = Harness::new("30m", "0")
.liveness_sweep_secs(1)
.daemon_env("CLAUDE_CODE_OAUTH_TOKEN", "daemon-stale-token")
.daemon_env("ANTHROPIC_API_KEY", "daemon-key");
let record = h.env_recording_final_runner();
h.write_fork(
"handover.md",
"---\nfork: true\nrun_on:\n - idle: 30m\n---\nHAND OVER",
);
h.start_daemon();
let session_env = autofork_core::runenv::Snapshot {
names: vec!["CLAUDE_CODE_OAUTH_TOKEN".into(), "ANTHROPIC_API_KEY".into()],
vars: vec![("CLAUDE_CODE_OAUTH_TOKEN".into(), "session-token".into())],
};
let mut client = FakeClient::spawn();
let mut ev = h.event(EventKind::SessionStart, "s-env");
ev.harness = Some(client.harness());
ev.env = Some(session_env.clone());
assert_ack(h.send_event(ev));
let _rx = h.park_stop_wait({
let mut ev = h.event(EventKind::Stop, "s-env");
ev.harness = Some(client.harness());
ev.env = Some(session_env);
ev
});
std::thread::sleep(Duration::from_millis(300));
client.kill();
let start = Instant::now();
let recorded = loop {
let line = std::fs::read_to_string(&record).unwrap_or_default();
if !line.is_empty() {
break line.trim().to_string();
}
assert!(
start.elapsed() < Duration::from_secs(15),
"the end-runner never ran"
);
std::thread::sleep(Duration::from_millis(100));
};
assert_eq!(recorded, "session-token|-", "end-runner credential env");
}
#[test]
fn a_lifecycle_hook_runs_with_the_session_credential_env() {
let mut h =
Harness::new("30m", "0").daemon_env("CLAUDE_CODE_OAUTH_TOKEN", "daemon-stale-token");
let log = h.project.join("cred.log");
let hook = h.project.join(".autofork/hooks/cred.md");
std::fs::create_dir_all(hook.parent().unwrap()).unwrap();
std::fs::write(
&hook,
format!(
"---\nhook: true\non: session_start\n\
command: printf '%s\\n' \"${{CLAUDE_CODE_OAUTH_TOKEN:--}}\" >> \"{}\"\n---\nx\n",
log.display()
),
)
.unwrap();
h.start_daemon();
let mut ev = h.event(EventKind::SessionStart, "s-hookenv");
ev.env = Some(autofork_core::runenv::Snapshot {
names: vec!["CLAUDE_CODE_OAUTH_TOKEN".into()],
vars: vec![("CLAUDE_CODE_OAUTH_TOKEN".into(), "session-token".into())],
});
assert_ack(h.send_event(ev));
let lines = h.wait_for_hook_lines(&log, 1, Duration::from_secs(10));
assert_eq!(lines[0], "session-token");
}
#[test]
fn a_late_session_end_hook_cannot_re_issue_the_batch_the_close_already_ran() {
let mut h = Harness::new("30m", "0").liveness_sweep_secs(1);
let record = h.recording_final_runner();
h.write_fork(
"handover.md",
"---\nfork: true\nrun_on:\n - idle: 30m\n---\nHAND OVER",
);
h.write_fork(
"journal.md",
"---\nfork: true\nrun_on:\n - idle: 30m\n---\nJOURNAL",
);
h.start_daemon();
let mut client = FakeClient::spawn();
let mut ev = h.event(EventKind::SessionStart, "s-double");
ev.harness = Some(client.harness());
assert_ack(h.send_event(ev));
let _rx = h.park_stop_wait({
let mut ev = h.event(EventKind::Stop, "s-double");
ev.harness = Some(client.harness());
ev
});
std::thread::sleep(Duration::from_millis(300));
client.kill();
let start = Instant::now();
loop {
if !std::fs::read_to_string(&record)
.unwrap_or_default()
.is_empty()
{
break;
}
assert!(
start.elapsed() < Duration::from_secs(15),
"the end-runner never ran"
);
std::thread::sleep(Duration::from_millis(100));
}
let ResponseBody::Due { forks } = h.request(RequestBody::TakeFinalRuns {
session_id: "s-double".into(),
}) else {
panic!("expected Due");
};
assert!(
forks.is_empty(),
"the batch was already claimed and is running: {forks:?}"
);
}
#[test]
fn an_orphaned_stop_wait_cannot_resurrect_a_dead_session() {
let mut h = Harness::new("30m", "0");
let log = h.write_logging_hook("cleanup.md", "[session_end]");
h.start_daemon();
assert_ack(h.send_event(h.event(EventKind::SessionStart, "s-orphan")));
let mut ev = h.event(EventKind::Stop, "s-orphan");
ev.harness = Some(dead_harness());
let rx = h.park_stop_wait(ev);
let resp = rx.recv_timeout(Duration::from_secs(5)).unwrap();
assert!(
matches!(resp, ResponseBody::Waited),
"expected Waited, got {resp:?}"
);
let lines = h.wait_for_hook_lines(&log, 1, Duration::from_secs(5));
assert_eq!(lines[0], "session_end||gone||s-orphan");
}
#[test]
fn a_session_inherited_dead_from_a_previous_daemon_closes_without_flushing() {
let mut h = Harness::new("30m", "0").liveness_sweep_secs(1);
let record = h.recording_final_runner();
let log = h.write_logging_hook("cleanup.md", "[session_end]");
h.write_fork(
"handover.md",
"---\nfork: true\nrun_on:\n - idle: 30m\n---\nHAND OVER",
);
h.start_daemon();
let mut client = FakeClient::spawn();
let mut ev = h.event(EventKind::SessionStart, "s-inherited");
ev.harness = Some(client.harness());
assert_ack(h.send_event(ev));
let _rx = h.park_stop_wait({
let mut ev = h.event(EventKind::Stop, "s-inherited");
ev.harness = Some(client.harness());
ev
});
std::thread::sleep(Duration::from_millis(300));
h.kill_daemon();
client.kill();
std::thread::sleep(Duration::from_secs(2));
h.start_daemon();
let lines = h.wait_for_hook_lines(&log, 1, Duration::from_secs(15));
assert_eq!(lines[0], "session_end||gone||s-inherited");
assert!(
!record.exists(),
"an inherited dead session must not flush: {:?}",
std::fs::read_to_string(&record)
);
}
#[test]
fn a_feed_spools_its_stdout_as_a_context_block() {
let mut h = Harness::new("30m", "0");
h.write_feed("brief.md", "[activity]", "context", "RECENT: one, two");
h.start_daemon();
assert_ack(h.send_event(h.event(EventKind::SessionStart, "s1")));
assert_ack(h.send_event(h.prompt_submit("s1", true)));
let blocks = h.wait_for_reports("s1", Duration::from_secs(10));
assert_eq!(blocks.len(), 1, "{blocks:?}");
assert!(blocks[0].contains("source: autofork"), "{}", blocks[0]);
assert!(
blocks[0].contains("feed: brief (activity)"),
"{}",
blocks[0]
);
assert!(blocks[0].contains("RECENT: one, two"), "{}", blocks[0]);
}
#[test]
fn a_drain_that_asks_waits_for_a_context_feed_still_running() {
let mut h = Harness::new("30m", "0");
let path = h.project.join(".autofork/hooks/slow.md");
std::fs::create_dir_all(path.parent().unwrap()).unwrap();
std::fs::write(
&path,
"---\nhook: true\non: [session_start]\ndeliver: context\n\
command: |-\n sleep 0.5; printf '%s' 'LATE'\n---\nfeed\n",
)
.unwrap();
h.start_daemon();
let mut ev = h.event(EventKind::SessionStart, "oc1");
ev.client = Some("opencode".into());
assert_ack(h.send_event(ev));
match h.request(RequestBody::TakeReports {
session_id: "oc1".into(),
wait_ms: None,
}) {
ResponseBody::Reports { blocks } => assert!(blocks.is_empty(), "{blocks:?}"),
other => panic!("unexpected {other:?}"),
}
let ResponseBody::Reports { blocks } = h.request(RequestBody::TakeReports {
session_id: "oc1".into(),
wait_ms: Some(5_000),
}) else {
panic!("expected Reports");
};
assert_eq!(blocks.len(), 1, "{blocks:?}");
assert!(blocks[0].contains("LATE"), "{}", blocks[0]);
}
#[test]
fn a_drain_that_asks_gives_up_on_a_feed_that_outlasts_its_budget() {
let mut h = Harness::new("30m", "0");
let path = h.project.join(".autofork/hooks/wedged.md");
std::fs::create_dir_all(path.parent().unwrap()).unwrap();
std::fs::write(
&path,
"---\nhook: true\non: [session_start]\ndeliver: context\n\
command: |-\n sleep 30; printf '%s' 'NEVER'\n---\nfeed\n",
)
.unwrap();
h.start_daemon();
let mut ev = h.event(EventKind::SessionStart, "oc2");
ev.client = Some("opencode".into());
assert_ack(h.send_event(ev));
let started = Instant::now();
match h.request(RequestBody::TakeReports {
session_id: "oc2".into(),
wait_ms: Some(300),
}) {
ResponseBody::Reports { blocks } => assert!(blocks.is_empty(), "{blocks:?}"),
other => panic!("unexpected {other:?}"),
}
assert!(
started.elapsed() < Duration::from_secs(5),
"the drain waited past its budget: {:?}",
started.elapsed()
);
}
#[test]
fn a_feed_does_not_deliver_the_same_block_twice() {
let mut h = Harness::new("30m", "0");
h.write_feed("brief.md", "[activity]", "context", "UNCHANGED");
h.start_daemon();
assert_ack(h.send_event(h.event(EventKind::SessionStart, "s1")));
assert_ack(h.send_event(h.prompt_submit("s1", true)));
let first = h.wait_for_reports("s1", Duration::from_secs(10));
assert_eq!(first.len(), 1);
assert_ack(h.send_event(h.prompt_submit("s1", true)));
std::thread::sleep(Duration::from_secs(2));
match h.request(RequestBody::TakeReports {
session_id: "s1".into(),
wait_ms: None,
}) {
ResponseBody::Reports { blocks } => {
assert!(
blocks.is_empty(),
"unchanged output was re-delivered: {blocks:?}"
)
}
other => panic!("unexpected {other:?}"),
}
}
#[test]
fn a_wake_feed_resolves_a_parked_poll_with_its_blocks() {
let mut h = Harness::new("30m", "0");
h.write_feed("urgent.md", "[activity]", "wake", "SOMETHING MOVED");
h.start_daemon();
assert_ack(h.send_event(h.event(EventKind::SessionStart, "s1")));
assert_ack(h.send_event(h.prompt_submit("s1", true)));
std::thread::sleep(Duration::from_secs(1));
let rx = h.park_stop_wait(h.event(EventKind::Stop, "s1"));
match rx.recv_timeout(Duration::from_secs(10)).unwrap() {
ResponseBody::Wake {
payload,
forks,
feed,
} => {
let feed = feed.expect("a feed wake carries its blocks");
assert!(feed.wake, "deliver: wake must ask for a turn");
assert_eq!(feed.blocks.len(), 1);
assert!(feed.blocks[0].contains("SOMETHING MOVED"));
assert!(payload.contains("SOMETHING MOVED"), "{payload}");
assert!(forks.is_none() || forks.as_ref().unwrap().is_empty());
}
other => panic!("expected a feed wake, got {other:?}"),
}
}
#[test]
fn a_watched_path_changing_fires_a_hook_with_the_paths() {
let mut h = Harness::new("30m", "0");
h.append_config("watch_interval = 1");
h.append_config("watch_debounce = 0");
let watched = h.watched_dir();
let log = h.project.join("changed.log");
let hook = h.project.join(".autofork/hooks/notes.md");
std::fs::create_dir_all(hook.parent().unwrap()).unwrap();
std::fs::write(
&hook,
format!(
"---\nhook: true\non:\n - \"changed: {}/**/*.md\"\n\
command: printf '%s|%s\\n' \"$AUTOFORK_EVENT\" \"$AUTOFORK_CHANGED_PATHS\" >> \"{}\"\n---\nwatcher\n",
watched.display(),
log.display()
),
)
.unwrap();
h.start_daemon();
assert_ack(h.send_event(h.event(EventKind::SessionStart, "s1")));
std::thread::sleep(Duration::from_secs(2));
std::fs::write(watched.join("new.md"), "hello").unwrap();
let lines = h.wait_for_hook_lines(&log, 1, Duration::from_secs(20));
assert!(lines[0].starts_with("changed|"), "{lines:?}");
assert!(lines[0].contains("new.md"), "{lines:?}");
}
#[test]
fn a_watched_path_changing_wakes_a_fork_once() {
let mut h = Harness::new("30m", "0");
h.append_config("watch_interval = 1");
h.append_config("watch_debounce = 0");
let watched = h.watched_dir();
h.write_fork(
"review.md",
&format!(
"---\nfork: true\nrun_on:\n - \"changed: {}/**/*.rs\"\n---\nReview what changed.\n",
watched.display()
),
);
h.start_daemon();
assert_ack(h.send_event(h.event(EventKind::SessionStart, "s1")));
std::thread::sleep(Duration::from_secs(2));
let rx = h.park_stop_wait(h.event(EventKind::Stop, "s1"));
std::fs::write(watched.join("lib.rs"), "fn main() {}").unwrap();
match rx.recv_timeout(Duration::from_secs(20)).unwrap() {
ResponseBody::Wake { payload, forks, .. } => {
assert!(payload.contains("review"), "{payload}");
let forks = forks.expect("structured specs");
assert_eq!(forks.len(), 1);
assert!(forks[0].trigger.starts_with("changed:"), "{:?}", forks[0]);
assert!(forks[0].prompt.contains("lib.rs"), "{}", forks[0].prompt);
}
other => panic!("expected a wake, got {other:?}"),
}
let rx2 = h.park_stop_wait(h.event(EventKind::Stop, "s1"));
match rx2.recv_timeout(Duration::from_secs(4)) {
Err(mpsc::RecvTimeoutError::Timeout) => {}
Ok(other) => panic!("the consumed trigger re-fired: {other:?}"),
Err(e) => panic!("{e:?}"),
}
}
#[test]
fn emit_triggers_a_listening_fork_and_hook() {
let mut h = Harness::new("30m", "0");
let log = h.write_logging_hook("announce.md", "[\"event: deploy\"]");
h.write_fork(
"on-deploy.md",
"---\nfork: true\nrun_on:\n - \"event: deploy\"\n---\nCheck the deploy.\n",
);
h.start_daemon();
assert_ack(h.send_event(h.event(EventKind::SessionStart, "s1")));
let rx = h.park_stop_wait(h.event(EventKind::Stop, "s1"));
std::thread::sleep(Duration::from_millis(300));
match h.request(RequestBody::Emit {
name: "deploy".into(),
payload: Some("build 412 is live".into()),
project_root: None,
session_id: None,
}) {
ResponseBody::Emitted { sessions } => assert_eq!(sessions, 1),
other => panic!("unexpected {other:?}"),
}
match rx.recv_timeout(Duration::from_secs(10)).unwrap() {
ResponseBody::Wake { forks, .. } => {
let forks = forks.expect("structured specs");
assert_eq!(forks.len(), 1);
assert_eq!(forks[0].name, "on-deploy");
assert_eq!(forks[0].trigger, "event:deploy");
assert!(
forks[0].prompt.contains("build 412 is live"),
"the emit payload must reach the fork: {}",
forks[0].prompt
);
}
other => panic!("expected a wake, got {other:?}"),
}
let lines = h.wait_for_hook_lines(&log, 1, Duration::from_secs(10));
assert!(lines[0].starts_with("event|"), "{lines:?}");
}
#[test]
fn emit_can_be_scoped_to_one_session() {
let mut h = Harness::new("30m", "0");
h.start_daemon();
assert_ack(h.send_event(h.event(EventKind::SessionStart, "s1")));
assert_ack(h.send_event(h.event(EventKind::SessionStart, "s2")));
match h.request(RequestBody::Emit {
name: "ping".into(),
payload: None,
project_root: None,
session_id: Some("s2".into()),
}) {
ResponseBody::Emitted { sessions } => assert_eq!(sessions, 1),
other => panic!("unexpected {other:?}"),
}
}
#[test]
fn a_run_finishing_after_the_user_spoke_is_stale() {
let mut h = Harness::new("1s", "0").wake_grace_secs(0);
h.write_fork(
"goal.md",
"---\nfork: true\nrun_on: [idle: 0s]\nchain: true\ngate: true\n---\nGOAL",
);
h.start_daemon();
assert_ack(h.send_event(oc_event(&h, EventKind::SessionStart, "oc1")));
let rx = h.park_stop_wait(oc_event(&h, EventKind::Stop, "oc1"));
let forks = wake_forks(rx.recv_timeout(Duration::from_secs(10)).unwrap());
assert_eq!(forks[0].name, "goal");
assert_ack(h.request(RequestBody::ForkSpawned {
session_id: "oc1".into(),
fork: "goal".into(),
run_ref: "ses_g1".into(),
}));
assert!(matches!(
h.request(RequestBody::RunState {
session_id: "oc1".into(),
run_ref: "ses_g1".into(),
}),
ResponseBody::RunState { stale: false }
));
assert_ack(h.send_event(h.prompt_submit("oc1", true)));
assert!(matches!(
h.request(RequestBody::RunState {
session_id: "oc1".into(),
run_ref: "ses_g1".into(),
}),
ResponseBody::RunState { stale: true }
));
let parked = h.park_stop_wait(oc_event(&h, EventKind::Stop, "oc1"));
assert!(
parked.recv_timeout(Duration::from_millis(1500)).is_err(),
"an in-flight run must hold the overlap gate"
);
assert_ack(h.request(RequestBody::ForkCompleted {
session_id: "oc1".into(),
fork: "goal".into(),
run_ref: "ses_g1".into(),
status: "completed".into(),
cont: Some(true),
}));
let forks = wake_forks(parked.recv_timeout(Duration::from_secs(10)).unwrap());
assert_eq!(
forks[0].name, "goal",
"the next goal run supersedes the stale one"
);
}
#[test]
fn a_stale_completion_leaves_the_new_pause_gate_alone() {
let mut h = Harness::new("1s", "0").wake_grace_secs(0);
h.write_fork(
"goal.md",
"---\nfork: true\nrun_on: [idle: 0s]\nchain: true\ngate: true\noverlap: true\n---\nGOAL",
);
h.write_fork("handover.md", "---\nfork: true\nrun_on: [idle: 1s]\n---\nH");
h.start_daemon();
assert_ack(h.send_event(oc_event(&h, EventKind::SessionStart, "oc1")));
let rx = h.park_stop_wait(oc_event(&h, EventKind::Stop, "oc1"));
assert_eq!(
wake_forks(rx.recv_timeout(Duration::from_secs(10)).unwrap())[0].name,
"goal"
);
assert_ack(h.request(RequestBody::ForkSpawned {
session_id: "oc1".into(),
fork: "goal".into(),
run_ref: "ses_g1".into(),
}));
assert_ack(h.send_event(h.prompt_submit("oc1", true)));
let rx2 = h.park_stop_wait(oc_event(&h, EventKind::Stop, "oc1"));
let forks = wake_forks(rx2.recv_timeout(Duration::from_secs(10)).unwrap());
assert_eq!(forks.len(), 1, "{forks:?}");
assert_eq!(forks[0].name, "goal");
assert_ack(h.request(RequestBody::ForkSpawned {
session_id: "oc1".into(),
fork: "goal".into(),
run_ref: "ses_g2".into(),
}));
assert_ack(h.request(RequestBody::ForkCompleted {
session_id: "oc1".into(),
fork: "goal".into(),
run_ref: "ses_g1".into(),
status: "completed".into(),
cont: None,
}));
let rx3 = h.park_stop_wait(oc_event(&h, EventKind::Stop, "oc1"));
assert!(
rx3.recv_timeout(Duration::from_millis(2500)).is_err(),
"a stale settlement must not release the gate of the run that holds it"
);
assert_ack(h.request(RequestBody::ForkCompleted {
session_id: "oc1".into(),
fork: "goal".into(),
run_ref: "ses_g2".into(),
status: "completed".into(),
cont: None,
}));
assert!(matches!(
rx3.recv_timeout(Duration::from_secs(5)).unwrap(),
ResponseBody::Waited
));
let rx4 = h.park_stop_wait(oc_event(&h, EventKind::Stop, "oc1"));
let forks = wake_forks(rx4.recv_timeout(Duration::from_secs(10)).unwrap());
assert_eq!(forks[0].name, "handover");
}