use super::poll;
use super::stop_common::{
HAPPY_SSE, lernie_bin, poll_for_conv_branch_with_diag, poll_for_path, reap, scaffold_repo,
spawn_prompt, write_brazen_config, write_global_models,
};
use crate::prompt::inbox::inbox_dir;
use crate::prompt::stop::{PgidFinder, ProcFsFinder};
use httpmock::Method::POST;
use httpmock::MockServer;
use std::fs;
use std::path::{Path, PathBuf};
use std::process::{Child, Command, Stdio};
use std::time::Duration;
use tempfile::TempDir;
struct Family {
_server: MockServer,
_tmp: TempDir,
dest: PathBuf,
parent: String,
child: String,
prompt: Child,
child_pgid: i32,
}
fn step_response(dest: &Path, agent: &str) -> PathBuf {
dest.join("steps")
.join(agent)
.join("001")
.join("response.json")
}
fn holder_pgid(dest: &Path, agent: &str) -> Option<i32> {
ProcFsFinder::default()
.find_holder_pgid(&inbox_dir(dest, agent))
.expect("scan /proc")
}
fn poll_until_no_holder(dest: &Path, agent: &str) {
let gone = poll::until(dest, || holder_pgid(dest, agent).is_none().then_some(()));
assert!(
gone.is_some(),
"executor for {agent} (pgid {:?}) outlived the stop, with {} untouched for {:?}",
holder_pgid(dest, agent),
dest.display(),
poll::patience()
);
}
fn run_stop(dest: &Path, agent: &str, stop_children: bool) {
let mut cmd = Command::new(lernie_bin());
cmd.arg("stop").arg(dest).arg(agent);
if stop_children {
cmd.arg("--stop-children");
}
let out = cmd
.stdout(Stdio::piped())
.stderr(Stdio::piped())
.output()
.expect("spawn lernie stop");
assert!(
out.status.success(),
"lernie stop: {}",
String::from_utf8_lossy(&out.stderr)
);
}
fn assert_no_terminal_end(path: &Path) {
let bytes = fs::read(path).expect("stopped response.json is readable");
let Some(last) = bytes.split(|b| *b == b'\n').rfind(|l| !l.is_empty()) else {
return; };
let v: serde_json::Value = serde_json::from_slice(last).expect("trailing line is JSON");
assert_ne!(
v["type"].as_str(),
Some("end"),
"a stopped step's response.json carries no terminal `end`; last: {v}"
);
}
fn live_family() -> Family {
let server = MockServer::start();
server.mock(|when, then| {
when.method(POST).path("/v1/messages");
then.status(200)
.header("content-type", "text/event-stream")
.delay(Duration::from_secs(120))
.body(HAPPY_SSE);
});
let tmp = TempDir::new().unwrap();
let harness = tmp.path().join("harness");
fs::create_dir_all(&harness).unwrap();
write_global_models(&harness);
let brazen_config = write_brazen_config(tmp.path(), &server.base_url());
let dest = tmp.path().join("conv");
scaffold_repo(&dest, &harness);
let mut prompt = spawn_prompt(&dest, &harness, &brazen_config, "ping");
let parent = poll_for_conv_branch_with_diag(&dest, &mut prompt);
poll_for_path(&dest, &step_response(&dest, &parent));
let out = Command::new(lernie_bin())
.args(["dispatch", "worker"])
.arg(&dest)
.arg(&parent)
.arg("--goal")
.arg("hold a model call open")
.env("LERNIE_HOME", &harness)
.env("BRAZEN_CONFIG", &brazen_config)
.output()
.expect("spawn lernie dispatch");
assert!(
out.status.success(),
"lernie dispatch worker: {}",
String::from_utf8_lossy(&out.stderr)
);
let child = String::from_utf8(out.stdout).unwrap().trim().to_owned();
assert!(
child.starts_with(&format!("{parent}-")),
"hyphenated descent (§2.3): {child} must extend {parent}"
);
poll_for_path(&dest, &step_response(&dest, &child));
let parent_pgid = holder_pgid(&dest, &parent).expect("parent executor holds its inbox lock");
let child_pgid = holder_pgid(&dest, &child).expect("child executor holds its inbox lock");
assert_ne!(
parent_pgid, child_pgid,
"every executor takes its own process group, child alike (§2.9)"
);
Family {
_server: server,
_tmp: tmp,
dest,
parent,
child,
prompt,
child_pgid,
}
}
#[test]
fn stop_children_fells_the_live_child_executor() {
let mut fam = live_family();
run_stop(&fam.dest, &fam.parent, true);
let status = reap(&fam.dest, &mut fam.prompt);
assert!(
status.success(),
"the stopped parent must exit cleanly (§2.9 step 3), got {status:?}"
);
poll_until_no_holder(&fam.dest, &fam.child);
assert_no_terminal_end(&step_response(&fam.dest, &fam.child));
let deposited = inbox_dir(&fam.dest, &fam.parent).join(format!("{}-001.md", fam.child));
poll_for_path(&fam.dest, &deposited);
let body = fs::read_to_string(&deposited).unwrap();
assert!(
body.contains("epitaph: stopped") && body.contains(&format!("from: {}", fam.child)),
"the felled child deposits its stopped result (§2.9 step 3): {body}"
);
}
#[test]
fn bare_stop_leaves_the_live_child_running() {
let mut fam = live_family();
run_stop(&fam.dest, &fam.parent, false);
let status = reap(&fam.dest, &mut fam.prompt);
assert!(
status.success(),
"the stopped parent must exit cleanly (§2.9 step 3), got {status:?}"
);
poll_until_no_holder(&fam.dest, &fam.parent);
assert_eq!(
holder_pgid(&fam.dest, &fam.child),
Some(fam.child_pgid),
"a bare stop must not cross into the child's group (§2.9)"
);
run_stop(&fam.dest, &fam.child, false);
poll_until_no_holder(&fam.dest, &fam.child);
}