use super::*;
use crate::shim::protocol::{Command, ShimState, socketpair};
use crate::team::config::{AllocationPolicy, AllocationStrategy, BoardConfig, WorkflowPolicy};
use crate::team::daemon::agent_handle::AgentHandle;
use crate::team::events;
use crate::team::events::{QualityMetricsInfo, TeamEvent};
use crate::team::inbox;
use crate::team::standup::MemberState;
use crate::team::task_loop::{
current_worktree_branch, engineer_base_branch_name, setup_engineer_worktree,
setup_multi_repo_worktree_from_trunk,
};
use crate::team::team_events_path;
use crate::team::telemetry_db;
use crate::team::test_support::{
TestDaemonBuilder, engineer_member, git, git_ok, git_stdout, init_git_repo, manager_member,
write_open_task_file, write_owned_task_file,
};
use std::collections::HashMap;
use std::time::{Duration, Instant};
fn write_task_file(project_root: &std::path::Path, file_name: &str, content: &str) {
let tasks_dir = project_root
.join(".batty")
.join("team_config")
.join("board")
.join("tasks");
std::fs::create_dir_all(&tasks_dir).unwrap();
std::fs::write(tasks_dir.join(file_name), content).unwrap();
}
fn seed_dispatch_telemetry(project_root: &std::path::Path, engineer: &str, completion_rate: f64) {
std::fs::create_dir_all(project_root.join(".batty")).unwrap();
let conn = telemetry_db::open(project_root).unwrap();
for task_id in 1..=5 {
telemetry_db::insert_event(
&conn,
&TeamEvent::task_assigned(engineer, &task_id.to_string()),
)
.unwrap();
if completion_rate >= 1.0 || task_id as f64 <= completion_rate * 5.0 {
telemetry_db::insert_event(
&conn,
&TeamEvent::task_completed(engineer, Some(&task_id.to_string())),
)
.unwrap();
}
telemetry_db::insert_event(
&conn,
&TeamEvent::quality_metrics_recorded(&QualityMetricsInfo {
backend: "codex",
role: engineer,
task: &task_id.to_string(),
narration_ratio: 0.1,
commit_frequency: 1.0,
first_pass_test_rate: 1.0,
retry_rate: 0.0,
time_to_completion_secs: 120,
}),
)
.unwrap();
}
}
#[test]
fn engineer_task_branch_name_uses_explicit_task_id() {
assert_eq!(
engineer_task_branch_name("eng-1-3", "freeform task body", Some(123)),
"eng-1-3/123"
);
}
#[test]
fn engineer_task_branch_name_extracts_task_id_from_assignment_text() {
assert_eq!(
engineer_task_branch_name("eng-1-3", "Task #456: fix move generation", None),
"eng-1-3/456"
);
}
#[test]
fn engineer_task_branch_name_falls_back_to_slugged_branch() {
let branch = engineer_task_branch_name("eng-1-3", "Fix castling rights sync", None);
assert!(branch.starts_with("eng-1-3/task-fix-castling-rights-sy"));
}
#[test]
fn summarize_assignment_uses_first_non_empty_line() {
assert_eq!(
summarize_assignment("\n\nTask #9: fix move ordering\n\nDetails below"),
"Task #9: fix move ordering"
);
}
#[test]
fn prepare_assignment_launch_rebuilds_task_branch_even_if_engineer_state_is_working() {
let tmp = tempfile::tempdir().unwrap();
let repo = init_git_repo(&tmp, "dispatch-working-prep");
let worktree_dir = repo.join(".batty").join("worktrees").join("eng-1");
let team_config_dir = repo.join(".batty").join("team_config");
setup_engineer_worktree(
&repo,
&worktree_dir,
&engineer_base_branch_name("eng-1"),
&team_config_dir,
)
.unwrap();
let mut daemon = TestDaemonBuilder::new(repo.as_path())
.members(vec![
manager_member("manager", None),
engineer_member("eng-1", Some("manager"), true),
])
.states(HashMap::from([("eng-1".to_string(), MemberState::Working)]))
.build();
daemon.active_tasks.insert("eng-1".to_string(), 40);
let launch = daemon
.prepare_assignment_launch("eng-1", "Task #41: new assignment", Some(41))
.unwrap();
assert_eq!(launch.branch.as_deref(), Some("eng-1/41"));
assert_eq!(current_worktree_branch(&worktree_dir).unwrap(), "eng-1/41");
}
#[test]
fn prepare_assignment_launch_refreshes_stale_worktree_before_dispatch() {
let tmp = tempfile::tempdir().unwrap();
let repo = init_git_repo(&tmp, "dispatch-stale-rebase");
let worktree_dir = repo.join(".batty").join("worktrees").join("eng-1");
let team_config_dir = repo.join(".batty").join("team_config");
let base_branch = engineer_base_branch_name("eng-1");
setup_engineer_worktree(&repo, &worktree_dir, &base_branch, &team_config_dir).unwrap();
std::fs::write(repo.join("main.txt"), "new main content\n").unwrap();
git_ok(&repo, &["add", "main.txt"]);
git_ok(&repo, &["commit", "-m", "advance main"]);
let mut daemon = TestDaemonBuilder::new(repo.as_path())
.members(vec![
manager_member("manager", None),
engineer_member("eng-1", Some("manager"), true),
])
.board(BoardConfig {
worktree_stale_rebase_threshold: 0,
..BoardConfig::default()
})
.build();
let launch = daemon
.prepare_assignment_launch("eng-1", "Task #41: new assignment", Some(41))
.unwrap();
assert_eq!(launch.branch.as_deref(), Some("eng-1/41"));
assert_eq!(current_worktree_branch(&worktree_dir).unwrap(), "eng-1/41");
assert_eq!(
git_stdout(&repo, &["rev-parse", "main"]),
git_stdout(&worktree_dir, &["rev-parse", "HEAD"])
);
assert_eq!(
std::fs::read_to_string(worktree_dir.join("main.txt")).unwrap(),
"new main content\n"
);
let events = events::read_events(&team_events_path(&repo)).unwrap();
let event = events
.iter()
.find(|event| event.event == "worktree_refreshed")
.unwrap();
assert_eq!(event.role.as_deref(), Some("eng-1"));
assert!(
event
.reason
.as_deref()
.unwrap()
.contains("rebased stale worktree"),
"unexpected reason: {:?}",
event.reason
);
}
#[test]
fn prepare_assignment_launch_resets_stale_existing_branch_to_current_main() {
let tmp = tempfile::tempdir().unwrap();
let repo = init_git_repo(&tmp, "dispatch-stale-branch");
let worktree_dir = repo.join(".batty").join("worktrees").join("eng-1");
let team_config_dir = repo.join(".batty").join("team_config");
setup_engineer_worktree(
&repo,
&worktree_dir,
&engineer_base_branch_name("eng-1"),
&team_config_dir,
)
.unwrap();
git_ok(&worktree_dir, &["checkout", "-B", "eng-1/300"]);
std::fs::write(worktree_dir.join("stale.txt"), "stale work\n").unwrap();
git_ok(&worktree_dir, &["add", "stale.txt"]);
git_ok(&worktree_dir, &["commit", "-m", "stale task work"]);
git_ok(&repo, &["checkout", "main"]);
std::fs::write(
repo.join("src").join("lib.rs"),
"pub fn smoke() -> bool { false }\n",
)
.unwrap();
git_ok(&repo, &["add", "src/lib.rs"]);
git_ok(&repo, &["commit", "-m", "advance main"]);
let mut daemon = TestDaemonBuilder::new(repo.as_path())
.members(vec![
manager_member("manager", None),
engineer_member("eng-1", Some("manager"), true),
])
.states(HashMap::from([("eng-1".to_string(), MemberState::Idle)]))
.build();
let launch = daemon
.prepare_assignment_launch("eng-1", "Task #301: current assignment", Some(301))
.unwrap();
assert_eq!(launch.branch.as_deref(), Some("eng-1/301"));
assert_eq!(current_worktree_branch(&worktree_dir).unwrap(), "eng-1/301");
assert_eq!(
git_stdout(&worktree_dir, &["rev-parse", "HEAD"]),
git_stdout(&repo, &["rev-parse", "main"])
);
assert!(!worktree_dir.join("stale.txt").exists());
}
#[test]
fn prepare_assignment_launch_falls_back_to_local_main_when_origin_main_is_frozen() {
let tmp = tempfile::tempdir().unwrap();
let repo = init_git_repo(&tmp, "dispatch-stale-origin");
let worktree_dir = repo.join(".batty").join("worktrees").join("eng-1");
let team_config_dir = repo.join(".batty").join("team_config");
let frozen = git_stdout(&repo, &["rev-parse", "main"]);
setup_engineer_worktree(
&repo,
&worktree_dir,
&engineer_base_branch_name("eng-1"),
&team_config_dir,
)
.unwrap();
git_ok(
&repo,
&["update-ref", "refs/remotes/origin/main", frozen.trim()],
);
for i in 1..=3 {
std::fs::write(repo.join(format!("local-{i}.txt")), format!("local {i}\n")).unwrap();
git_ok(&repo, &["add", "."]);
git_ok(&repo, &["commit", "-m", &format!("advance local main {i}")]);
}
let mut daemon = TestDaemonBuilder::new(repo.as_path())
.members(vec![
manager_member("manager", None),
engineer_member("eng-1", Some("manager"), true),
])
.build();
let launch = daemon
.prepare_assignment_launch("eng-1", "Task #41: new assignment", Some(41))
.unwrap();
assert_eq!(launch.branch.as_deref(), Some("eng-1/41"));
assert_eq!(current_worktree_branch(&worktree_dir).unwrap(), "eng-1/41");
assert_eq!(
git_stdout(&worktree_dir, &["rev-parse", "HEAD"]),
git_stdout(&repo, &["rev-parse", "main"])
);
assert_ne!(
git_stdout(&worktree_dir, &["rev-parse", "HEAD"]),
frozen,
"worktree should not reset to the frozen origin/main"
);
let events = events::read_events(&team_events_path(&repo)).unwrap();
let event = events
.iter()
.find(|event| {
event.event == "worktree_refreshed"
&& event
.reason
.as_deref()
.is_some_and(|reason| reason.contains("stale_origin_fallback ahead=3"))
})
.unwrap();
assert_eq!(event.role.as_deref(), Some("eng-1"));
}
#[test]
fn prepare_assignment_launch_recreates_missing_multi_repo_baseline_from_configured_trunk() {
let tmp = tempfile::tempdir().unwrap();
let project_root = tmp.path().join("workspace");
let repo_name = "pkg-a";
let repo_root = project_root.join(repo_name);
let team_config_dir = project_root.join(".batty").join("team_config");
std::fs::create_dir_all(&team_config_dir).unwrap();
std::fs::create_dir_all(repo_root.join("src")).unwrap();
std::fs::write(repo_root.join("src").join("lib.rs"), "pub fn smoke() {}\n").unwrap();
git_ok(
&project_root,
&["init", "-b", "mainline", repo_root.to_str().unwrap()],
);
git_ok(&repo_root, &["config", "user.email", "batty@example.com"]);
git_ok(&repo_root, &["config", "user.name", "Batty Tests"]);
git_ok(&repo_root, &["add", "."]);
git_ok(&repo_root, &["commit", "-m", "initial mainline"]);
let base_branch = engineer_base_branch_name("eng-1");
let worktree_dir = project_root.join(".batty").join("worktrees").join("eng-1");
setup_multi_repo_worktree_from_trunk(
&project_root,
&worktree_dir,
&base_branch,
&team_config_dir,
&[repo_name.to_string()],
"mainline",
)
.unwrap();
let base_ref_path = repo_root
.join(".git")
.join("refs")
.join("heads")
.join("eng-main")
.join("eng-1");
assert!(base_ref_path.exists());
std::fs::remove_file(&base_ref_path).unwrap();
let missing = git(
&repo_root,
&[
"show-ref",
"--verify",
"--quiet",
"refs/heads/eng-main/eng-1",
],
);
assert!(!missing.status.success());
let mut daemon = TestDaemonBuilder::new(project_root.as_path())
.members(vec![
manager_member("manager", None),
engineer_member("eng-1", Some("manager"), true),
])
.build();
daemon.config.team_config.trunk_branch = "mainline".to_string();
let launch = daemon
.prepare_assignment_launch("eng-1", "Task #683: multi repo assignment", Some(683))
.unwrap();
let sub_worktree = worktree_dir.join(repo_name);
assert_eq!(launch.branch.as_deref(), Some("eng-1/683"));
assert_eq!(current_worktree_branch(&sub_worktree).unwrap(), "eng-1/683");
assert_eq!(
git_stdout(&repo_root, &["rev-parse", "eng-main/eng-1"]),
git_stdout(&repo_root, &["rev-parse", "mainline"])
);
assert_eq!(
git_stdout(&sub_worktree, &["rev-parse", "HEAD"]),
git_stdout(&repo_root, &["rev-parse", "mainline"])
);
let events = events::read_events(&team_events_path(&project_root)).unwrap();
let event = events
.iter()
.find(|event| {
event.event == "worktree_refreshed"
&& event.reason.as_deref().is_some_and(|reason| {
reason.contains("recreated missing baseline branch 'eng-main/eng-1'")
&& reason.contains("from 'mainline'")
&& reason.contains("repo=pkg-a")
})
})
.unwrap();
assert_eq!(event.role.as_deref(), Some("eng-1"));
}
#[test]
fn prepare_assignment_launch_resets_stale_worktree_after_rebase_conflict() {
let tmp = tempfile::tempdir().unwrap();
let repo = init_git_repo(&tmp, "dispatch-stale-reset");
let worktree_dir = repo.join(".batty").join("worktrees").join("eng-1");
let team_config_dir = repo.join(".batty").join("team_config");
let base_branch = engineer_base_branch_name("eng-1");
std::fs::write(repo.join("file.txt"), "base\n").unwrap();
git_ok(&repo, &["add", "file.txt"]);
git_ok(&repo, &["commit", "-m", "add file"]);
setup_engineer_worktree(&repo, &worktree_dir, &base_branch, &team_config_dir).unwrap();
std::fs::write(worktree_dir.join("file.txt"), "engineer change\n").unwrap();
git_ok(&worktree_dir, &["add", "file.txt"]);
git_ok(&worktree_dir, &["commit", "-m", "engineer change"]);
std::fs::write(repo.join("file.txt"), "main change\n").unwrap();
git_ok(&repo, &["add", "file.txt"]);
git_ok(&repo, &["commit", "-m", "main change"]);
let mut daemon = TestDaemonBuilder::new(repo.as_path())
.members(vec![
manager_member("manager", None),
engineer_member("eng-1", Some("manager"), true),
])
.board(BoardConfig {
worktree_stale_rebase_threshold: 0,
..BoardConfig::default()
})
.build();
let launch = daemon
.prepare_assignment_launch("eng-1", "Task #42: conflicted assignment", Some(42))
.unwrap();
assert_eq!(launch.branch.as_deref(), Some("eng-1/42"));
assert_eq!(current_worktree_branch(&worktree_dir).unwrap(), "eng-1/42");
assert_eq!(
git_stdout(&repo, &["rev-parse", "main"]),
git_stdout(&worktree_dir, &["rev-parse", "HEAD"])
);
assert_eq!(
std::fs::read_to_string(worktree_dir.join("file.txt")).unwrap(),
"main change\n"
);
let events = events::read_events(&team_events_path(&repo)).unwrap();
let event = events
.iter()
.find(|event| event.event == "worktree_refreshed")
.unwrap();
assert_eq!(event.role.as_deref(), Some("eng-1"));
assert!(
event
.reason
.as_deref()
.unwrap()
.contains("reset stale worktree after rebase failed"),
"unexpected reason: {:?}",
event.reason
);
}
#[test]
fn prepare_assignment_launch_refuses_when_claimed_task_branch_mismatch_exists() {
let tmp = tempfile::tempdir().unwrap();
let repo = init_git_repo(&tmp, "dispatch-branch-mismatch");
let worktree_dir = repo.join(".batty").join("worktrees").join("eng-1");
let team_config_dir = repo.join(".batty").join("team_config");
let base_branch = engineer_base_branch_name("eng-1");
setup_engineer_worktree(&repo, &worktree_dir, &base_branch, &team_config_dir).unwrap();
git_ok(&worktree_dir, &["checkout", "-b", "eng-1/41"]);
write_task_file(
&repo,
"042-claimed-task.md",
"---\nid: 42\ntitle: claimed-task\nstatus: in-progress\npriority: critical\nclaimed_by: eng-1\nbranch: eng-1/42\nworktree_path: .batty/worktrees/eng-1\nclass: standard\n---\n\nTask description.\n",
);
let mut daemon = TestDaemonBuilder::new(repo.as_path())
.members(vec![
manager_member("manager", None),
engineer_member("eng-1", Some("manager"), true),
])
.states(HashMap::from([("eng-1".to_string(), MemberState::Idle)]))
.build();
let err = daemon
.prepare_assignment_launch("eng-1", "Task #99: new assignment", Some(99))
.unwrap_err();
assert!(err.to_string().contains("claimed task #42"));
assert_eq!(current_worktree_branch(&worktree_dir).unwrap(), "eng-1/41");
let inbox_root = inbox::inboxes_root(&repo);
let manager_messages = inbox::pending_messages(&inbox_root, "manager").unwrap();
assert!(manager_messages.iter().any(|message| {
message.body.contains("Reconciliation alert")
&& message.body.contains("eng-1/41")
&& message.body.contains("eng-1/42")
}));
}
#[test]
fn shim_assignment_sends_message_to_existing_engineer() {
let tmp = tempfile::tempdir().unwrap();
std::fs::create_dir_all(
tmp.path()
.join(".batty")
.join("team_config")
.join("board")
.join("tasks"),
)
.unwrap();
let mut daemon = TestDaemonBuilder::new(tmp.path())
.members(vec![
manager_member("manager", None),
engineer_member("eng-1", Some("manager"), false),
])
.build();
daemon.config.team_config.use_shim = true;
let (parent_sock, child_sock) = socketpair().unwrap();
let parent_channel = crate::shim::protocol::Channel::new(parent_sock);
let mut child_channel = crate::shim::protocol::Channel::new(child_sock);
let mut handle = AgentHandle::new(
"eng-1".to_string(),
parent_channel,
12345,
"codex".to_string(),
"codex".to_string(),
tmp.path().to_path_buf(),
);
handle.apply_state_change(ShimState::Idle);
daemon.shim_handles.insert("eng-1".to_string(), handle);
let launch = daemon
.assign_task_with_task_id_as("manager", "eng-1", "Task #42: fix it", Some(42))
.unwrap();
let cmd: Command = child_channel.recv().unwrap().unwrap();
match cmd {
Command::SendMessage { from, body, .. } => {
assert_eq!(from, "manager");
assert_eq!(body, "Task #42: fix it");
}
other => panic!("expected SendMessage, got {other:?}"),
}
assert_ne!(daemon.states.get("eng-1"), Some(&MemberState::Working));
assert_eq!(launch.branch, None);
assert_eq!(launch.work_dir, tmp.path());
}
#[test]
fn assign_task_with_task_id_records_assignment_context_for_worktree_assignments() {
let tmp = tempfile::tempdir().unwrap();
let repo = init_git_repo(&tmp, "dispatch-record-branch");
let worktree_dir = repo.join(".batty").join("worktrees").join("eng-1");
let team_config_dir = repo.join(".batty").join("team_config");
setup_engineer_worktree(
&repo,
&worktree_dir,
&engineer_base_branch_name("eng-1"),
&team_config_dir,
)
.unwrap();
write_task_file(
&repo,
"042-active-task.md",
"---\nid: 42\ntitle: active-task\nstatus: todo\npriority: high\nclass: standard\n---\n\nTask body.\n",
);
let mut daemon = TestDaemonBuilder::new(repo.as_path())
.members(vec![
manager_member("manager", None),
engineer_member("eng-1", Some("manager"), true),
])
.build();
daemon.config.team_config.use_shim = true;
let (parent_sock, child_sock) = socketpair().unwrap();
let parent_channel = crate::shim::protocol::Channel::new(parent_sock);
let mut child_channel = crate::shim::protocol::Channel::new(child_sock);
let mut handle = AgentHandle::new(
"eng-1".to_string(),
parent_channel,
12345,
"codex".to_string(),
"codex".to_string(),
worktree_dir.clone(),
);
handle.apply_state_change(ShimState::Idle);
daemon.shim_handles.insert("eng-1".to_string(), handle);
let launch = daemon
.assign_task_with_task_id_as("manager", "eng-1", "Task #42: fix it", Some(42))
.unwrap();
let cmd: Command = child_channel.recv().unwrap().unwrap();
match cmd {
Command::SendMessage { from, body, .. } => {
assert_eq!(from, "manager");
assert!(body.contains("Task #42: active-task"));
assert!(body.contains("Assignment Packet:"));
}
other => panic!("expected SendMessage, got {other:?}"),
}
let task = crate::task::Task::from_file(
&repo
.join(".batty")
.join("team_config")
.join("board")
.join("tasks")
.join("042-active-task.md"),
)
.unwrap();
assert_eq!(launch.branch.as_deref(), Some("eng-1/42"));
assert_eq!(task.branch.as_deref(), Some("eng-1/42"));
assert_eq!(
task.worktree_path.as_deref(),
Some(".batty/worktrees/eng-1")
);
}
#[test]
fn assignment_guard_rejects_second_active_task_for_engineer() {
let tmp = tempfile::tempdir().unwrap();
crate::team::test_support::write_owned_task_file(
tmp.path(),
91,
"active-task",
"in-progress",
"eng-1",
);
let mut daemon = TestDaemonBuilder::new(tmp.path())
.members(vec![
manager_member("manager", None),
engineer_member("eng-1", Some("manager"), false),
])
.build();
daemon.config.team_config.use_shim = true;
let (parent_sock, _child_sock) = socketpair().unwrap();
let parent_channel = crate::shim::protocol::Channel::new(parent_sock);
let mut handle = AgentHandle::new(
"eng-1".to_string(),
parent_channel,
12345,
"codex".to_string(),
"codex".to_string(),
tmp.path().to_path_buf(),
);
handle.apply_state_change(ShimState::Idle);
daemon.shim_handles.insert("eng-1".to_string(), handle);
let error = daemon
.assign_task_with_task_id_as("manager", "eng-1", "Task #42: fix it", Some(42))
.unwrap_err()
.to_string();
assert!(error.contains("already owns active board task(s) #91"));
}
#[test]
fn assignment_guard_does_not_mutate_blocked_task_assignment_context() {
let tmp = tempfile::tempdir().unwrap();
let repo = init_git_repo(&tmp, "dispatch-guard-no-mutation");
write_task_file(
&repo,
"091-active-task.md",
"---\nid: 91\ntitle: active-task\nstatus: in-progress\npriority: high\nclaimed_by: eng-1\nclass: standard\n---\n\nTask body.\n",
);
write_task_file(
&repo,
"092-next-task.md",
"---\nid: 92\ntitle: next-task\nstatus: todo\npriority: high\nclass: standard\n---\n\nTask body.\n",
);
let mut daemon = TestDaemonBuilder::new(repo.as_path())
.members(vec![
manager_member("manager", None),
engineer_member("eng-1", Some("manager"), true),
])
.build();
daemon.config.team_config.use_shim = true;
let error = daemon
.assign_task_with_task_id_as("manager", "eng-1", "Task #92: do next", Some(92))
.unwrap_err()
.to_string();
assert!(error.contains("already owns active board task(s) #91"));
let task = crate::task::Task::from_file(
&repo
.join(".batty")
.join("team_config")
.join("board")
.join("tasks")
.join("092-next-task.md"),
)
.unwrap();
assert!(task.branch.is_none());
assert!(task.worktree_path.is_none());
}
#[test]
fn assignment_guard_allows_resuming_same_active_task() {
let tmp = tempfile::tempdir().unwrap();
crate::team::test_support::write_owned_task_file(
tmp.path(),
91,
"active-task",
"in-progress",
"eng-1",
);
let mut daemon = TestDaemonBuilder::new(tmp.path())
.members(vec![
manager_member("manager", None),
engineer_member("eng-1", Some("manager"), false),
])
.build();
daemon.config.team_config.use_shim = true;
let (parent_sock, child_sock) = socketpair().unwrap();
let parent_channel = crate::shim::protocol::Channel::new(parent_sock);
let mut child_channel = crate::shim::protocol::Channel::new(child_sock);
let mut handle = AgentHandle::new(
"eng-1".to_string(),
parent_channel,
12345,
"codex".to_string(),
"codex".to_string(),
tmp.path().to_path_buf(),
);
handle.apply_state_change(ShimState::Idle);
daemon.shim_handles.insert("eng-1".to_string(), handle);
daemon
.assign_task_with_task_id_as("manager", "eng-1", "Task #91: continue", Some(91))
.unwrap();
let cmd: Command = child_channel.recv().unwrap().unwrap();
match cmd {
Command::SendMessage { from, body, .. } => {
assert_eq!(from, "manager");
assert!(body.contains("Task #91: active-task"));
assert!(body.contains("Assignment Packet:"));
assert!(body.contains("scope_ack_required: false"));
}
other => panic!("expected SendMessage, got {other:?}"),
}
}
#[test]
fn stabilization_delay_prevents_premature_dispatch() {
let tmp = tempfile::tempdir().unwrap();
write_open_task_file(tmp.path(), 101, "queued-task", "todo");
let members = vec![
manager_member("manager", None),
engineer_member("eng-1", Some("manager"), false),
];
let mut daemon = TestDaemonBuilder::new(tmp.path())
.members(members)
.board(BoardConfig {
auto_dispatch: true,
dispatch_stabilization_delay_secs: 30,
..BoardConfig::default()
})
.states(HashMap::from([("eng-1".to_string(), MemberState::Idle)]))
.build();
daemon.last_auto_dispatch = Instant::now() - Duration::from_secs(30);
daemon
.idle_started_at
.insert("eng-1".to_string(), Instant::now() - Duration::from_secs(5));
daemon.enqueue_dispatch_candidates().unwrap();
daemon.process_dispatch_queue().unwrap();
assert_eq!(daemon.dispatch_queue.len(), 1);
assert_eq!(daemon.dispatch_queue[0].validation_failures, 0);
assert_eq!(daemon.dispatch_queue[0].task_id, 101);
}
#[test]
fn wip_gate_blocks_double_assignment() {
let tmp = tempfile::tempdir().unwrap();
write_open_task_file(tmp.path(), 101, "queued-task", "todo");
crate::team::test_support::write_owned_task_file(
tmp.path(),
91,
"active-task",
"in-progress",
"eng-1",
);
let members = vec![
manager_member("manager", None),
engineer_member("eng-1", Some("manager"), false),
];
let mut daemon = TestDaemonBuilder::new(tmp.path())
.members(members)
.workflow_policy(WorkflowPolicy {
wip_limit_per_engineer: Some(1),
..WorkflowPolicy::default()
})
.board(BoardConfig {
auto_dispatch: true,
dispatch_stabilization_delay_secs: 0,
..BoardConfig::default()
})
.states(HashMap::from([("eng-1".to_string(), MemberState::Idle)]))
.build();
daemon.last_auto_dispatch = Instant::now() - Duration::from_secs(30);
daemon.idle_started_at.insert(
"eng-1".to_string(),
Instant::now() - Duration::from_secs(60),
);
daemon.enqueue_dispatch_candidates().unwrap();
daemon.process_dispatch_queue().unwrap();
assert_eq!(daemon.dispatch_queue.len(), 1);
assert_eq!(daemon.dispatch_queue[0].validation_failures, 1);
assert!(
daemon.dispatch_queue[0]
.last_failure
.as_deref()
.unwrap_or_default()
.contains("Dispatch guard")
);
}
#[test]
fn dispatch_guard_blocks_claimed_todo_assignment() {
let tmp = tempfile::tempdir().unwrap();
write_open_task_file(tmp.path(), 101, "queued-task", "todo");
crate::team::test_support::write_owned_task_file(
tmp.path(),
91,
"claimed-todo",
"todo",
"eng-1",
);
let members = vec![
manager_member("manager", None),
engineer_member("eng-1", Some("manager"), false),
];
let mut daemon = TestDaemonBuilder::new(tmp.path())
.members(members)
.board(BoardConfig {
auto_dispatch: true,
dispatch_stabilization_delay_secs: 0,
..BoardConfig::default()
})
.states(HashMap::from([("eng-1".to_string(), MemberState::Idle)]))
.build();
daemon.last_auto_dispatch = Instant::now() - Duration::from_secs(30);
daemon.idle_started_at.insert(
"eng-1".to_string(),
Instant::now() - Duration::from_secs(60),
);
daemon.enqueue_dispatch_candidates().unwrap();
daemon.process_dispatch_queue().unwrap();
assert_eq!(daemon.dispatch_queue.len(), 1);
assert_eq!(daemon.dispatch_queue[0].validation_failures, 1);
assert!(
daemon.dispatch_queue[0]
.last_failure
.as_deref()
.unwrap_or_default()
.contains("Dispatch guard blocked assignment")
);
}
#[test]
fn active_board_item_count_includes_todo_in_progress_and_review() {
let tmp = tempfile::tempdir().unwrap();
crate::team::test_support::write_owned_task_file(tmp.path(), 11, "todo-task", "todo", "eng-1");
crate::team::test_support::write_owned_task_file(
tmp.path(),
12,
"working-task",
"in-progress",
"eng-1",
);
let tasks_dir = tmp
.path()
.join(".batty")
.join("team_config")
.join("board")
.join("tasks");
std::fs::create_dir_all(&tasks_dir).unwrap();
std::fs::write(
tasks_dir.join("013-review-task.md"),
"---\nid: 13\ntitle: review-task\nstatus: review\npriority: critical\nclaimed_by: manager\nreview_owner: eng-1\nclass: standard\n---\n\nTask description.\n",
)
.unwrap();
let daemon = TestDaemonBuilder::new(tmp.path()).build();
let board_dir = tmp.path().join(".batty").join("team_config").join("board");
assert_eq!(
daemon
.engineer_active_board_item_count(&board_dir, "eng-1")
.unwrap(),
3
);
}
#[test]
fn worktree_gate_blocks_dirty_worktrees() {
let tmp = tempfile::tempdir().unwrap();
let repo = init_git_repo(&tmp, "dispatch-queue");
write_open_task_file(&repo, 101, "queued-task", "todo");
let team_config_dir = repo.join(".batty").join("team_config");
let worktree_dir = repo.join(".batty").join("worktrees").join("eng-1");
setup_engineer_worktree(
&repo,
&worktree_dir,
&engineer_base_branch_name("eng-1"),
&team_config_dir,
)
.unwrap();
std::fs::write(worktree_dir.join("DIRTY.txt"), "dirty\n").unwrap();
let members = vec![
manager_member("manager", None),
engineer_member("eng-1", Some("manager"), true),
];
let mut daemon = TestDaemonBuilder::new(&repo)
.members(members)
.board(BoardConfig {
auto_dispatch: true,
dispatch_stabilization_delay_secs: 0,
..BoardConfig::default()
})
.states(HashMap::from([("eng-1".to_string(), MemberState::Idle)]))
.build();
daemon.last_auto_dispatch = Instant::now() - Duration::from_secs(30);
daemon.idle_started_at.insert(
"eng-1".to_string(),
Instant::now() - Duration::from_secs(60),
);
daemon.enqueue_dispatch_candidates().unwrap();
daemon.process_dispatch_queue().unwrap();
assert_eq!(daemon.dispatch_queue.len(), 1);
assert_eq!(
daemon.dispatch_queue[0].validation_failures, 2,
"dirty worktree recovery should record the readiness failure and the skipped reset"
);
assert!(
daemon.dispatch_queue[0]
.last_failure
.as_deref()
.unwrap_or("")
.contains("could not safely auto-save dirty worktree"),
"dirty worktree recovery should retain the preservation blocker"
);
assert!(
crate::team::test_support::git_stdout(&worktree_dir, &["status", "--short"])
.contains("?? DIRTY.txt"),
"dirty file should remain untouched for manual recovery"
);
}
#[test]
fn worktree_gate_resets_base_branch_that_remains_ahead_after_rebase() {
let tmp = tempfile::tempdir().unwrap();
let repo = init_git_repo(&tmp, "dispatch-queue-ahead-base");
write_open_task_file(&repo, 101, "queued-task", "todo");
let team_config_dir = repo.join(".batty").join("team_config");
let worktree_dir = repo.join(".batty").join("worktrees").join("eng-1");
let base_branch = engineer_base_branch_name("eng-1");
setup_engineer_worktree(&repo, &worktree_dir, &base_branch, &team_config_dir).unwrap();
std::fs::write(worktree_dir.join("AHEAD.txt"), "ahead\n").unwrap();
assert!(
std::process::Command::new("git")
.args(["add", "AHEAD.txt"])
.current_dir(&worktree_dir)
.status()
.unwrap()
.success()
);
assert!(
std::process::Command::new("git")
.args(["commit", "-m", "ahead base branch commit"])
.current_dir(&worktree_dir)
.status()
.unwrap()
.success()
);
let members = vec![
manager_member("manager", None),
engineer_member("eng-1", Some("manager"), true),
];
let mut daemon = TestDaemonBuilder::new(&repo)
.members(members)
.board(BoardConfig {
auto_dispatch: true,
dispatch_stabilization_delay_secs: 0,
..BoardConfig::default()
})
.states(HashMap::from([("eng-1".to_string(), MemberState::Idle)]))
.build();
daemon.last_auto_dispatch = Instant::now() - Duration::from_secs(30);
daemon.idle_started_at.insert(
"eng-1".to_string(),
Instant::now() - Duration::from_secs(60),
);
daemon.maybe_auto_dispatch().unwrap();
assert_eq!(daemon.dispatch_queue.len(), 1);
assert_eq!(daemon.dispatch_queue[0].validation_failures, 0);
assert!(daemon.dispatch_queue[0].last_failure.is_none());
assert_eq!(
crate::worktree::git_current_branch(&worktree_dir).unwrap(),
base_branch,
"recovery should leave the engineer on the cleaned base branch"
);
let ahead = std::process::Command::new("git")
.args(["rev-list", "--count", "main..HEAD"])
.current_dir(&worktree_dir)
.output()
.unwrap();
assert!(ahead.status.success());
assert_eq!(String::from_utf8_lossy(&ahead.stdout).trim(), "0");
}
#[test]
fn queue_escalates_after_repeated_validation_failures() {
let tmp = tempfile::tempdir().unwrap();
write_open_task_file(tmp.path(), 101, "queued-task", "todo");
crate::team::test_support::write_owned_task_file(
tmp.path(),
91,
"active-task",
"in-progress",
"eng-1",
);
let members = vec![
manager_member("manager", None),
engineer_member("eng-1", Some("manager"), false),
];
let mut daemon = TestDaemonBuilder::new(tmp.path())
.members(members)
.workflow_policy(WorkflowPolicy {
wip_limit_per_engineer: Some(1),
..WorkflowPolicy::default()
})
.board(BoardConfig {
auto_dispatch: true,
dispatch_stabilization_delay_secs: 0,
..BoardConfig::default()
})
.states(HashMap::from([
("eng-1".to_string(), MemberState::Idle),
("manager".to_string(), MemberState::Idle),
]))
.build();
daemon.last_auto_dispatch = Instant::now() - Duration::from_secs(30);
daemon.idle_started_at.insert(
"eng-1".to_string(),
Instant::now() - Duration::from_secs(60),
);
for _ in 0..DISPATCH_QUEUE_FAILURE_LIMIT {
daemon.last_auto_dispatch = Instant::now() - Duration::from_secs(30);
daemon.maybe_auto_dispatch().unwrap();
}
assert!(
daemon.dispatch_queue.is_empty(),
"queue should be drained after failure limit"
);
let inbox_root = inbox::inboxes_root(tmp.path());
let manager_messages = inbox::pending_messages(&inbox_root, "manager").unwrap();
assert_eq!(manager_messages.len(), 0);
}
#[test]
fn cwd_correction_handles_symlinks() {
let tmp = tempfile::tempdir().unwrap();
let canonical = tmp.path().canonicalize().unwrap();
let from_canonical = normalized_assignment_dir(&canonical);
let from_raw = normalized_assignment_dir(tmp.path());
assert_eq!(from_canonical, from_raw);
#[cfg(target_os = "macos")]
{
let path_str = canonical.to_string_lossy();
if path_str.starts_with("/private/") {
let without_private = PathBuf::from(&path_str["/private".len()..]);
let normalized = normalized_assignment_dir(&without_private);
assert_eq!(normalized, from_canonical);
}
}
}
#[test]
fn cwd_correction_normalizes_nonexistent_paths_to_self() {
let bogus = PathBuf::from("/nonexistent/path/that/does/not/exist");
let normalized = normalized_assignment_dir(&bogus);
assert_eq!(normalized, bogus);
}
#[test]
fn cwd_correction_retries_on_stale_read() {
let tmp = tempfile::tempdir().unwrap();
let expected = tmp.path().to_path_buf();
let normalized_expected = normalized_assignment_dir(&expected);
let project_root = tmp.path().parent().unwrap_or(tmp.path());
let stale_path = normalized_assignment_dir(project_root);
assert_ne!(
stale_path, normalized_expected,
"stale path should differ from expected"
);
let corrected = normalized_assignment_dir(&expected);
assert_eq!(
corrected, normalized_expected,
"corrected path should match expected"
);
let codex_dir = expected.join(".batty").join("codex-context").join("eng-1");
std::fs::create_dir_all(&codex_dir).unwrap();
let codex_normalized = normalized_assignment_dir(&codex_dir);
assert_ne!(codex_normalized, normalized_expected);
}
fn write_scheduled_task_file(
project_root: &Path,
id: u32,
title: &str,
status: &str,
scheduled_for: &str,
) {
let tasks_dir = project_root
.join(".batty")
.join("team_config")
.join("board")
.join("tasks");
std::fs::create_dir_all(&tasks_dir).unwrap();
std::fs::write(
tasks_dir.join(format!("{id:03}-{title}.md")),
format!(
"---\nid: {id}\ntitle: {title}\nstatus: {status}\npriority: high\nscheduled_for: \"{scheduled_for}\"\nclass: standard\n---\n\nTask description.\n"
),
)
.unwrap();
}
#[test]
fn dispatch_skips_future_scheduled_task() {
let future = (chrono::Utc::now() + chrono::Duration::hours(1)).to_rfc3339();
let tmp = tempfile::tempdir().unwrap();
write_scheduled_task_file(tmp.path(), 101, "future-task", "todo", &future);
let members = vec![
manager_member("manager", None),
engineer_member("eng-1", Some("manager"), false),
];
let mut daemon = TestDaemonBuilder::new(tmp.path())
.members(members)
.board(BoardConfig {
auto_dispatch: true,
dispatch_stabilization_delay_secs: 0,
..BoardConfig::default()
})
.states(HashMap::from([("eng-1".to_string(), MemberState::Idle)]))
.build();
daemon.last_auto_dispatch = Instant::now() - Duration::from_secs(30);
daemon.idle_started_at.insert(
"eng-1".to_string(),
Instant::now() - Duration::from_secs(60),
);
daemon.maybe_auto_dispatch().unwrap();
assert!(
daemon.dispatch_queue.is_empty(),
"future-scheduled task should not be dispatched"
);
}
#[test]
fn dispatch_includes_past_scheduled_task() {
let past = (chrono::Utc::now() - chrono::Duration::hours(1)).to_rfc3339();
let tmp = tempfile::tempdir().unwrap();
write_scheduled_task_file(tmp.path(), 101, "past-task", "todo", &past);
let members = vec![
manager_member("manager", None),
engineer_member("eng-1", Some("manager"), false),
];
let mut daemon = TestDaemonBuilder::new(tmp.path())
.members(members)
.board(BoardConfig {
auto_dispatch: true,
dispatch_stabilization_delay_secs: 0,
..BoardConfig::default()
})
.states(HashMap::from([("eng-1".to_string(), MemberState::Idle)]))
.build();
daemon.last_auto_dispatch = Instant::now() - Duration::from_secs(30);
daemon.idle_started_at.insert(
"eng-1".to_string(),
Instant::now() - Duration::from_secs(60),
);
daemon.maybe_auto_dispatch().unwrap();
assert_eq!(daemon.dispatch_queue.len(), 1);
assert_eq!(daemon.dispatch_queue[0].task_id, 101);
}
#[test]
fn completion_frees_engineer_and_dispatch_skips_already_claimed_tasks() {
let tmp = tempfile::tempdir().unwrap();
write_owned_task_file(tmp.path(), 90, "finished-task", "done", "eng-1");
write_owned_task_file(tmp.path(), 101, "claimed-next", "todo", "eng-2");
write_open_task_file(tmp.path(), 102, "dispatchable-next", "todo");
let members = vec![
manager_member("manager", None),
engineer_member("eng-1", Some("manager"), false),
engineer_member("eng-2", Some("manager"), false),
];
let mut daemon = TestDaemonBuilder::new(tmp.path())
.members(members)
.board(BoardConfig {
auto_dispatch: true,
dispatch_stabilization_delay_secs: 0,
..BoardConfig::default()
})
.states(HashMap::from([
("eng-1".to_string(), MemberState::Idle),
("eng-2".to_string(), MemberState::Idle),
]))
.build();
daemon.last_auto_dispatch = Instant::now() - Duration::from_secs(30);
daemon.idle_started_at.insert(
"eng-1".to_string(),
Instant::now() - Duration::from_secs(60),
);
daemon.idle_started_at.insert(
"eng-2".to_string(),
Instant::now() - Duration::from_secs(60),
);
daemon.maybe_auto_dispatch().unwrap();
assert_eq!(
daemon.dispatch_queue.len(),
1,
"only the next unclaimed todo task should be queued"
);
assert_eq!(daemon.dispatch_queue[0].engineer, "eng-1");
assert_eq!(
daemon.dispatch_queue[0].task_id, 102,
"claimed todo task should not be re-dispatched after completion"
);
}
#[test]
fn dispatch_preps_worktree_for_idle_engineer() {
let tmp = tempfile::tempdir().unwrap();
let repo = init_git_repo(&tmp, "idle-prep");
write_open_task_file(&repo, 301, "idle-task", "todo");
let team_config_dir = repo.join(".batty").join("team_config");
let worktree_dir = repo.join(".batty").join("worktrees").join("eng-1");
setup_engineer_worktree(
&repo,
&worktree_dir,
&engineer_base_branch_name("eng-1"),
&team_config_dir,
)
.unwrap();
let members = vec![
manager_member("manager", None),
engineer_member("eng-1", Some("manager"), true),
];
let mut daemon = TestDaemonBuilder::new(&repo)
.members(members)
.pane_map(HashMap::from([("eng-1".to_string(), "%99".to_string())]))
.board(BoardConfig {
auto_dispatch: true,
dispatch_stabilization_delay_secs: 0,
..BoardConfig::default()
})
.states(HashMap::from([("eng-1".to_string(), MemberState::Idle)]))
.build();
daemon.last_auto_dispatch = Instant::now() - Duration::from_secs(30);
daemon.idle_started_at.insert(
"eng-1".to_string(),
Instant::now() - Duration::from_secs(60),
);
daemon.maybe_auto_dispatch().unwrap();
let task_branch = "eng-1/301";
let branches = git_stdout(&repo, &["branch", "--list"]);
assert!(
branches.contains(task_branch),
"idle engineer should have task branch created by worktree prep; branches: {branches}"
);
}
#[test]
fn e2e_past_scheduled_for_is_dispatchable() {
use crate::team::resolver::{ResolutionStatus, resolve_board};
let past = (chrono::Utc::now() - chrono::Duration::hours(1)).to_rfc3339();
let tmp = tempfile::tempdir().unwrap();
write_scheduled_task_file(tmp.path(), 501, "past-sched", "todo", &past);
let board_dir = tmp.path().join(".batty").join("team_config").join("board");
let members_list = vec![
manager_member("manager", None),
engineer_member("eng-1", Some("manager"), false),
];
let resolutions = resolve_board(&board_dir, &members_list).unwrap();
assert_eq!(resolutions.len(), 1);
assert_eq!(
resolutions[0].status,
ResolutionStatus::Runnable,
"past scheduled_for task should be Runnable"
);
let mut daemon = TestDaemonBuilder::new(tmp.path())
.members(members_list)
.board(BoardConfig {
auto_dispatch: true,
dispatch_stabilization_delay_secs: 0,
..BoardConfig::default()
})
.states(HashMap::from([("eng-1".to_string(), MemberState::Idle)]))
.build();
daemon.last_auto_dispatch = Instant::now() - Duration::from_secs(30);
daemon.idle_started_at.insert(
"eng-1".to_string(),
Instant::now() - Duration::from_secs(60),
);
daemon.maybe_auto_dispatch().unwrap();
assert_eq!(
daemon.dispatch_queue.len(),
1,
"past-scheduled task should be dispatched"
);
assert_eq!(daemon.dispatch_queue[0].task_id, 501);
}
#[test]
fn e2e_future_scheduled_for_is_blocked() {
use crate::team::resolver::{ResolutionStatus, resolve_board};
let future = (chrono::Utc::now() + chrono::Duration::hours(1)).to_rfc3339();
let tmp = tempfile::tempdir().unwrap();
write_scheduled_task_file(tmp.path(), 502, "future-sched", "todo", &future);
let board_dir = tmp.path().join(".batty").join("team_config").join("board");
let members_list = vec![
manager_member("manager", None),
engineer_member("eng-1", Some("manager"), false),
];
let resolutions = resolve_board(&board_dir, &members_list).unwrap();
assert_eq!(resolutions.len(), 1);
assert_eq!(
resolutions[0].status,
ResolutionStatus::Blocked,
"future scheduled_for task should be Blocked"
);
assert!(
resolutions[0]
.blocking_reason
.as_ref()
.unwrap()
.contains("scheduled for"),
"blocking reason should mention 'scheduled for'"
);
let mut daemon = TestDaemonBuilder::new(tmp.path())
.members(members_list)
.board(BoardConfig {
auto_dispatch: true,
dispatch_stabilization_delay_secs: 0,
..BoardConfig::default()
})
.states(HashMap::from([("eng-1".to_string(), MemberState::Idle)]))
.build();
daemon.last_auto_dispatch = Instant::now() - Duration::from_secs(30);
daemon.idle_started_at.insert(
"eng-1".to_string(),
Instant::now() - Duration::from_secs(60),
);
daemon.maybe_auto_dispatch().unwrap();
assert!(
daemon.dispatch_queue.is_empty(),
"future-scheduled task should NOT be dispatched"
);
}
#[test]
fn e2e_no_scheduled_for_always_runnable() {
use crate::team::resolver::{ResolutionStatus, resolve_board};
let tmp = tempfile::tempdir().unwrap();
write_open_task_file(tmp.path(), 503, "no-schedule", "todo");
let board_dir = tmp.path().join(".batty").join("team_config").join("board");
let members_list = vec![
manager_member("manager", None),
engineer_member("eng-1", Some("manager"), false),
];
let resolutions = resolve_board(&board_dir, &members_list).unwrap();
assert_eq!(resolutions.len(), 1);
assert_eq!(
resolutions[0].status,
ResolutionStatus::Runnable,
"task without scheduled_for should be Runnable"
);
assert!(
resolutions[0].blocking_reason.is_none(),
"no blocking reason expected"
);
let mut daemon = TestDaemonBuilder::new(tmp.path())
.members(members_list)
.board(BoardConfig {
auto_dispatch: true,
dispatch_stabilization_delay_secs: 0,
..BoardConfig::default()
})
.states(HashMap::from([("eng-1".to_string(), MemberState::Idle)]))
.build();
daemon.last_auto_dispatch = Instant::now() - Duration::from_secs(30);
daemon.idle_started_at.insert(
"eng-1".to_string(),
Instant::now() - Duration::from_secs(60),
);
daemon.maybe_auto_dispatch().unwrap();
assert_eq!(
daemon.dispatch_queue.len(),
1,
"task without scheduled_for should be dispatched"
);
assert_eq!(daemon.dispatch_queue[0].task_id, 503);
}
#[test]
fn dedup_window_prevents_duplicate_enqueue() {
let tmp = tempfile::tempdir().unwrap();
write_open_task_file(tmp.path(), 101, "queued-task", "todo");
let members = vec![
manager_member("manager", None),
engineer_member("eng-1", Some("manager"), false),
];
let mut daemon = TestDaemonBuilder::new(tmp.path())
.members(members)
.board(BoardConfig {
auto_dispatch: true,
dispatch_stabilization_delay_secs: 0,
dispatch_dedup_window_secs: 60,
..BoardConfig::default()
})
.states(HashMap::from([("eng-1".to_string(), MemberState::Idle)]))
.build();
daemon.last_auto_dispatch = Instant::now() - Duration::from_secs(30);
daemon.idle_started_at.insert(
"eng-1".to_string(),
Instant::now() - Duration::from_secs(60),
);
daemon
.recent_dispatches
.insert((101, "eng-1".to_string()), Instant::now());
daemon.maybe_auto_dispatch().unwrap();
assert!(
daemon.dispatch_queue.is_empty(),
"task should be skipped due to dedup window"
);
}
#[test]
fn dedup_window_allows_different_task() {
let tmp = tempfile::tempdir().unwrap();
write_open_task_file(tmp.path(), 101, "first-task", "todo");
write_open_task_file(tmp.path(), 102, "second-task", "todo");
let members = vec![
manager_member("manager", None),
engineer_member("eng-1", Some("manager"), false),
];
let mut daemon = TestDaemonBuilder::new(tmp.path())
.members(members)
.board(BoardConfig {
auto_dispatch: true,
dispatch_stabilization_delay_secs: 0,
dispatch_dedup_window_secs: 60,
..BoardConfig::default()
})
.states(HashMap::from([("eng-1".to_string(), MemberState::Idle)]))
.build();
daemon.last_auto_dispatch = Instant::now() - Duration::from_secs(30);
daemon.idle_started_at.insert(
"eng-1".to_string(),
Instant::now() - Duration::from_secs(60),
);
daemon
.recent_dispatches
.insert((102, "eng-1".to_string()), Instant::now());
daemon.maybe_auto_dispatch().unwrap();
assert_eq!(daemon.dispatch_queue.len(), 1);
assert_eq!(
daemon.dispatch_queue[0].task_id, 101,
"a different task should still be dispatched"
);
}
#[test]
fn dedup_window_expires_and_allows_reassignment() {
let tmp = tempfile::tempdir().unwrap();
write_open_task_file(tmp.path(), 101, "queued-task", "todo");
let members = vec![
manager_member("manager", None),
engineer_member("eng-1", Some("manager"), false),
];
let mut daemon = TestDaemonBuilder::new(tmp.path())
.members(members)
.board(BoardConfig {
auto_dispatch: true,
dispatch_stabilization_delay_secs: 0,
dispatch_dedup_window_secs: 60,
..BoardConfig::default()
})
.states(HashMap::from([("eng-1".to_string(), MemberState::Idle)]))
.build();
daemon.last_auto_dispatch = Instant::now() - Duration::from_secs(30);
daemon.idle_started_at.insert(
"eng-1".to_string(),
Instant::now() - Duration::from_secs(60),
);
daemon.recent_dispatches.insert(
(101, "eng-1".to_string()),
Instant::now() - Duration::from_secs(120),
);
daemon.maybe_auto_dispatch().unwrap();
assert_eq!(
daemon.dispatch_queue.len(),
1,
"expired dedup entry should allow reassignment"
);
assert_eq!(daemon.dispatch_queue[0].task_id, 101);
}
#[test]
fn dedup_window_zero_disables_dedup() {
let tmp = tempfile::tempdir().unwrap();
write_open_task_file(tmp.path(), 101, "queued-task", "todo");
let members = vec![
manager_member("manager", None),
engineer_member("eng-1", Some("manager"), false),
];
let mut daemon = TestDaemonBuilder::new(tmp.path())
.members(members)
.board(BoardConfig {
auto_dispatch: true,
dispatch_stabilization_delay_secs: 0,
dispatch_dedup_window_secs: 0,
..BoardConfig::default()
})
.states(HashMap::from([("eng-1".to_string(), MemberState::Idle)]))
.build();
daemon.last_auto_dispatch = Instant::now() - Duration::from_secs(30);
daemon.idle_started_at.insert(
"eng-1".to_string(),
Instant::now() - Duration::from_secs(60),
);
daemon
.recent_dispatches
.insert((101, "eng-1".to_string()), Instant::now());
daemon.maybe_auto_dispatch().unwrap();
assert_eq!(
daemon.dispatch_queue.len(),
1,
"dedup_window_secs=0 should effectively disable dedup"
);
assert_eq!(daemon.dispatch_queue[0].task_id, 101);
}
#[test]
fn manual_cooldown_blocks_dispatch() {
let tmp = tempfile::tempdir().unwrap();
write_open_task_file(tmp.path(), 101, "queued-task", "todo");
let members = vec![
manager_member("manager", None),
engineer_member("eng-1", Some("manager"), false),
];
let mut daemon = TestDaemonBuilder::new(tmp.path())
.members(members)
.board(BoardConfig {
auto_dispatch: true,
dispatch_stabilization_delay_secs: 0,
dispatch_manual_cooldown_secs: 30,
..BoardConfig::default()
})
.states(HashMap::from([("eng-1".to_string(), MemberState::Idle)]))
.build();
daemon.last_auto_dispatch = Instant::now() - Duration::from_secs(30);
daemon.idle_started_at.insert(
"eng-1".to_string(),
Instant::now() - Duration::from_secs(60),
);
daemon
.manual_assign_cooldowns
.insert("eng-1".to_string(), Instant::now());
daemon.maybe_auto_dispatch().unwrap();
assert!(
daemon.dispatch_queue.is_empty(),
"engineer within manual cooldown should not be enqueued"
);
}
#[test]
fn manual_cooldown_expires_and_allows_dispatch() {
let tmp = tempfile::tempdir().unwrap();
write_open_task_file(tmp.path(), 101, "queued-task", "todo");
let members = vec![
manager_member("manager", None),
engineer_member("eng-1", Some("manager"), false),
];
let mut daemon = TestDaemonBuilder::new(tmp.path())
.members(members)
.board(BoardConfig {
auto_dispatch: true,
dispatch_stabilization_delay_secs: 0,
dispatch_manual_cooldown_secs: 30,
..BoardConfig::default()
})
.states(HashMap::from([("eng-1".to_string(), MemberState::Idle)]))
.build();
daemon.last_auto_dispatch = Instant::now() - Duration::from_secs(30);
daemon.idle_started_at.insert(
"eng-1".to_string(),
Instant::now() - Duration::from_secs(60),
);
daemon.manual_assign_cooldowns.insert(
"eng-1".to_string(),
Instant::now() - Duration::from_secs(60),
);
daemon.maybe_auto_dispatch().unwrap();
assert_eq!(
daemon.dispatch_queue.len(),
1,
"expired manual cooldown should allow dispatch"
);
assert_eq!(daemon.dispatch_queue[0].task_id, 101);
}
#[test]
fn manual_cooldown_only_affects_assigned_engineer() {
let tmp = tempfile::tempdir().unwrap();
write_open_task_file(tmp.path(), 101, "queued-task", "todo");
let members = vec![
manager_member("manager", None),
engineer_member("eng-1", Some("manager"), false),
engineer_member("eng-2", Some("manager"), false),
];
let mut daemon = TestDaemonBuilder::new(tmp.path())
.members(members)
.board(BoardConfig {
auto_dispatch: true,
dispatch_stabilization_delay_secs: 0,
dispatch_manual_cooldown_secs: 30,
..BoardConfig::default()
})
.states(HashMap::from([
("eng-1".to_string(), MemberState::Idle),
("eng-2".to_string(), MemberState::Idle),
]))
.build();
daemon.last_auto_dispatch = Instant::now() - Duration::from_secs(30);
daemon.idle_started_at.insert(
"eng-1".to_string(),
Instant::now() - Duration::from_secs(60),
);
daemon.idle_started_at.insert(
"eng-2".to_string(),
Instant::now() - Duration::from_secs(60),
);
daemon
.manual_assign_cooldowns
.insert("eng-1".to_string(), Instant::now());
daemon.maybe_auto_dispatch().unwrap();
assert_eq!(
daemon.dispatch_queue.len(),
1,
"only the non-cooldown engineer should be enqueued"
);
assert_eq!(daemon.dispatch_queue[0].engineer, "eng-2");
}
#[test]
fn scored_dispatch_prefers_engineer_with_matching_tag_history() {
let tmp = tempfile::tempdir().unwrap();
seed_dispatch_telemetry(tmp.path(), "eng-1", 0.2);
seed_dispatch_telemetry(tmp.path(), "eng-2", 0.2);
write_task_file(
tmp.path(),
"001-history.md",
"---\nid: 1\ntitle: prior dispatch work\nstatus: done\npriority: high\nclaimed_by: eng-2\ntags:\n - dispatch\nclass: standard\n---\n\nCompleted dispatch work.\n",
);
write_task_file(
tmp.path(),
"101-new-task.md",
"---\nid: 101\ntitle: new dispatch task\nstatus: todo\npriority: high\ntags:\n - dispatch\nclass: standard\n---\n\nImplement dispatch scoring.\n",
);
let members = vec![
manager_member("manager", None),
engineer_member("eng-1", Some("manager"), false),
engineer_member("eng-2", Some("manager"), false),
];
let mut daemon = TestDaemonBuilder::new(tmp.path())
.members(members)
.workflow_policy(WorkflowPolicy {
allocation: AllocationPolicy {
strategy: AllocationStrategy::Scored,
..AllocationPolicy::default()
},
..WorkflowPolicy::default()
})
.board(BoardConfig {
auto_dispatch: true,
dispatch_stabilization_delay_secs: 0,
..BoardConfig::default()
})
.states(HashMap::from([
("eng-1".to_string(), MemberState::Idle),
("eng-2".to_string(), MemberState::Idle),
]))
.build();
daemon.last_auto_dispatch = Instant::now() - Duration::from_secs(30);
daemon.idle_started_at.insert(
"eng-1".to_string(),
Instant::now() - Duration::from_secs(60),
);
daemon.idle_started_at.insert(
"eng-2".to_string(),
Instant::now() - Duration::from_secs(60),
);
daemon.maybe_auto_dispatch().unwrap();
assert_eq!(daemon.dispatch_queue.len(), 1);
assert_eq!(daemon.dispatch_queue[0].engineer, "eng-2");
assert_eq!(daemon.dispatch_queue[0].task_id, 101);
}
#[test]
fn scored_dispatch_prefers_engineer_with_matching_changed_paths() {
let tmp = tempfile::tempdir().unwrap();
seed_dispatch_telemetry(tmp.path(), "eng-1", 0.2);
seed_dispatch_telemetry(tmp.path(), "eng-2", 0.2);
write_task_file(
tmp.path(),
"001-history.md",
"---\nid: 1\ntitle: prior queue work\nstatus: done\npriority: high\nclaimed_by: eng-2\nchanged_paths:\n - src/team/dispatch/queue.rs\nclass: standard\n---\n\nCompleted queue work.\n",
);
write_task_file(
tmp.path(),
"101-new-task.md",
"---\nid: 101\ntitle: new dispatch task\nstatus: todo\npriority: high\nclass: standard\n---\n\nUpdate src/team/dispatch/mod.rs to wire scoring.\n",
);
let members = vec![
manager_member("manager", None),
engineer_member("eng-1", Some("manager"), false),
engineer_member("eng-2", Some("manager"), false),
];
let mut daemon = TestDaemonBuilder::new(tmp.path())
.members(members)
.workflow_policy(WorkflowPolicy {
allocation: AllocationPolicy {
strategy: AllocationStrategy::Scored,
..AllocationPolicy::default()
},
..WorkflowPolicy::default()
})
.board(BoardConfig {
auto_dispatch: true,
dispatch_stabilization_delay_secs: 0,
..BoardConfig::default()
})
.states(HashMap::from([
("eng-1".to_string(), MemberState::Idle),
("eng-2".to_string(), MemberState::Idle),
]))
.build();
daemon.last_auto_dispatch = Instant::now() - Duration::from_secs(30);
daemon.idle_started_at.insert(
"eng-1".to_string(),
Instant::now() - Duration::from_secs(60),
);
daemon.idle_started_at.insert(
"eng-2".to_string(),
Instant::now() - Duration::from_secs(60),
);
daemon.maybe_auto_dispatch().unwrap();
assert_eq!(daemon.dispatch_queue.len(), 1);
assert_eq!(daemon.dispatch_queue[0].engineer, "eng-2");
}
#[test]
fn round_robin_allocation_ignores_profile_scoring() {
let tmp = tempfile::tempdir().unwrap();
write_task_file(
tmp.path(),
"001-history.md",
"---\nid: 1\ntitle: prior dispatch work\nstatus: done\npriority: high\nclaimed_by: eng-2\ntags:\n - dispatch\nclass: standard\n---\n\nCompleted dispatch work.\n",
);
write_task_file(
tmp.path(),
"101-new-task.md",
"---\nid: 101\ntitle: new dispatch task\nstatus: todo\npriority: high\ntags:\n - dispatch\nclass: standard\n---\n\nImplement dispatch scoring.\n",
);
let members = vec![
manager_member("manager", None),
engineer_member("eng-1", Some("manager"), false),
engineer_member("eng-2", Some("manager"), false),
];
let mut daemon = TestDaemonBuilder::new(tmp.path())
.members(members)
.workflow_policy(WorkflowPolicy {
allocation: AllocationPolicy {
strategy: AllocationStrategy::RoundRobin,
..AllocationPolicy::default()
},
..WorkflowPolicy::default()
})
.board(BoardConfig {
auto_dispatch: true,
dispatch_stabilization_delay_secs: 0,
..BoardConfig::default()
})
.states(HashMap::from([
("eng-1".to_string(), MemberState::Idle),
("eng-2".to_string(), MemberState::Idle),
]))
.build();
daemon.last_auto_dispatch = Instant::now() - Duration::from_secs(30);
daemon.idle_started_at.insert(
"eng-1".to_string(),
Instant::now() - Duration::from_secs(60),
);
daemon.idle_started_at.insert(
"eng-2".to_string(),
Instant::now() - Duration::from_secs(60),
);
daemon.maybe_auto_dispatch().unwrap();
assert_eq!(daemon.dispatch_queue.len(), 1);
assert_eq!(daemon.dispatch_queue[0].engineer, "eng-1");
}
#[test]
fn dispatch_queue_prefers_original_owner_for_verification_retry_rework() {
let tmp = tempfile::tempdir().unwrap();
write_task_file(
tmp.path(),
"042-retry.md",
"---\nid: 42\ntitle: retry rework\nstatus: in-progress\npriority: high\nclaimed_by: eng-1\nclass: standard\ntests_passed: false\noutcome: verification_retry_required\nartifacts:\n - artifacts/retry-42.log\n---\n\nFix failed verification.\n",
);
let members = vec![
manager_member("manager", None),
engineer_member("eng-1", Some("manager"), false),
engineer_member("eng-2", Some("manager"), false),
];
let mut daemon = TestDaemonBuilder::new(tmp.path())
.members(members)
.workflow_policy(WorkflowPolicy {
allocation: AllocationPolicy {
strategy: AllocationStrategy::RoundRobin,
..AllocationPolicy::default()
},
..WorkflowPolicy::default()
})
.board(BoardConfig {
auto_dispatch: true,
dispatch_stabilization_delay_secs: 0,
..BoardConfig::default()
})
.states(HashMap::from([
("eng-1".to_string(), MemberState::Idle),
("eng-2".to_string(), MemberState::Idle),
]))
.build();
daemon.last_auto_dispatch = Instant::now() - Duration::from_secs(30);
daemon.idle_started_at.insert(
"eng-1".to_string(),
Instant::now() - Duration::from_secs(60),
);
daemon.idle_started_at.insert(
"eng-2".to_string(),
Instant::now() - Duration::from_secs(60),
);
daemon.enqueue_dispatch_candidates().unwrap();
assert_eq!(daemon.dispatch_queue.len(), 1);
assert_eq!(daemon.dispatch_queue[0].task_id, 42);
assert_eq!(daemon.dispatch_queue[0].engineer, "eng-1");
}
#[test]
fn dispatch_queue_ignores_done_verification_retry_metadata() {
let tmp = tempfile::tempdir().unwrap();
write_task_file(
tmp.path(),
"042-done-retry.md",
"---\nid: 42\ntitle: done retry metadata\nstatus: done\npriority: high\nclaimed_by: eng-1\nclass: standard\ntests_passed: false\noutcome: verification_retry_required\nartifacts:\n - artifacts/retry-42.log\n---\n\nAlready merged.\n",
);
write_open_task_file(tmp.path(), 101, "normal-work", "todo");
let members = vec![
manager_member("manager", None),
engineer_member("eng-1", Some("manager"), false),
engineer_member("eng-2", Some("manager"), false),
];
let mut daemon = TestDaemonBuilder::new(tmp.path())
.members(members)
.workflow_policy(WorkflowPolicy {
allocation: AllocationPolicy {
strategy: AllocationStrategy::RoundRobin,
..AllocationPolicy::default()
},
..WorkflowPolicy::default()
})
.board(BoardConfig {
auto_dispatch: true,
dispatch_stabilization_delay_secs: 0,
..BoardConfig::default()
})
.states(HashMap::from([
("eng-1".to_string(), MemberState::Idle),
("eng-2".to_string(), MemberState::Idle),
]))
.build();
daemon.last_auto_dispatch = Instant::now() - Duration::from_secs(30);
daemon.idle_started_at.insert(
"eng-1".to_string(),
Instant::now() - Duration::from_secs(60),
);
daemon.idle_started_at.insert(
"eng-2".to_string(),
Instant::now() - Duration::from_secs(60),
);
daemon.enqueue_dispatch_candidates().unwrap();
assert_eq!(daemon.dispatch_queue.len(), 1);
assert_eq!(daemon.dispatch_queue[0].task_id, 101);
}
#[test]
fn dispatch_queue_sends_idle_peer_to_normal_work_before_retry_rework() {
let tmp = tempfile::tempdir().unwrap();
write_task_file(
tmp.path(),
"042-retry.md",
"---\nid: 42\ntitle: retry rework\nstatus: in-progress\npriority: high\nclaimed_by: eng-1\nclass: standard\ntests_passed: false\noutcome: verification_retry_required\nartifacts:\n - artifacts/retry-42.log\n---\n\nFix failed verification.\n",
);
write_open_task_file(tmp.path(), 101, "normal-work", "todo");
write_owned_task_file(tmp.path(), 777, "busy-owner", "in-progress", "eng-1");
let members = vec![
manager_member("manager", None),
engineer_member("eng-1", Some("manager"), false),
engineer_member("eng-2", Some("manager"), false),
];
let mut daemon = TestDaemonBuilder::new(tmp.path())
.members(members)
.workflow_policy(WorkflowPolicy {
allocation: AllocationPolicy {
strategy: AllocationStrategy::RoundRobin,
..AllocationPolicy::default()
},
..WorkflowPolicy::default()
})
.board(BoardConfig {
auto_dispatch: true,
dispatch_stabilization_delay_secs: 0,
..BoardConfig::default()
})
.states(HashMap::from([
("eng-1".to_string(), MemberState::Working),
("eng-2".to_string(), MemberState::Idle),
]))
.build();
daemon.active_tasks.insert("eng-1".to_string(), 777);
daemon.last_auto_dispatch = Instant::now() - Duration::from_secs(30);
daemon.idle_started_at.insert(
"eng-2".to_string(),
Instant::now() - Duration::from_secs(60),
);
daemon.enqueue_dispatch_candidates().unwrap();
assert_eq!(daemon.dispatch_queue.len(), 1);
assert_eq!(daemon.dispatch_queue[0].task_id, 101);
assert_eq!(daemon.dispatch_queue[0].engineer, "eng-2");
}
#[test]
fn priority_rank_critical() {
assert_eq!(super::dispatch_priority_rank("critical"), 0);
}
#[test]
fn priority_rank_high() {
assert_eq!(super::dispatch_priority_rank("high"), 1);
}
#[test]
fn priority_rank_medium() {
assert_eq!(super::dispatch_priority_rank("medium"), 2);
}
#[test]
fn priority_rank_low() {
assert_eq!(super::dispatch_priority_rank("low"), 3);
}
#[test]
fn priority_rank_unknown_defaults_highest_number() {
assert_eq!(super::dispatch_priority_rank(""), 4);
assert_eq!(super::dispatch_priority_rank("urgent"), 4);
}
#[test]
fn parse_task_id_standard_format() {
assert_eq!(
super::parse_assignment_task_id("Task #42: do stuff"),
Some(42)
);
}
#[test]
fn parse_task_id_no_match() {
assert_eq!(super::parse_assignment_task_id("just a message"), None);
}
#[test]
fn parse_task_id_case_insensitive() {
assert_eq!(
super::parse_assignment_task_id("TASK #99: uppercase"),
Some(99)
);
assert_eq!(
super::parse_assignment_task_id("task #77: lowercase"),
Some(77)
);
}
#[test]
fn parse_task_id_multiple_returns_first() {
assert_eq!(
super::parse_assignment_task_id("Task #10: depends on Task #5"),
Some(10)
);
}
#[test]
fn parse_task_id_empty_digits() {
assert_eq!(super::parse_assignment_task_id("Task #: no digits"), None);
}
#[test]
fn slugify_basic_text() {
assert_eq!(
super::slugify_task_branch("Fix castling rights"),
"fix-castling-rights"
);
}
#[test]
fn slugify_special_characters() {
assert_eq!(
super::slugify_task_branch("Add feature (v2) — fix!"),
"add-feature-v2-fix"
);
}
#[test]
fn slugify_empty_returns_task() {
assert_eq!(super::slugify_task_branch(""), "task");
}
#[test]
fn slugify_only_special_chars() {
assert_eq!(super::slugify_task_branch("--- !!!"), "task");
}
#[test]
fn slugify_preserves_numbers() {
assert_eq!(super::slugify_task_branch("Task 42 fix"), "task-42-fix");
}
#[test]
fn summarize_empty_body() {
assert_eq!(super::summarize_assignment(""), "task");
}
#[test]
fn summarize_only_whitespace() {
assert_eq!(super::summarize_assignment(" \n \n "), "task");
}
#[test]
fn summarize_truncates_long_line() {
let long = "a".repeat(200);
let result = super::summarize_assignment(&long);
assert_eq!(result.len(), 120);
assert!(result.ends_with("..."));
}
#[test]
fn summarize_preserves_short_line() {
assert_eq!(super::summarize_assignment("Short task"), "Short task");
}
#[test]
fn dispatch_queue_entry_serde_roundtrip() {
use super::DispatchQueueEntry;
let entry = DispatchQueueEntry {
engineer: "eng-1".to_string(),
task_id: 42,
task_title: "Test task".to_string(),
queued_at: 1234567890,
validation_failures: 2,
last_failure: Some("worktree dirty".to_string()),
};
let json = serde_json::to_string(&entry).unwrap();
let restored: DispatchQueueEntry = serde_json::from_str(&json).unwrap();
assert_eq!(restored, entry);
}
#[test]
fn dispatch_queue_entry_serde_no_failure() {
use super::DispatchQueueEntry;
let entry = DispatchQueueEntry {
engineer: "eng-2".to_string(),
task_id: 1,
task_title: "Clean task".to_string(),
queued_at: 0,
validation_failures: 0,
last_failure: None,
};
let json = serde_json::to_string(&entry).unwrap();
let restored: DispatchQueueEntry = serde_json::from_str(&json).unwrap();
assert_eq!(restored, entry);
}