use super::exit_launch::PROBE_RETRIES;
use super::fixtures::*;
use crate::config::models::Capabilities;
use crate::config::{Model, Workflow};
use crate::prompt::Error;
use crate::prompt::dispatch::advance::{AdvanceOutcome, run};
use crate::prompt::inbox::{self, Launcher, inbox_dir, try_acquire};
use crate::prompt::resolve::WorkerConfig;
use brazen::{Content, FinishReason};
use std::cell::RefCell;
use std::io;
use std::path::{Path, PathBuf};
use tempfile::TempDir;
pub(super) const AGENT: &str = "20260101-a1";
#[derive(Default)]
pub(super) struct RecLauncher {
pub(super) invocations: RefCell<Vec<String>>,
}
impl Launcher for RecLauncher {
fn launch(&self, _ws: &Path, agent: &str) -> io::Result<()> {
self.invocations.borrow_mut().push(agent.to_string());
Ok(())
}
}
pub(super) fn worker_config() -> WorkerConfig {
WorkerConfig {
role: "worker".into(),
model: Model {
provider: "anthropic".into(),
model_id: "claude-sonnet-5".into(),
capabilities: Capabilities(vec![]),
context_window: 200_000,
},
provider_row: "anthropic".into(),
tools: vec![],
soul: "be helpful".into(),
binary: "bz".into(),
workflow: Workflow::parse("events: {}\n", std::path::Path::new("workflow.yaml")).unwrap(),
manifest: None,
expect_handshake: false,
}
}
pub(super) fn workspace_with_tail(entries: &[(&str, String)]) -> (TempDir, PathBuf) {
let ws = TempDir::new().unwrap();
let wt = crate::workspace::agent_worktree(ws.path(), AGENT);
std::fs::create_dir_all(wt.join("messages")).unwrap();
std::fs::write(wt.join("goal.md"), "the goal").unwrap();
for (name, body) in entries {
std::fs::write(wt.join("messages").join(name), body).unwrap();
}
(ws, wt)
}
pub(super) fn model_entry(blocks: &[Content]) -> String {
serde_json::to_string(blocks).unwrap()
}
pub(super) fn terminal_tail() -> Vec<(&'static str, String)> {
vec![
("001-user.md", "hi".to_string()),
(
"002-claude-sonnet-5.json",
model_entry(&[Content::Text("final".into())]),
),
]
}
pub(super) fn no_resolve() -> Result<WorkerConfig, Error> {
panic!("resolve must not run on a no-op hop")
}
#[test]
#[should_panic(expected = "resolve must not run")]
fn no_resolve_is_a_tripwire() {
let _ = no_resolve();
}
pub(super) fn eventually_free(ws: &Path, agent: &str) -> bool {
free_within(ws, agent, PROBE_RETRIES)
}
pub(super) fn free_within(ws: &Path, agent: &str, retries: u32) -> bool {
for attempt in 0..retries {
if attempt > 0 {
std::thread::sleep(std::time::Duration::from_millis(5));
}
if try_acquire(&inbox_dir(ws, agent)).unwrap().is_some() {
return true;
}
}
false
}
#[test]
fn free_within_gives_up_on_a_genuinely_held_lease() {
let ws = TempDir::new().unwrap();
let _held = try_acquire(&inbox_dir(ws.path(), AGENT)).unwrap().unwrap();
assert!(!free_within(ws.path(), AGENT, 2));
}
#[test]
fn already_driven_is_a_clean_noop_without_resolving() {
let (ws, _wt) = workspace_with_tail(&terminal_tail());
let _held = try_acquire(&inbox_dir(ws.path(), AGENT)).unwrap().unwrap();
let (adapter, sleeper, git) = (unreachable_adapter(), StubSleeper::default(), StubGit::ok());
let (clock, id) = (FixedClock::default(), FixedIdGen);
let tools = StubToolExecutor::ok();
let deps = valid_deps(&adapter, &sleeper, &git, &clock, &id, &tools, ws.path());
let out = run(ws.path(), AGENT, None, &deps, &mut no_resolve).unwrap();
assert!(matches!(out, AdvanceOutcome::AlreadyDriven));
}
#[test]
fn empty_workspace_is_nothing_to_do() {
let ws = TempDir::new().unwrap();
let (adapter, sleeper, git) = (unreachable_adapter(), StubSleeper::default(), StubGit::ok());
let (clock, id) = (FixedClock::default(), FixedIdGen);
let tools = StubToolExecutor::ok();
let deps = valid_deps(&adapter, &sleeper, &git, &clock, &id, &tools, ws.path());
let out = run(ws.path(), AGENT, None, &deps, &mut no_resolve).unwrap();
assert!(matches!(out, AdvanceOutcome::NothingToDo));
}
#[test]
fn terminal_tail_with_empty_inbox_is_the_pin_1_silent_exit() {
let (ws, wt) = workspace_with_tail(&terminal_tail());
let (adapter, sleeper, git) = (unreachable_adapter(), StubSleeper::default(), StubGit::ok());
let (clock, id) = (FixedClock::default(), FixedIdGen);
let tools = StubToolExecutor::ok();
let deps = valid_deps(&adapter, &sleeper, &git, &clock, &id, &tools, ws.path());
let out = run(ws.path(), AGENT, None, &deps, &mut no_resolve).unwrap();
assert!(matches!(out, AdvanceOutcome::NothingToDo));
assert!(!ws.path().join("steps").exists());
assert_eq!(std::fs::read_dir(wt.join("messages")).unwrap().count(), 2);
}
#[test]
fn a_deposit_steps_the_branch_to_a_new_final_response() {
let (ws, wt) = workspace_with_tail(&terminal_tail());
let (clock, id) = (FixedClock::default(), FixedIdGen);
inbox::deposit(ws.path(), AGENT, "user", "again", &clock).unwrap();
let adapter = StubAdapter::scripted([StubAdapter::reply_ok(&happy_response_bytes())]);
let (sleeper, git) = (StubSleeper::default(), StubGit::ok());
let tools = StubToolExecutor::ok();
let rec = RecLauncher::default();
let mut deps = valid_deps(&adapter, &sleeper, &git, &clock, &id, &tools, ws.path());
deps.launcher = &rec;
let out = run(ws.path(), AGENT, None, &deps, &mut || Ok(worker_config())).unwrap();
assert!(matches!(out, AdvanceOutcome::Terminal));
let delivered = std::fs::read_to_string(wt.join("messages/003-user.md")).unwrap();
assert!(delivered.contains("again"), "got {delivered:?}");
assert!(wt.join("messages/004-claude-sonnet-5.json").exists());
assert!(
ws.path()
.join(format!("steps/{AGENT}/001/response.json"))
.exists()
);
assert_eq!(*rec.invocations.borrow(), vec![AGENT.to_string()]);
assert!(eventually_free(ws.path(), AGENT));
}
#[test]
fn a_tool_use_step_hands_off_as_tools_pending_with_the_lease_held() {
let (ws, wt) = workspace_with_tail(&terminal_tail());
let (clock, id) = (FixedClock::default(), FixedIdGen);
inbox::deposit(ws.path(), AGENT, "user", "run it", &clock).unwrap();
let tool_stream = stream_of(
FinishReason::ToolUse,
&[Block::ToolUse {
id: "t1",
name: "bash",
input: serde_json::json!({"command": "true"}),
}],
);
let adapter = StubAdapter::scripted([StubAdapter::reply_ok(&tool_stream)]);
let (sleeper, git) = (StubSleeper::default(), StubGit::ok());
let tools = StubToolExecutor::ok();
let rec = RecLauncher::default();
let mut deps = valid_deps(&adapter, &sleeper, &git, &clock, &id, &tools, ws.path());
deps.launcher = &rec;
let out = run(ws.path(), AGENT, None, &deps, &mut || Ok(worker_config())).unwrap();
let AdvanceOutcome::ToolsPending(lease) = out else {
panic!("expected ToolsPending");
};
assert_eq!(tools.invocations.borrow().len(), 1);
assert!(wt.join("messages/005-tool.json").exists());
assert!(try_acquire(&inbox_dir(ws.path(), AGENT)).unwrap().is_none());
drop(lease);
assert!(eventually_free(ws.path(), AGENT));
assert!(rec.invocations.borrow().is_empty());
}