use super::advance::{RecLauncher, worker_config};
use super::fixtures::*;
use super::parent_revival::DescentClock;
use crate::prompt::dispatch::advance::run;
use crate::prompt::inbox::{self, Launcher, inbox_dir};
use crate::prompt::{Clock, PinnedDocs};
use crate::template::RealGit;
use crate::workspace::agent_name::mint::test_rng;
use crate::workspace::fixture;
use std::cell::RefCell;
use std::path::{Path, PathBuf};
use std::{io, sync::atomic::AtomicUsize, sync::atomic::Ordering};
const HOP_CAP: usize = 12;
struct SecondClock;
impl Clock for SecondClock {
fn now_iso8601(&self) -> String {
"iso".into()
}
fn now_compact(&self) -> String {
"ct2".into()
}
}
struct ExchangeLauncher {
hops: AtomicUsize,
invocations: RefCell<Vec<String>>,
}
impl ExchangeLauncher {
fn new() -> Self {
Self {
hops: AtomicUsize::new(0),
invocations: RefCell::new(Vec::new()),
}
}
}
impl Launcher for ExchangeLauncher {
fn launch(&self, ws: &Path, agent: &str) -> io::Result<()> {
self.invocations.borrow_mut().push(agent.to_string());
if self.hops.fetch_add(1, Ordering::SeqCst) < HOP_CAP {
advance_with(ws, agent, self);
}
Ok(())
}
}
fn advance_with(ws: &Path, agent: &str, launcher: &dyn Launcher) {
let adapter = StubAdapter::scripted([StubAdapter::reply_ok(&happy_response_bytes())]);
let (sleeper, tools, stub_git) = (
StubSleeper::default(),
StubToolExecutor::ok(),
StubGit::ok(),
);
let (clock, id) = (FixedClock::default(), FixedIdGen);
let git = RealGit::new();
let mut deps = valid_deps(&adapter, &sleeper, &stub_git, &clock, &id, &tools, ws);
deps.git = &git;
deps.launcher = launcher;
run(ws, agent, None, &deps, &mut || Ok(worker_config())).unwrap();
}
fn two_siblings() -> (tempfile::TempDir, PathBuf, String, String) {
use crate::prompt::child_dispatch::{ChildDispatchRequest, run as dispatch_child};
let (holder, ws) = fixture::workspace();
let parent = "20260101-a1";
let parent_wt = fixture::spawn_root(&ws, parent);
let child = |clock: &dyn Clock| {
dispatch_child(
&ChildDispatchRequest {
repo: &ws,
parent_branch: parent,
parent_worktree: &parent_wt,
role: "worker",
goal: "cooperate",
name: None,
fork_point: None,
cwd: None,
pins: PinnedDocs::none(),
},
&RealGit::new(),
clock,
&FixedIdGen,
no_launch(),
test_rng(),
)
.unwrap()
};
let speccer = child(&DescentClock);
let builder = child(&SecondClock);
let inert = RecLauncher::default();
advance_with(&ws, &speccer, &inert);
advance_with(&ws, &builder, &inert);
(holder, ws, speccer, builder)
}
fn pending_results(ws: &Path, agent: &str) -> Vec<String> {
std::fs::read_dir(inbox_dir(ws, agent))
.into_iter()
.flatten()
.flatten()
.map(|e| std::fs::read_to_string(e.path()).unwrap())
.filter(|b| b.contains("terminal_ref:"))
.collect()
}
#[test]
fn one_exchange_between_two_agents_terminates() {
let (_holder, ws, speccer, builder) = two_siblings();
inbox::deposit(
&ws,
&builder,
&speccer,
"the spec",
&DescentClock,
&RealGit::new(),
)
.unwrap();
inbox::deposit(
&ws,
&speccer,
&builder,
"DONE",
&DescentClock,
&RealGit::new(),
)
.unwrap();
let launcher = ExchangeLauncher::new();
advance_with(&ws, &builder, &launcher);
let hops = launcher.hops.load(Ordering::SeqCst);
assert!(
hops < HOP_CAP,
"the exchange did not terminate — {hops} hops: {:?}",
launcher.invocations.borrow()
);
assert!(
pending_results(&ws, &speccer).is_empty() && pending_results(&ws, &builder).is_empty(),
"an exchange that settled leaves no reply behind: {:?} / {:?}",
pending_results(&ws, &speccer),
pending_results(&ws, &builder)
);
}