use std::collections::{BTreeSet, HashMap, HashSet};
use std::path::Path;
use std::time::Instant;
use anyhow::Result;
use regex::Regex;
use tracing::{debug, info, warn};
use super::super::super::policy::check_wip_limit;
use super::super::super::task_loop::engineer_worktree_ready_for_dispatch_from_trunk;
use super::super::task_cmd::{
StatusTransitionAttribution, append_task_dependencies, assign_task_owners,
transition_task_with_attribution,
};
use super::super::*;
use crate::team::allocation::{
EngineerProfile, load_engineer_profiles, predict_task_file_paths, rank_engineers_for_task,
};
use crate::team::config::AllocationStrategy;
use serde::Deserialize;
fn split_acyclic_blocking_ids(
board_dir: &Path,
candidate_id: u32,
blocking_task_ids: &[u32],
) -> Result<(Vec<u32>, Vec<u32>)> {
let tasks = crate::task::load_tasks_from_dir(&board_dir.join("tasks"))?;
let deps_by_id: HashMap<u32, Vec<u32>> = tasks
.iter()
.map(|task| (task.id, task.depends_on.clone()))
.collect();
let reaches_candidate = |start: u32| -> bool {
if start == candidate_id {
return true;
}
let mut visited: HashSet<u32> = HashSet::new();
let mut stack: Vec<u32> = vec![start];
while let Some(node) = stack.pop() {
if !visited.insert(node) {
continue;
}
let Some(neighbors) = deps_by_id.get(&node) else {
continue;
};
for next in neighbors {
if *next == candidate_id {
return true;
}
if !visited.contains(next) {
stack.push(*next);
}
}
}
false
};
let mut safe = Vec::new();
let mut rejected = Vec::new();
for id in blocking_task_ids {
if reaches_candidate(*id) {
rejected.push(*id);
} else {
safe.push(*id);
}
}
Ok((safe, rejected))
}
const NON_ROLE_HYPHEN_TOKENS: &[&str] = &[
"role-flexible",
"strategic-analysis",
"narrative-audit",
"north-star",
"any-engineer",
];
fn role_name_seed_tags(role_name: &str) -> Vec<String> {
let mut seeds = vec![role_name.to_string()];
let Some(suffix) = role_name.rsplit('-').next() else {
return seeds;
};
if suffix == role_name || suffix.is_empty() {
return seeds;
}
seeds.push(suffix.to_string());
if let Some(stem) = suffix.strip_suffix("er")
&& stem.len() >= 3
{
seeds.push(stem.to_string());
seeds.push(format!("{stem}ing"));
}
seeds
}
fn first_role_token_after(text: &str) -> Option<String> {
let bytes = text.as_bytes();
let mut i = 0;
while i < bytes.len() {
while i < bytes.len() && !bytes[i].is_ascii_lowercase() {
i += 1;
}
if i >= bytes.len() {
break;
}
let start = i;
while i < bytes.len() && (bytes[i].is_ascii_lowercase() || bytes[i] == b'-') {
i += 1;
}
let token = text[start..i].trim_end_matches('-');
if token.contains('-')
&& !token.starts_with('-')
&& !NON_ROLE_HYPHEN_TOKENS.contains(&token)
{
return Some(token.to_string());
}
}
None
}
pub(crate) fn parse_body_owner_role(body: &str) -> Option<String> {
const ROUTING_CUES: &[&str] = &[
"Owner:",
"OWNER:",
"Route:",
"Primary:",
"Owner routing",
"route to",
"dispatch to",
"assign to",
];
for line in body.lines() {
for cue in ROUTING_CUES {
let mut search_from = 0;
while let Some(rel_idx) = line[search_from..].find(cue) {
let idx = search_from + rel_idx;
search_from = idx + cue.len();
let boundary_ok = idx == 0
|| line[..idx]
.chars()
.next_back()
.map(|c| {
c.is_whitespace()
|| matches!(c, '-' | '*' | '_' | '(' | '[' | '.' | ';' | ':' | ',')
})
.unwrap_or(true);
if !boundary_ok {
continue;
}
if let Some(role) = first_role_token_after(&line[search_from..]) {
return Some(role);
}
}
}
}
None
}
pub(crate) fn parse_body_dependency_ids(body: &str) -> Option<Vec<u32>> {
for line in body.lines() {
if let Some(trimmed) = body_dependency_reference(line) {
let ids: Vec<u32> = trimmed
.split('#')
.skip(1)
.filter_map(|s| {
s.chars()
.take_while(|c| c.is_ascii_digit())
.collect::<String>()
.parse()
.ok()
})
.collect();
return Some(ids);
}
}
None
}
fn body_dependency_reference(line: &str) -> Option<&str> {
let trimmed = line.trim().trim_start_matches('-').trim();
let lower = trimmed.to_ascii_lowercase();
if lower.starts_with("blocked on:") || lower.starts_with("depends on:") {
Some(trimmed)
} else {
None
}
}
use super::{DISPATCH_QUEUE_FAILURE_LIMIT, DispatchQueueEntry, dispatch_priority_rank};
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct OverlapConflict {
pub task_id: String,
pub conflicting_files: Vec<String>,
pub in_progress_engineer: String,
}
#[derive(Debug, Default, Deserialize)]
struct ChangedPathsFrontmatter {
#[serde(default)]
changed_paths: Vec<String>,
}
fn board_tasks_dir(project_root: &Path) -> std::path::PathBuf {
project_root
.join(".batty")
.join("team_config")
.join("board")
.join("tasks")
}
fn extract_frontmatter(content: &str) -> Option<&str> {
let trimmed = content.trim_start();
if !trimmed.starts_with("---") {
return None;
}
let after_open = trimmed[3..].strip_prefix('\n').unwrap_or(&trimmed[3..]);
let end = after_open.find("\n---")?;
Some(&after_open[..end])
}
fn load_changed_paths(path: &Path) -> Vec<String> {
let Ok(content) = std::fs::read_to_string(path) else {
return Vec::new();
};
let Some(frontmatter) = extract_frontmatter(&content) else {
return Vec::new();
};
serde_yaml::from_str::<ChangedPathsFrontmatter>(frontmatter)
.map(|parsed| parsed.changed_paths)
.unwrap_or_default()
}
fn normalize_predicted_path(path: &str) -> String {
path.trim_matches(|ch: char| matches!(ch, '.' | ',' | ';' | ':' | ')' | ']'))
.to_string()
}
fn has_glob_magic(path: &str) -> bool {
path.contains('*') || path.contains('?')
}
fn glob_to_regex(pattern: &str) -> Option<Regex> {
let mut regex = String::from("^");
let mut chars = pattern.chars().peekable();
while let Some(ch) = chars.next() {
match ch {
'*' => {
if chars.peek() == Some(&'*') {
chars.next();
if chars.peek() == Some(&'/') {
chars.next();
regex.push_str("(?:.*/)?");
} else {
regex.push_str(".*");
}
} else {
regex.push_str("[^/]*");
}
}
'?' => regex.push_str("[^/]"),
'.' | '+' | '(' | ')' | '[' | ']' | '{' | '}' | '^' | '$' | '|' | '\\' => {
regex.push('\\');
regex.push(ch);
}
_ => regex.push(ch),
}
}
regex.push('$');
Regex::new(®ex).ok()
}
fn glob_matches_path(pattern: &str, path: &str) -> bool {
if !has_glob_magic(pattern) {
return pattern == path;
}
glob_to_regex(pattern)
.map(|regex| regex.is_match(path))
.unwrap_or(false)
}
fn glob_literal_prefix(pattern: &str) -> Option<&str> {
let idx = pattern
.char_indices()
.find_map(|(idx, ch)| matches!(ch, '*' | '?').then_some(idx))
.unwrap_or(pattern.len());
let prefix = pattern[..idx].trim_end_matches('/');
(!prefix.is_empty()).then_some(prefix)
}
fn paths_overlap(left: &str, right: &str) -> bool {
match (has_glob_magic(left), has_glob_magic(right)) {
(false, false) => left == right,
(true, false) => glob_matches_path(left, right),
(false, true) => glob_matches_path(right, left),
(true, true) => {
if left == right {
return true;
}
match (glob_literal_prefix(left), glob_literal_prefix(right)) {
(Some(left_prefix), Some(right_prefix)) => {
left_prefix.starts_with(right_prefix) || right_prefix.starts_with(left_prefix)
}
_ => true,
}
}
}
}
fn describe_overlap(left: &str, right: &str) -> String {
match (has_glob_magic(left), has_glob_magic(right)) {
(false, false) => left.to_string(),
(true, false) => right.to_string(),
(false, true) => left.to_string(),
(true, true) if left == right => left.to_string(),
(true, true) => format!("{left} <> {right}"),
}
}
pub fn predicted_files(task: &crate::task::Task, project_root: &Path) -> Vec<String> {
let mut paths = predict_task_file_paths(project_root, task)
.unwrap_or_default()
.into_iter()
.map(|path| normalize_predicted_path(&path))
.collect::<Vec<_>>();
if let Ok(tasks) = crate::task::load_tasks_from_dir(&board_tasks_dir(project_root)) {
for historical in tasks {
if historical.id == task.id || historical.tags.is_empty() {
continue;
}
if !task
.tags
.iter()
.any(|tag| historical.tags.iter().any(|candidate| candidate == tag))
{
continue;
}
paths.extend(
load_changed_paths(historical.source_path.as_path())
.into_iter()
.map(|path| normalize_predicted_path(&path)),
);
}
}
paths.retain(|path| !path.is_empty());
paths.sort();
paths.dedup();
paths
}
fn overlapping_files(candidate_paths: &[String], active_paths: &[String]) -> Vec<String> {
let mut overlaps = BTreeSet::new();
for candidate in candidate_paths {
for active in active_paths {
if paths_overlap(candidate, active) {
overlaps.insert(describe_overlap(candidate, active));
}
}
}
overlaps.into_iter().collect()
}
pub fn find_overlapping_tasks(
candidate: &crate::task::Task,
in_progress: &[crate::task::Task],
project_root: &Path,
) -> Vec<OverlapConflict> {
let candidate_paths = predicted_files(candidate, project_root);
let mut conflicts = Vec::new();
for active_task in in_progress {
if active_task.id == candidate.id {
continue;
}
let active_paths = predicted_files(active_task, project_root);
let conflicting_files = overlapping_files(&candidate_paths, &active_paths);
if conflicting_files.is_empty() {
continue;
}
conflicts.push(OverlapConflict {
task_id: active_task.id.to_string(),
conflicting_files,
in_progress_engineer: active_task
.claimed_by
.clone()
.unwrap_or_else(|| "unknown".to_string()),
});
}
conflicts.sort_by(|left, right| left.task_id.cmp(&right.task_id));
conflicts
}
fn available_dispatch_tasks(
board_dir: &Path,
queued_task_ids: &HashSet<u32>,
excluded_tags: &[String],
non_engineer_assignees: &HashSet<String>,
rescued_task_ids: &HashSet<u32>,
verification_retry_task_ids: &HashSet<u32>,
) -> Result<Vec<crate::task::Task>> {
let tasks = crate::task::load_tasks_from_dir(&board_dir.join("tasks"))?;
let task_status_by_id: HashMap<u32, String> = tasks
.iter()
.map(|task| (task.id, task.status.clone()))
.collect();
let mut available: Vec<crate::task::Task> = tasks
.into_iter()
.filter(|task| {
matches!(task.status.as_str(), "backlog" | "todo")
|| verification_retry_task_ids.contains(&task.id)
})
.filter(|task| task.claimed_by.is_none() || verification_retry_task_ids.contains(&task.id))
.filter(|task| task.blocked.is_none())
.filter(|task| task.blocked_on.is_none())
.filter(|task| !task.is_schedule_blocked())
.filter(|task| !queued_task_ids.contains(&task.id))
.filter(|task| !task_has_excluded_tag(task, excluded_tags))
.filter(|task| !rescued_task_ids.contains(&task.id))
.filter(|task| {
task.assignee
.as_deref()
.is_none_or(|name| !non_engineer_assignees.contains(name))
})
.filter(|task| {
if task.assignee.is_some() {
return true;
}
match parse_body_owner_role(&task.description) {
Some(owner) => !non_engineer_assignees.contains(&owner),
None => true,
}
})
.filter(|task| {
task.depends_on.iter().all(|dep_id| {
task_status_by_id
.get(dep_id)
.is_none_or(|status| dep_status_satisfied(status))
})
})
.filter(|task| body_dependencies_satisfied(task, &task_status_by_id))
.collect();
available.sort_by_key(|task| (dispatch_priority_rank(&task.priority), task.id));
Ok(available)
}
fn dep_status_satisfied(status: &str) -> bool {
matches!(status, "done" | "archived")
}
fn body_dependencies_satisfied(
task: &crate::task::Task,
task_status_by_id: &HashMap<u32, String>,
) -> bool {
let Some(blocked_ids) = parse_body_dependency_ids(&task.description) else {
return true;
};
!blocked_ids.is_empty()
&& blocked_ids.iter().all(|dep_id| {
task_status_by_id
.get(dep_id)
.is_some_and(|status| dep_status_satisfied(status))
})
}
fn task_has_excluded_tag(task: &crate::task::Task, excluded_tags: &[String]) -> bool {
if excluded_tags.is_empty() {
return false;
}
task.tags.iter().any(|task_tag| {
excluded_tags
.iter()
.any(|excluded| excluded.eq_ignore_ascii_case(task_tag))
})
}
fn verification_retry_required_metadata(
task: &crate::task::Task,
) -> Option<crate::team::board::WorkflowMetadata> {
let metadata = crate::team::board::read_workflow_metadata(&task.source_path).ok()?;
(metadata.tests_passed == Some(false)
&& metadata.outcome.as_deref() == Some("verification_retry_required")
&& !metadata.artifacts.is_empty())
.then_some(metadata)
}
fn verification_retry_assignment_context(task: &crate::task::Task) -> Option<String> {
let metadata = verification_retry_required_metadata(task)?;
let mut lines = vec![
"Verification retry required.".to_string(),
"Outcome: verification_retry_required.".to_string(),
];
if let Some(owner) = task.claimed_by.as_deref() {
lines.push(format!("Previous owner: {owner}."));
}
if let Some(artifact) = metadata.artifacts.last() {
lines.push(format!("Latest verification artifact: {artifact}."));
}
Some(lines.join("\n"))
}
impl TeamDaemon {
fn verification_retry_dispatchable_task_ids(
&self,
board_dir: &Path,
allow_peer_pickup: bool,
) -> Result<HashSet<u32>> {
Ok(crate::task::load_tasks_from_dir(&board_dir.join("tasks"))?
.into_iter()
.filter(|task| self.verification_retry_dispatchable_task(task, allow_peer_pickup))
.map(|task| task.id)
.collect())
}
fn verification_retry_dispatchable_task(
&self,
task: &crate::task::Task,
allow_peer_pickup: bool,
) -> bool {
if matches!(task.status.as_str(), "done" | "archived") {
return false;
}
if verification_retry_required_metadata(task).is_none() {
return false;
}
let Some(owner) = task.claimed_by.as_deref() else {
return true;
};
if self.active_tasks.get(owner) == Some(&task.id)
&& self.states.get(owner) == Some(&MemberState::Working)
{
return false;
}
self.idle_engineer_names()
.iter()
.any(|engineer| engineer == owner)
|| allow_peer_pickup
}
fn serialize_overlapping_candidate(
&mut self,
board_dir: &Path,
candidate: &crate::task::Task,
conflicts: &[OverlapConflict],
persist_dependency: bool,
) -> Result<bool> {
if conflicts.is_empty() {
return Ok(false);
}
let mut blocking_task_ids: Vec<u32> = conflicts
.iter()
.filter_map(|conflict| conflict.task_id.parse::<u32>().ok())
.collect();
blocking_task_ids.sort_unstable();
blocking_task_ids.dedup();
let overlap_details = conflicts
.iter()
.map(|conflict| {
format!(
"#{} [{}]",
conflict.task_id,
conflict.conflicting_files.join(", ")
)
})
.collect::<Vec<_>>();
let (safe_blocking_ids, rejected_blocking_ids) = if persist_dependency {
split_acyclic_blocking_ids(board_dir, candidate.id, &blocking_task_ids)?
} else {
(blocking_task_ids.clone(), Vec::new())
};
if !rejected_blocking_ids.is_empty() {
warn!(
task_id = candidate.id,
rejected = ?rejected_blocking_ids,
"dispatch queue: skipped cycle-creating overlap dependency"
);
}
let updated_dependencies = if persist_dependency && !safe_blocking_ids.is_empty() {
Some(append_task_dependencies(
board_dir,
candidate.id,
&safe_blocking_ids,
)?)
} else {
None
};
let details = if persist_dependency {
format!(
"serialized task #{} behind {} due to predicted file overlap",
candidate.id,
overlap_details.join("; ")
)
} else {
format!(
"deferred task #{} in file_lock_wait behind {} due to predicted file overlap",
candidate.id,
overlap_details.join("; ")
)
};
self.emit_event(TeamEvent::dispatch_overlap_prevented(
candidate.id,
&blocking_task_ids,
&details,
));
self.record_orchestrator_action(format!("dispatch overlap: {details}"));
info!(
task_id = candidate.id,
blocking = ?updated_dependencies,
persist_dependency,
"dispatch queue: prevented overlapping dispatch"
);
Ok(true)
}
pub(in super::super) fn idle_engineer_names(&self) -> Vec<String> {
self.config
.members
.iter()
.filter(|member| member.role_type == RoleType::Engineer)
.filter(|member| {
let state = self.states.get(&member.name);
match state {
Some(&MemberState::Idle) => true,
Some(&MemberState::Working) => !self.active_tasks.contains_key(&member.name),
_ => false,
}
})
.map(|member| member.name.clone())
.collect()
}
#[cfg_attr(not(test), allow(dead_code))]
fn next_dispatch_task(
&self,
board_dir: &Path,
queued_task_ids: &HashSet<u32>,
) -> Result<Option<crate::task::Task>> {
let normal_available = available_dispatch_tasks(
board_dir,
queued_task_ids,
&self.config.team_config.board.dispatch_excluded_tags,
&self.non_engineer_member_names(),
&self.rescued_task_ids(),
&HashSet::new(),
)?;
let allow_peer_pickup =
normal_available.is_empty() && !self.idle_engineer_names().is_empty();
Ok(available_dispatch_tasks(
board_dir,
queued_task_ids,
&self.config.team_config.board.dispatch_excluded_tags,
&self.non_engineer_member_names(),
&self.rescued_task_ids(),
&self.verification_retry_dispatchable_task_ids(board_dir, allow_peer_pickup)?,
)?
.into_iter()
.next())
}
pub(super) fn rescued_task_ids(&self) -> HashSet<u32> {
let base = Duration::from_secs(self.config.team_config.board.orphan_rescue_cooldown_secs);
self.recently_rescued_tasks
.iter()
.filter(|(_, record)| record.dispatch_blocked(base))
.map(|(task_id, _)| *task_id)
.collect()
}
pub(in super::super) fn record_task_rescue(&mut self, task_id: u32) {
let base = Duration::from_secs(self.config.team_config.board.orphan_rescue_cooldown_secs);
let now = Instant::now();
self.recently_rescued_tasks
.entry(task_id)
.and_modify(|record| {
if record.in_cascade_window(base) {
record.count = record.count.saturating_add(1);
} else {
record.count = 1;
}
record.last_rescued_at = now;
})
.or_insert(crate::team::daemon::RescueRecord {
last_rescued_at: now,
count: 1,
});
}
pub(in super::super) fn release_exclusion_window(&self) -> Duration {
Duration::from_secs(
self.config
.team_config
.board
.dispatch_release_exclusion_secs,
)
}
pub(in super::super) fn record_task_release_by(&mut self, task_id: u32, engineer: &str) {
let base = self.release_exclusion_window();
let now = Instant::now();
self.recently_released_by
.entry((task_id, engineer.to_string()))
.and_modify(|record| {
if record.in_cascade_window(base) {
record.count = record.count.saturating_add(1);
} else {
record.count = 1;
}
record.last_released_at = now;
})
.or_insert(crate::team::daemon::ReleaseRecord {
last_released_at: now,
count: 1,
});
}
pub(in super::super) fn is_release_excluded(&self, task_id: u32, engineer: &str) -> bool {
let base = self.release_exclusion_window();
self.recently_released_by
.get(&(task_id, engineer.to_string()))
.map(|record| record.dispatch_excluded(base))
.unwrap_or(false)
}
fn non_engineer_member_names(&self) -> HashSet<String> {
self.config
.members
.iter()
.filter(|member| member.role_type != RoleType::Engineer)
.map(|member| member.name.clone())
.collect()
}
#[cfg(test)]
pub(super) fn test_next_dispatch_task(
&self,
board_dir: &std::path::Path,
queued: &HashSet<u32>,
) -> Result<Option<crate::task::Task>> {
self.next_dispatch_task(board_dir, queued)
}
pub(in super::super) fn enqueue_dispatch_candidates(&mut self) -> Result<()> {
let board_dir = self.board_dir();
let board_tasks = crate::task::load_tasks_from_dir(&board_dir.join("tasks"))?;
let benched_engineers = crate::team::bench::benched_engineer_names(self.project_root())?;
let dedup_window =
Duration::from_secs(self.config.team_config.board.dispatch_dedup_window_secs);
self.recent_dispatches
.retain(|_, dispatched_at| dispatched_at.elapsed() < dedup_window);
let release_base_window = self.release_exclusion_window();
self.recently_released_by
.retain(|_, record| record.in_cascade_window(release_base_window));
let rescue_base_cooldown =
Duration::from_secs(self.config.team_config.board.orphan_rescue_cooldown_secs);
self.recently_rescued_tasks
.retain(|_, record| record.in_cascade_window(rescue_base_cooldown));
let rescued_task_ids: HashSet<u32> = self
.recently_rescued_tasks
.iter()
.filter(|(_, record)| record.dispatch_blocked(rescue_base_cooldown))
.map(|(task_id, _)| *task_id)
.collect();
let mut queued_task_ids: HashSet<u32> = self
.dispatch_queue
.iter()
.map(|entry| entry.task_id)
.collect();
let mut queued_engineers: HashSet<String> = self
.dispatch_queue
.iter()
.map(|entry| entry.engineer.clone())
.collect();
let mut file_locked_task_ids = HashSet::new();
let manual_cooldown =
Duration::from_secs(self.config.team_config.board.dispatch_manual_cooldown_secs);
let all_engineers: Vec<String> = self
.config
.members
.iter()
.filter(|member| member.role_type == RoleType::Engineer)
.map(|member| member.name.clone())
.collect();
let non_engineer_names = self.non_engineer_member_names();
let mut profiles =
load_engineer_profiles(self.project_root(), &all_engineers, &board_tasks)?;
for member in &self.config.members {
if member.role_type != RoleType::Engineer {
continue;
}
if let Some(profile) = profiles.get_mut(&member.name) {
for seed in role_name_seed_tags(&member.role_name) {
profile.domain_tags.insert(seed);
}
}
}
let allow_peer_retry_pickup = available_dispatch_tasks(
&board_dir,
&queued_task_ids,
&self.config.team_config.board.dispatch_excluded_tags,
&non_engineer_names,
&rescued_task_ids,
&HashSet::new(),
)?
.is_empty()
&& !self.idle_engineer_names().is_empty();
let mut eligibility_excluded_task_ids: HashSet<u32> = HashSet::new();
loop {
let mut unavailable_task_ids = queued_task_ids.clone();
unavailable_task_ids.extend(file_locked_task_ids.iter().copied());
unavailable_task_ids.extend(eligibility_excluded_task_ids.iter().copied());
let verification_retry_task_ids =
self.verification_retry_dispatchable_task_ids(&board_dir, allow_peer_retry_pickup)?;
let available_tasks = available_dispatch_tasks(
&board_dir,
&unavailable_task_ids,
&self.config.team_config.board.dispatch_excluded_tags,
&non_engineer_names,
&rescued_task_ids,
&verification_retry_task_ids,
)?;
if available_tasks.is_empty() {
break;
}
let in_progress_tasks: Vec<crate::task::Task> =
crate::task::load_tasks_from_dir(&board_dir.join("tasks"))?
.into_iter()
.filter(|task| task.status == "in-progress")
.collect();
let mut selected_task = None;
let mut least_conflicted: Option<(crate::task::Task, Vec<OverlapConflict>)> = None;
let file_level_locks_enabled = self.config.team_config.workflow_policy.file_level_locks;
let all_engineers_use_worktrees = self
.config
.team_config
.roles
.iter()
.filter(|r| r.role_type == crate::team::config::RoleType::Engineer)
.all(|r| r.use_worktrees);
let skip_overlap_checks = all_engineers_use_worktrees && !file_level_locks_enabled;
for task in available_tasks {
if skip_overlap_checks {
selected_task = Some(task);
break;
}
let conflicts =
find_overlapping_tasks(&task, &in_progress_tasks, self.project_root());
if conflicts.is_empty() {
selected_task = Some(task);
break;
}
for conflict in &conflicts {
self.emit_event(TeamEvent::dispatch_overlap_skipped(
task.id,
&conflict.task_id,
&conflict.conflicting_files,
));
}
if file_level_locks_enabled {
self.serialize_overlapping_candidate(&board_dir, &task, &conflicts, false)?;
file_locked_task_ids.insert(task.id);
continue;
}
let replace = least_conflicted
.as_ref()
.is_none_or(|(_, existing)| conflicts.len() < existing.len());
if replace {
least_conflicted = Some((task, conflicts));
}
}
let task = if let Some(task) = selected_task {
task
} else if let Some((task, conflicts)) = least_conflicted {
self.serialize_overlapping_candidate(&board_dir, &task, &conflicts, true)?;
continue;
} else {
break;
};
let ranked_engineers = self.rank_dispatch_engineers(
&task,
&queued_engineers,
&benched_engineers,
manual_cooldown,
&profiles,
);
let retry_previous_owner = self
.verification_retry_dispatchable_task(&task, allow_peer_retry_pickup)
.then(|| task.claimed_by.clone())
.flatten();
let mut ranked_engineers = ranked_engineers;
if let Some(owner) = retry_previous_owner.as_deref() {
ranked_engineers
.retain(|engineer_name| engineer_name == owner || allow_peer_retry_pickup);
if let Some(index) = ranked_engineers
.iter()
.position(|engineer_name| engineer_name == owner)
{
let owner = ranked_engineers.remove(index);
ranked_engineers.insert(0, owner);
}
}
let Some(engineer_name) = ranked_engineers.into_iter().find(|engineer_name| {
!self
.recent_dispatches
.contains_key(&(task.id, engineer_name.clone()))
&& !self.is_release_excluded(task.id, engineer_name)
}) else {
eligibility_excluded_task_ids.insert(task.id);
continue;
};
queued_task_ids.insert(task.id);
queued_engineers.insert(engineer_name.clone());
self.dispatch_queue.push(DispatchQueueEntry {
engineer: engineer_name,
task_id: task.id,
task_title: task.title,
queued_at: now_unix(),
validation_failures: 0,
last_failure: None,
});
}
Ok(())
}
fn task_for_dispatch_entry(
&self,
board_dir: &Path,
entry: &DispatchQueueEntry,
) -> Result<Option<crate::task::Task>> {
let tasks = crate::task::load_tasks_from_dir(&board_dir.join("tasks"))?;
let task_status_by_id: HashMap<u32, String> = tasks
.iter()
.map(|task| (task.id, task.status.clone()))
.collect();
Ok(tasks.into_iter().find(|task| {
let retry_dispatchable = self.verification_retry_dispatchable_task(task, true);
task.id == entry.task_id
&& (matches!(task.status.as_str(), "backlog" | "todo") || retry_dispatchable)
&& (task.claimed_by.is_none() || retry_dispatchable)
&& task.blocked.is_none()
&& task.blocked_on.is_none()
&& !task.is_schedule_blocked()
&& task.depends_on.iter().all(|dep_id| {
task_status_by_id
.get(dep_id)
.is_none_or(|status| dep_status_satisfied(status))
})
&& body_dependencies_satisfied(task, &task_status_by_id)
}))
}
pub(in super::super) fn process_dispatch_queue(&mut self) -> Result<()> {
self.reconcile_active_tasks()?;
let board_dir = self.board_dir();
let benched_engineers = crate::team::bench::benched_engineer_names(self.project_root())?;
let rescued_task_ids = self.rescued_task_ids();
let mut pending: Vec<DispatchQueueEntry> = std::mem::take(&mut self.dispatch_queue);
let mut retained = Vec::new();
for mut entry in pending.drain(..) {
let task_still_dispatchable =
self.task_for_dispatch_entry(&board_dir, &entry)?.is_some();
if !task_still_dispatchable {
debug!(
engineer = %entry.engineer,
task_id = entry.task_id,
"dispatch queue: pruning stale entry (task done/claimed/missing)"
);
continue;
}
if rescued_task_ids.contains(&entry.task_id) {
info!(
engineer = %entry.engineer,
task_id = entry.task_id,
"dispatch queue: pruning entry — task re-entered orphan-rescue cooldown"
);
continue;
}
if self.is_release_excluded(entry.task_id, &entry.engineer) {
info!(
engineer = %entry.engineer,
task_id = entry.task_id,
"dispatch queue: pruning entry — engineer recently released this task"
);
continue;
}
if benched_engineers.contains(&entry.engineer) {
debug!(
engineer = %entry.engineer,
task_id = entry.task_id,
"dispatch queue: pruning benched engineer entry"
);
continue;
}
if self.states.get(&entry.engineer) == Some(&MemberState::Working)
&& !self.active_tasks.contains_key(&entry.engineer)
{
info!(
engineer = %entry.engineer,
task_id = entry.task_id,
"dispatch queue: recovering engineer stuck in Working with no active task"
);
self.states
.insert(entry.engineer.clone(), MemberState::Idle);
self.update_automation_timers_for_state(&entry.engineer, MemberState::Idle);
}
if self.states.get(&entry.engineer) != Some(&MemberState::Idle) {
retained.push(entry);
continue;
}
if self.should_hold_dispatch_for_stabilization(&entry.engineer) {
retained.push(entry);
continue;
}
let Some(task) = self.task_for_dispatch_entry(&board_dir, &entry)? else {
continue;
};
let retry_dispatchable = self.verification_retry_dispatchable_task(&task, true);
if task.status == "in-progress" && !retry_dispatchable {
info!(
engineer = %entry.engineer,
task_id = task.id,
"dispatch queue: task already in-progress, skipping"
);
continue;
}
if let Some(blocked_ids) = parse_body_dependency_ids(&task.description) {
let all_tasks =
crate::task::load_tasks_from_dir(&board_dir.join("tasks")).unwrap_or_default();
let unmet: Vec<u32> = blocked_ids
.iter()
.filter(|id| {
!all_tasks
.iter()
.any(|t| t.id == **id && dep_status_satisfied(&t.status))
})
.copied()
.collect();
if blocked_ids.is_empty() || !unmet.is_empty() {
warn!(
engineer = %entry.engineer,
task_id = task.id,
?unmet,
"dispatch queue: task has unmet body dependencies, skipping"
);
let _ = crate::team::task_cmd::transition_task_with_attribution(
&board_dir,
task.id,
"blocked",
StatusTransitionAttribution::daemon("daemon.dispatch.queue.dependencies"),
);
continue;
}
}
let active_count =
self.engineer_active_board_item_count(&board_dir, &entry.engineer)?;
let retry_same_owner =
retry_dispatchable && task.claimed_by.as_deref() == Some(entry.engineer.as_str());
let effective_active_count = if retry_same_owner {
active_count.saturating_sub(1)
} else {
active_count
};
if effective_active_count > 0 {
let retained_engineers: HashSet<&str> =
retained.iter().map(|e| e.engineer.as_str()).collect();
let alt = self.idle_engineer_names().into_iter().find(|name| {
name != &entry.engineer
&& !retained_engineers.contains(name.as_str())
&& self
.engineer_active_board_item_count(&board_dir, name)
.unwrap_or(1)
== 0
});
if let Some(alt_engineer) = alt {
debug!(
from = %entry.engineer,
to = %alt_engineer,
task_id = entry.task_id,
"dispatch queue: reassigning to idle engineer"
);
entry.engineer = alt_engineer;
entry.validation_failures = 0;
entry.last_failure = None;
retained.push(entry);
continue;
}
entry.validation_failures += 1;
entry.last_failure = Some(format!(
"Dispatch guard blocked assignment for '{}' with {} active board item(s); no idle alternative",
entry.engineer, effective_active_count
));
if entry.validation_failures >= DISPATCH_QUEUE_FAILURE_LIMIT {
debug!(
engineer = %entry.engineer,
task_id = entry.task_id,
"dispatch queue: all engineers busy, dropping entry (will re-queue)"
);
} else {
retained.push(entry);
}
continue;
}
if !check_wip_limit(
&self.config.team_config.workflow_policy,
RoleType::Engineer,
effective_active_count,
) {
entry.validation_failures += 1;
entry.last_failure = Some(format!(
"WIP gate blocked dispatch for '{}' with {} active board task(s)",
entry.engineer, effective_active_count
));
warn!(
engineer = %entry.engineer,
task_id = entry.task_id,
failures = entry.validation_failures,
"dispatch queue: WIP limit blocked dispatch"
);
if entry.validation_failures >= DISPATCH_QUEUE_FAILURE_LIMIT {
self.escalate_dispatch_queue_entry(
&entry,
entry
.last_failure
.as_deref()
.unwrap_or("wip gate blocked dispatch"),
)?;
} else {
retained.push(entry);
}
continue;
}
let member_uses_worktrees = self.member_uses_worktrees(&entry.engineer);
if member_uses_worktrees {
let worktree_dir = self.worktree_dir(&entry.engineer);
if let Err(error) = engineer_worktree_ready_for_dispatch_from_trunk(
&self.config.project_root,
&worktree_dir,
&entry.engineer,
self.config.team_config.trunk_branch(),
) {
entry.validation_failures += 1;
entry.last_failure = Some(error.to_string());
warn!(
engineer = %entry.engineer,
task_id = entry.task_id,
failures = entry.validation_failures,
error = %error,
"dispatch queue: worktree not ready for dispatch"
);
let base_branch = format!("eng-main/{}", entry.engineer);
let has_work = crate::worktree::commits_ahead(
&worktree_dir,
self.config.team_config.trunk_branch(),
)
.map(|n| n > 0)
.unwrap_or(false)
|| crate::worktree::has_uncommitted_changes(&worktree_dir).unwrap_or(false);
if has_work {
info!(
engineer = %entry.engineer,
"dispatch queue: worktree has work; trying rebase instead of reset"
);
let rebase_result = std::process::Command::new("git")
.args(["rebase", self.config.team_config.trunk_branch()])
.current_dir(&worktree_dir)
.output();
if rebase_result.map(|o| o.status.success()).unwrap_or(false) {
match crate::team::task_loop::engineer_worktree_ready_for_dispatch_from_trunk(
&self.config.project_root,
&worktree_dir,
&entry.engineer,
self.config.team_config.trunk_branch(),
) {
Ok(()) => {
info!(
engineer = %entry.engineer,
"dispatch queue: rebase succeeded; retrying dispatch"
);
entry.validation_failures = 0;
entry.last_failure = None;
retained.push(entry);
continue;
}
Err(error) => {
warn!(
engineer = %entry.engineer,
error = %error,
"dispatch queue: rebase succeeded but worktree is still not ready; falling through to reset"
);
}
}
}
let _ = std::process::Command::new("git")
.args(["rebase", "--abort"])
.current_dir(&worktree_dir)
.output();
warn!(
engineer = %entry.engineer,
"dispatch queue: rebase failed; falling through to reset (work may be lost)"
);
}
info!(
engineer = %entry.engineer,
base_branch = %base_branch,
"dispatch queue: auto-resetting worktree to base branch"
);
match crate::worktree::reset_worktree_to_base_if_clean_from_trunk(
&worktree_dir,
&base_branch,
"dispatch/reset recovery",
self.config.team_config.trunk_branch(),
) {
Err(reset_err) => {
warn!(
engineer = %entry.engineer,
error = %reset_err,
"dispatch queue: worktree auto-reset failed; escalating"
);
entry.validation_failures += 1;
entry.last_failure = Some(reset_err.to_string());
self.report_preserve_failure(
&entry.engineer,
None,
"dispatch/reset recovery",
&reset_err.to_string(),
);
if entry.validation_failures >= DISPATCH_QUEUE_FAILURE_LIMIT {
self.escalate_dispatch_queue_entry(
&entry,
entry
.last_failure
.as_deref()
.unwrap_or("worktree readiness validation failed"),
)?;
} else {
retained.push(entry);
}
}
Ok(reason) if reason.reset_performed() => {
info!(
engineer = %entry.engineer,
reset_reason = reason.as_str(),
"dispatch queue: worktree auto-reset succeeded; retrying dispatch"
);
entry.validation_failures = 0;
entry.last_failure = None;
retained.push(entry);
}
Ok(reason) => {
warn!(
engineer = %entry.engineer,
reset_reason = reason.as_str(),
"dispatch queue: worktree auto-reset skipped"
);
entry.validation_failures += 1;
entry.last_failure = Some(
crate::team::task_loop::dirty_worktree_preservation_blocked_reason(
&worktree_dir,
"dispatch/reset recovery",
),
);
self.report_preserve_failure(
&entry.engineer,
None,
"dispatch/reset recovery",
reason.as_str(),
);
retained.push(entry);
}
}
continue;
}
}
if task.status == "backlog" {
let _ = transition_task_with_attribution(
&board_dir,
task.id,
"todo",
StatusTransitionAttribution::daemon("daemon.dispatch.queue"),
);
}
if let Err(e) = transition_task_with_attribution(
&board_dir,
task.id,
"in-progress",
StatusTransitionAttribution::daemon("daemon.dispatch.queue"),
) {
entry.validation_failures += 1;
entry.last_failure = Some(format!("board transition failed: {e}"));
warn!(
engineer = %entry.engineer,
task_id = task.id,
error = %e,
"dispatch queue: cannot transition task to in-progress, deferring"
);
if entry.validation_failures >= DISPATCH_QUEUE_FAILURE_LIMIT {
self.escalate_dispatch_queue_entry(
&entry,
entry
.last_failure
.as_deref()
.unwrap_or("board transition failed"),
)?;
} else {
retained.push(entry);
}
continue;
}
assign_task_owners(&board_dir, task.id, Some(&entry.engineer), None)?;
let assignment_message =
format!("Task #{}: {}\n\n{}", task.id, task.title, task.description);
let assignment_message =
if let Some(context) = verification_retry_assignment_context(&task) {
format!("{assignment_message}\n\n{context}")
} else {
assignment_message
};
match self.assign_task_with_task_id(&entry.engineer, &assignment_message, Some(task.id))
{
Ok(_) => {
self.active_tasks.insert(entry.engineer.clone(), task.id);
self.retry_counts.remove(&entry.engineer);
self.recent_dispatches
.insert((task.id, entry.engineer.clone()), Instant::now());
self.record_orchestrator_action(format!(
"dispatch queue: selected runnable task #{} ({}) and dispatched it to {}",
task.id, task.title, entry.engineer
));
info!(
engineer = %entry.engineer,
task_id = task.id,
task_title = %task.title,
"queued task dispatched"
);
}
Err(error) => {
entry.validation_failures += 1;
entry.last_failure = Some(error.to_string());
warn!(
engineer = %entry.engineer,
task_id = entry.task_id,
failures = entry.validation_failures,
error = %error,
"dispatch queue: assignment launch failed"
);
if entry.validation_failures >= DISPATCH_QUEUE_FAILURE_LIMIT {
self.escalate_dispatch_queue_entry(
&entry,
entry
.last_failure
.as_deref()
.unwrap_or("assignment launch failed"),
)?;
} else {
retained.push(entry);
}
}
}
}
self.dispatch_queue = retained;
Ok(())
}
fn rank_dispatch_engineers(
&self,
task: &crate::task::Task,
queued_engineers: &HashSet<String>,
benched_engineers: &std::collections::BTreeSet<String>,
manual_cooldown: Duration,
profiles: &HashMap<String, EngineerProfile>,
) -> Vec<String> {
let mut eligible: Vec<String> = self
.idle_engineer_names()
.into_iter()
.filter(|engineer_name| !queued_engineers.contains(engineer_name))
.filter(|engineer_name| !benched_engineers.contains(engineer_name))
.filter(|engineer_name| {
task.assignee
.as_deref()
.is_none_or(|preferred| preferred == engineer_name)
})
.filter(|engineer_name| {
if self.member_backend_parked(engineer_name) {
debug!(
engineer = %engineer_name,
"skipping dispatch — backend quota parked"
);
return false;
}
let Some(assigned_at) = self.manual_assign_cooldowns.get(engineer_name) else {
return true;
};
if assigned_at.elapsed() < manual_cooldown {
debug!(
engineer = %engineer_name,
"skipping dispatch — within manual assignment cooldown"
);
false
} else {
true
}
})
.collect();
eligible.sort();
let body_owner_role = parse_body_owner_role(&task.description);
if let Some(ref owner_role) = body_owner_role {
let engineers_with_role: HashSet<String> = self
.config
.members
.iter()
.filter(|m| m.role_type == RoleType::Engineer && &m.role_name == owner_role)
.map(|m| m.name.clone())
.collect();
if !engineers_with_role.is_empty() {
eligible.retain(|name| engineers_with_role.contains(name));
}
}
if self.config.team_config.workflow_policy.allocation.strategy
== AllocationStrategy::RoundRobin
{
return eligible;
}
let task_for_ranking = match body_owner_role {
Some(owner_role) if !task.tags.iter().any(|tag| tag == &owner_role) => {
let mut synth = task.clone();
synth.tags.push(owner_role);
std::borrow::Cow::Owned(synth)
}
_ => std::borrow::Cow::Borrowed(task),
};
rank_engineers_for_task(
&eligible,
profiles,
&task_for_ranking,
&self.config.team_config.workflow_policy.allocation,
)
}
}
#[cfg(test)]
mod tests {
use std::collections::{HashMap, HashSet};
use std::path::Path;
use super::{
OverlapConflict, find_overlapping_tasks, parse_body_owner_role, predicted_files,
role_name_seed_tags, split_acyclic_blocking_ids,
};
use crate::team::config::RoleType;
use crate::team::hierarchy::MemberInstance;
use crate::team::standup::MemberState;
use crate::team::task_loop::{
current_worktree_branch, engineer_base_branch_name, setup_engineer_worktree,
};
use crate::team::test_support::{
TestDaemonBuilder, architect_member, engineer_member, git_ok, git_stdout, init_git_repo,
manager_member, write_open_task_file, write_owned_task_file,
};
fn write_task_with_priority(project_root: &Path, id: u32, title: &str, priority: &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: todo\npriority: {priority}\nclass: standard\n---\n\nTask.\n"
),
)
.unwrap();
}
fn write_task_with_deps(project_root: &Path, id: u32, title: &str, depends_on: &[u32]) {
let tasks_dir = project_root
.join(".batty")
.join("team_config")
.join("board")
.join("tasks");
std::fs::create_dir_all(&tasks_dir).unwrap();
let mut content = format!("---\nid: {id}\ntitle: {title}\nstatus: todo\npriority: high\n");
if !depends_on.is_empty() {
content.push_str("depends_on:\n");
for dep in depends_on {
content.push_str(&format!(" - {dep}\n"));
}
}
content.push_str("class: standard\n---\n\nTask.\n");
std::fs::write(tasks_dir.join(format!("{id:03}-{title}.md")), content).unwrap();
}
fn write_task_with_body(
project_root: &Path,
id: u32,
title: &str,
status: &str,
claimed_by: Option<&str>,
body: &str,
) {
let tasks_dir = project_root
.join(".batty")
.join("team_config")
.join("board")
.join("tasks");
std::fs::create_dir_all(&tasks_dir).unwrap();
let mut content =
format!("---\nid: {id}\ntitle: {title}\nstatus: {status}\npriority: high\n");
if let Some(claimed_by) = claimed_by {
content.push_str(&format!("claimed_by: {claimed_by}\n"));
}
content.push_str("class: standard\n---\n\n");
content.push_str(body);
content.push('\n');
std::fs::write(tasks_dir.join(format!("{id:03}-{title}.md")), content).unwrap();
}
fn write_task_with_files(
project_root: &Path,
id: u32,
title: &str,
status: &str,
claimed_by: Option<&str>,
files: &[&str],
body: &str,
) {
let tasks_dir = project_root
.join(".batty")
.join("team_config")
.join("board")
.join("tasks");
std::fs::create_dir_all(&tasks_dir).unwrap();
let mut content =
format!("---\nid: {id}\ntitle: {title}\nstatus: {status}\npriority: high\n");
if let Some(claimed_by) = claimed_by {
content.push_str(&format!("claimed_by: {claimed_by}\n"));
}
if !files.is_empty() {
content.push_str("files:\n");
for file in files {
content.push_str(&format!(" - {file}\n"));
}
}
content.push_str("class: standard\n---\n\n");
content.push_str(body);
content.push('\n');
std::fs::write(tasks_dir.join(format!("{id:03}-{title}.md")), content).unwrap();
}
fn write_task_with_assignee(project_root: &Path, id: u32, title: &str, assignee: &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: todo\npriority: high\nassignee: {assignee}\nclass: standard\n---\n\nTask.\n"
),
)
.unwrap();
}
fn write_task_with_tags(project_root: &Path, id: u32, title: &str, tags: &[&str]) {
let tasks_dir = project_root
.join(".batty")
.join("team_config")
.join("board")
.join("tasks");
std::fs::create_dir_all(&tasks_dir).unwrap();
let tags_block = tags
.iter()
.map(|tag| format!(" - {tag}"))
.collect::<Vec<_>>()
.join("\n");
std::fs::write(
tasks_dir.join(format!("{id:03}-{title}.md")),
format!(
"---\nid: {id}\ntitle: {title}\nstatus: todo\npriority: high\ntags:\n{tags_block}\nclass: standard\n---\n\nTask.\n"
),
)
.unwrap();
}
fn write_task_with_tags_and_body(
project_root: &Path,
id: u32,
title: &str,
tags: &[&str],
body: &str,
) {
let tasks_dir = project_root
.join(".batty")
.join("team_config")
.join("board")
.join("tasks");
std::fs::create_dir_all(&tasks_dir).unwrap();
let tags_block = tags
.iter()
.map(|tag| format!(" - {tag}"))
.collect::<Vec<_>>()
.join("\n");
std::fs::write(
tasks_dir.join(format!("{id:03}-{title}.md")),
format!(
"---\nid: {id}\ntitle: {title}\nstatus: todo\npriority: high\ntags:\n{tags_block}\nclass: standard\n---\n\n{body}\n"
),
)
.unwrap();
}
fn write_bench_test_team_config(project_root: &Path, engineer_instances: u32) {
let team_dir = project_root.join(".batty").join("team_config");
std::fs::create_dir_all(&team_dir).unwrap();
std::fs::write(
team_dir.join("team.yaml"),
format!(
"name: test\nagent: codex\nroles:\n - name: eng\n role_type: engineer\n instances: {engineer_instances}\n"
),
)
.unwrap();
}
#[test]
fn idle_engineers_returns_only_idle() {
let tmp = tempfile::tempdir().unwrap();
let mut daemon = TestDaemonBuilder::new(tmp.path())
.members(vec![
manager_member("mgr", None),
engineer_member("eng-1", Some("mgr"), false),
engineer_member("eng-2", Some("mgr"), false),
engineer_member("eng-3", Some("mgr"), false),
])
.states(HashMap::from([
("eng-1".to_string(), MemberState::Idle),
("eng-2".to_string(), MemberState::Working),
("eng-3".to_string(), MemberState::Idle),
]))
.build();
daemon.active_tasks.insert("eng-2".to_string(), 42);
let idle = daemon.idle_engineer_names();
assert_eq!(idle, vec!["eng-1", "eng-3"]);
}
#[test]
fn idle_engineers_empty_when_all_working_with_tasks() {
let tmp = tempfile::tempdir().unwrap();
let mut daemon = TestDaemonBuilder::new(tmp.path())
.members(vec![
manager_member("mgr", None),
engineer_member("eng-1", Some("mgr"), false),
])
.states(HashMap::from([("eng-1".to_string(), MemberState::Working)]))
.build();
daemon.active_tasks.insert("eng-1".to_string(), 10);
assert!(daemon.idle_engineer_names().is_empty());
}
#[test]
fn idle_engineers_includes_working_without_active_task() {
let tmp = tempfile::tempdir().unwrap();
let daemon = TestDaemonBuilder::new(tmp.path())
.members(vec![
manager_member("mgr", None),
engineer_member("eng-1", Some("mgr"), false),
engineer_member("eng-2", Some("mgr"), false),
])
.states(HashMap::from([
("eng-1".to_string(), MemberState::Working),
("eng-2".to_string(), MemberState::Idle),
]))
.build();
let idle = daemon.idle_engineer_names();
assert_eq!(idle, vec!["eng-1", "eng-2"]);
}
#[test]
fn idle_engineers_working_no_task_mixed_with_working_with_task() {
let tmp = tempfile::tempdir().unwrap();
let mut daemon = TestDaemonBuilder::new(tmp.path())
.members(vec![
manager_member("mgr", None),
engineer_member("eng-1", Some("mgr"), false),
engineer_member("eng-2", Some("mgr"), false),
engineer_member("eng-3", Some("mgr"), false),
])
.states(HashMap::from([
("eng-1".to_string(), MemberState::Working),
("eng-2".to_string(), MemberState::Working),
("eng-3".to_string(), MemberState::Idle),
]))
.build();
daemon.active_tasks.insert("eng-1".to_string(), 50);
let idle = daemon.idle_engineer_names();
assert_eq!(idle, vec!["eng-2", "eng-3"]);
}
#[test]
fn idle_engineers_excludes_managers() {
let tmp = tempfile::tempdir().unwrap();
let daemon = TestDaemonBuilder::new(tmp.path())
.members(vec![
manager_member("mgr", None),
engineer_member("eng-1", Some("mgr"), false),
])
.states(HashMap::from([
("mgr".to_string(), MemberState::Idle),
("eng-1".to_string(), MemberState::Idle),
]))
.build();
let idle = daemon.idle_engineer_names();
assert_eq!(idle, vec!["eng-1"]);
}
#[test]
fn next_task_picks_highest_priority() {
let tmp = tempfile::tempdir().unwrap();
write_task_with_priority(tmp.path(), 10, "low-pri", "low");
write_task_with_priority(tmp.path(), 11, "critical-pri", "critical");
write_task_with_priority(tmp.path(), 12, "medium-pri", "medium");
let daemon = TestDaemonBuilder::new(tmp.path()).build();
let board_dir = tmp.path().join(".batty").join("team_config").join("board");
let task = daemon
.test_next_dispatch_task(&board_dir, &HashSet::new())
.unwrap()
.unwrap();
assert_eq!(task.id, 11, "should pick the critical-priority task");
}
#[test]
fn next_task_breaks_ties_by_id() {
let tmp = tempfile::tempdir().unwrap();
write_task_with_priority(tmp.path(), 20, "second", "high");
write_task_with_priority(tmp.path(), 10, "first", "high");
let daemon = TestDaemonBuilder::new(tmp.path()).build();
let board_dir = tmp.path().join(".batty").join("team_config").join("board");
let task = daemon
.test_next_dispatch_task(&board_dir, &HashSet::new())
.unwrap()
.unwrap();
assert_eq!(task.id, 10, "should pick lower id when priority is equal");
}
#[test]
fn next_task_skips_claimed_tasks() {
let tmp = tempfile::tempdir().unwrap();
write_owned_task_file(tmp.path(), 10, "claimed-task", "todo", "eng-2");
write_open_task_file(tmp.path(), 11, "open-task", "todo");
let daemon = TestDaemonBuilder::new(tmp.path()).build();
let board_dir = tmp.path().join(".batty").join("team_config").join("board");
let task = daemon
.test_next_dispatch_task(&board_dir, &HashSet::new())
.unwrap()
.unwrap();
assert_eq!(task.id, 11, "should skip claimed task");
}
#[test]
fn next_task_skips_done_tasks() {
let tmp = tempfile::tempdir().unwrap();
write_open_task_file(tmp.path(), 10, "done-task", "done");
write_open_task_file(tmp.path(), 11, "open-task", "todo");
let daemon = TestDaemonBuilder::new(tmp.path()).build();
let board_dir = tmp.path().join(".batty").join("team_config").join("board");
let task = daemon
.test_next_dispatch_task(&board_dir, &HashSet::new())
.unwrap()
.unwrap();
assert_eq!(task.id, 11);
}
#[test]
fn next_task_skips_already_queued() {
let tmp = tempfile::tempdir().unwrap();
write_open_task_file(tmp.path(), 10, "queued", "todo");
write_open_task_file(tmp.path(), 11, "available", "todo");
let daemon = TestDaemonBuilder::new(tmp.path()).build();
let board_dir = tmp.path().join(".batty").join("team_config").join("board");
let queued: HashSet<u32> = [10].into();
let task = daemon
.test_next_dispatch_task(&board_dir, &queued)
.unwrap()
.unwrap();
assert_eq!(task.id, 11, "should skip task already in queue set");
}
#[test]
fn next_task_skips_blocked_dependencies() {
let tmp = tempfile::tempdir().unwrap();
write_open_task_file(tmp.path(), 9, "dep-task", "in-progress");
write_task_with_deps(tmp.path(), 10, "blocked-task", &[9]);
write_open_task_file(tmp.path(), 11, "free-task", "todo");
let daemon = TestDaemonBuilder::new(tmp.path()).build();
let board_dir = tmp.path().join(".batty").join("team_config").join("board");
let task = daemon
.test_next_dispatch_task(&board_dir, &HashSet::new())
.unwrap()
.unwrap();
assert_eq!(task.id, 11, "should skip task with unmet dependency");
}
#[test]
fn next_task_skips_unmet_body_dependency() {
let tmp = tempfile::tempdir().unwrap();
write_open_task_file(tmp.path(), 9, "dep-task", "in-progress");
write_task_with_body(
tmp.path(),
10,
"body-blocked",
"todo",
None,
"Blocked on: #9",
);
write_open_task_file(tmp.path(), 11, "free-task", "todo");
let daemon = TestDaemonBuilder::new(tmp.path()).build();
let board_dir = tmp.path().join(".batty").join("team_config").join("board");
let task = daemon
.test_next_dispatch_task(&board_dir, &HashSet::new())
.unwrap()
.unwrap();
assert_eq!(task.id, 11, "should skip task with unmet body dependency");
}
#[test]
fn next_task_allows_met_body_dependency() {
let tmp = tempfile::tempdir().unwrap();
write_open_task_file(tmp.path(), 9, "dep-done", "done");
write_task_with_body(
tmp.path(),
10,
"body-unblocked",
"todo",
None,
"Blocked on: #9",
);
let daemon = TestDaemonBuilder::new(tmp.path()).build();
let board_dir = tmp.path().join(".batty").join("team_config").join("board");
let task = daemon
.test_next_dispatch_task(&board_dir, &HashSet::new())
.unwrap()
.unwrap();
assert_eq!(task.id, 10, "should allow satisfied body dependency");
}
#[test]
fn next_task_skips_body_blocker_without_task_id() {
let tmp = tempfile::tempdir().unwrap();
write_task_with_body(
tmp.path(),
10,
"body-blocked",
"todo",
None,
"Blocked on: provider-console token",
);
write_open_task_file(tmp.path(), 11, "free-task", "todo");
let daemon = TestDaemonBuilder::new(tmp.path()).build();
let board_dir = tmp.path().join(".batty").join("team_config").join("board");
let task = daemon
.test_next_dispatch_task(&board_dir, &HashSet::new())
.unwrap()
.unwrap();
assert_eq!(task.id, 11, "should skip non-task-id body blocker");
}
#[test]
fn next_task_skips_blocked_or_reworked_parent_dependencies() {
let tmp = tempfile::tempdir().unwrap();
write_open_task_file(tmp.path(), 8, "blocked-parent", "blocked");
write_open_task_file(tmp.path(), 9, "rework-parent", "rework");
write_task_with_deps(tmp.path(), 10, "blocked-child", &[8]);
write_task_with_deps(tmp.path(), 11, "rework-child", &[9]);
write_open_task_file(tmp.path(), 12, "free-task", "todo");
let daemon = TestDaemonBuilder::new(tmp.path()).build();
let board_dir = tmp.path().join(".batty").join("team_config").join("board");
let task = daemon
.test_next_dispatch_task(&board_dir, &HashSet::new())
.unwrap()
.unwrap();
assert_eq!(
task.id, 12,
"should only dispatch work whose parents are done or archived"
);
}
#[test]
fn next_task_allows_met_dependencies() {
let tmp = tempfile::tempdir().unwrap();
write_open_task_file(tmp.path(), 9, "dep-done", "done");
write_task_with_deps(tmp.path(), 10, "unblocked", &[9]);
let daemon = TestDaemonBuilder::new(tmp.path()).build();
let board_dir = tmp.path().join(".batty").join("team_config").join("board");
let task = daemon
.test_next_dispatch_task(&board_dir, &HashSet::new())
.unwrap()
.unwrap();
assert_eq!(task.id, 10, "should pick task with satisfied dependency");
}
#[test]
fn next_task_returns_none_when_empty() {
let tmp = tempfile::tempdir().unwrap();
let tasks_dir = tmp
.path()
.join(".batty")
.join("team_config")
.join("board")
.join("tasks");
std::fs::create_dir_all(&tasks_dir).unwrap();
let daemon = TestDaemonBuilder::new(tmp.path()).build();
let board_dir = tmp.path().join(".batty").join("team_config").join("board");
assert!(
daemon
.test_next_dispatch_task(&board_dir, &HashSet::new())
.unwrap()
.is_none()
);
}
#[test]
fn next_task_accepts_backlog_status() {
let tmp = tempfile::tempdir().unwrap();
write_open_task_file(tmp.path(), 10, "backlog-task", "backlog");
let daemon = TestDaemonBuilder::new(tmp.path()).build();
let board_dir = tmp.path().join(".batty").join("team_config").join("board");
let task = daemon
.test_next_dispatch_task(&board_dir, &HashSet::new())
.unwrap()
.unwrap();
assert_eq!(task.id, 10, "backlog status should be dispatchable");
}
#[test]
fn process_queue_prunes_entry_for_done_task_even_when_engineer_not_idle() {
use super::DispatchQueueEntry;
let tmp = tempfile::tempdir().unwrap();
write_owned_task_file(tmp.path(), 10, "finished", "done", "other-eng");
let mut daemon = TestDaemonBuilder::new(tmp.path())
.members(vec![
manager_member("mgr", None),
engineer_member("eng-1", Some("mgr"), false),
])
.states(HashMap::from([("eng-1".to_string(), MemberState::Working)]))
.build();
daemon.dispatch_queue.push(DispatchQueueEntry {
engineer: "eng-1".to_string(),
task_id: 10,
task_title: "finished".to_string(),
queued_at: 0,
validation_failures: 0,
last_failure: None,
});
daemon.process_dispatch_queue().unwrap();
assert!(
daemon.dispatch_queue.is_empty(),
"entry for done task should be pruned even when engineer is Working"
);
}
#[test]
fn process_queue_retains_valid_entry_for_non_idle_engineer() {
use super::DispatchQueueEntry;
let tmp = tempfile::tempdir().unwrap();
write_open_task_file(tmp.path(), 10, "pending-work", "todo");
let mut daemon = TestDaemonBuilder::new(tmp.path())
.members(vec![
manager_member("mgr", None),
engineer_member("eng-1", Some("mgr"), false),
])
.states(HashMap::from([("eng-1".to_string(), MemberState::Working)]))
.build();
daemon.dispatch_queue.push(DispatchQueueEntry {
engineer: "eng-1".to_string(),
task_id: 10,
task_title: "pending-work".to_string(),
queued_at: 0,
validation_failures: 0,
last_failure: None,
});
daemon.process_dispatch_queue().unwrap();
assert_eq!(
daemon.dispatch_queue.len(),
1,
"entry for valid todo task should be retained while engineer is Working"
);
}
#[test]
fn process_queue_blocks_dirty_worktree_instead_of_auto_preserving() {
use super::DispatchQueueEntry;
let tmp = tempfile::tempdir().unwrap();
let repo = init_git_repo(&tmp, "dispatch-preserve-reset");
write_open_task_file(&repo, 42, "dispatch-reset", "todo");
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"]);
std::fs::write(worktree_dir.join("tracked.txt"), "tracked dispatch work\n").unwrap();
git_ok(&worktree_dir, &["add", "tracked.txt"]);
std::fs::write(
worktree_dir.join("untracked.txt"),
"untracked dispatch work\n",
)
.unwrap();
let mut daemon = TestDaemonBuilder::new(repo.as_path())
.members(vec![
manager_member("mgr", None),
engineer_member("eng-1", Some("mgr"), true),
])
.board(crate::team::config::BoardConfig {
dispatch_stabilization_delay_secs: 0,
..crate::team::config::BoardConfig::default()
})
.states(HashMap::from([("eng-1".to_string(), MemberState::Idle)]))
.build();
daemon.idle_started_at.insert(
"eng-1".to_string(),
std::time::Instant::now() - std::time::Duration::from_secs(1),
);
daemon.dispatch_queue.push(DispatchQueueEntry {
engineer: "eng-1".to_string(),
task_id: 42,
task_title: "dispatch-reset".to_string(),
queued_at: 0,
validation_failures: 0,
last_failure: None,
});
daemon.process_dispatch_queue().unwrap();
assert_eq!(daemon.dispatch_queue.len(), 1);
assert_eq!(daemon.dispatch_queue[0].validation_failures, 2);
assert!(
daemon.dispatch_queue[0]
.last_failure
.as_deref()
.unwrap_or("")
.contains("could not safely auto-save dirty worktree")
);
assert_eq!(current_worktree_branch(&worktree_dir).unwrap(), "eng-1/41");
let status = git_stdout(&worktree_dir, &["status", "--short"]);
assert!(
status.contains("A tracked.txt"),
"tracked work should remain staged instead of being auto-preserved: {status}"
);
assert!(
status.contains("?? untracked.txt"),
"untracked work should remain untouched instead of being auto-preserved: {status}"
);
assert!(
git_stdout(&repo, &["branch", "--list", "eng-1/41"]).contains("eng-1/41"),
"dirty task branch should remain in place for manual recovery"
);
}
#[test]
fn process_queue_blocks_dirty_worktree_when_preserve_fails() {
use super::DispatchQueueEntry;
let tmp = tempfile::tempdir().unwrap();
let repo = init_git_repo(&tmp, "dispatch-preserve-blocked");
write_open_task_file(&repo, 42, "dispatch-reset", "todo");
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"]);
std::fs::write(worktree_dir.join("tracked.txt"), "tracked dispatch work\n").unwrap();
git_ok(&worktree_dir, &["add", "tracked.txt"]);
std::fs::write(worktree_dir.join("unstaged.txt"), "leave unstaged\n").unwrap();
let git_dir =
std::path::PathBuf::from(git_stdout(&worktree_dir, &["rev-parse", "--git-dir"]));
let git_dir = if git_dir.is_absolute() {
git_dir
} else {
worktree_dir.join(git_dir)
};
std::fs::write(git_dir.join("index.lock"), "locked\n").unwrap();
let mut daemon = TestDaemonBuilder::new(repo.as_path())
.members(vec![
manager_member("mgr", None),
engineer_member("eng-1", Some("mgr"), true),
])
.board(crate::team::config::BoardConfig {
dispatch_stabilization_delay_secs: 0,
..crate::team::config::BoardConfig::default()
})
.states(HashMap::from([("eng-1".to_string(), MemberState::Idle)]))
.build();
daemon.idle_started_at.insert(
"eng-1".to_string(),
std::time::Instant::now() - std::time::Duration::from_secs(1),
);
daemon.dispatch_queue.push(DispatchQueueEntry {
engineer: "eng-1".to_string(),
task_id: 42,
task_title: "dispatch-reset".to_string(),
queued_at: 0,
validation_failures: 0,
last_failure: None,
});
daemon.process_dispatch_queue().unwrap();
assert_eq!(daemon.dispatch_queue.len(), 1);
assert_eq!(daemon.dispatch_queue[0].validation_failures, 2);
assert!(
daemon.dispatch_queue[0]
.last_failure
.as_deref()
.unwrap_or("")
.contains("could not safely auto-save dirty worktree")
);
assert_eq!(current_worktree_branch(&worktree_dir).unwrap(), "eng-1/41");
let status = git_stdout(&worktree_dir, &["status", "--short"]);
assert!(
status.contains("A tracked.txt"),
"pre-existing staged work should remain staged: {status}"
);
assert!(
status.contains("?? unstaged.txt"),
"idle dispatch recovery must not stage new files: {status}"
);
}
fn write_blocked_task(project_root: &Path, id: u32, title: &str, block_reason: &str) {
let tasks_dir = project_root
.join(".batty")
.join("team_config")
.join("board")
.join("tasks");
std::fs::create_dir_all(&tasks_dir).unwrap();
let content = format!(
"---\nid: {id}\ntitle: {title}\nstatus: todo\npriority: high\nblocked: true\nblock_reason: \"{block_reason}\"\nclass: standard\n---\n\nBody.\n"
);
std::fs::write(tasks_dir.join(format!("{id:03}-{title}.md")), content).unwrap();
}
#[test]
fn enqueue_dispatch_candidates_skips_kanban_md_blocked_tasks() {
let tmp = tempfile::tempdir().unwrap();
write_blocked_task(
tmp.path(),
30,
"kanban-md-blocked",
"Deferred per architect",
);
write_task_with_body(
tmp.path(),
31,
"runnable-candidate",
"todo",
None,
"Touch src/team/telemetry_db.rs only.",
);
let mut daemon = TestDaemonBuilder::new(tmp.path())
.members(vec![
manager_member("mgr", None),
engineer_member("eng-1", Some("mgr"), false),
])
.states(HashMap::from([("eng-1".to_string(), MemberState::Idle)]))
.build();
daemon.enqueue_dispatch_candidates().unwrap();
assert_eq!(daemon.dispatch_queue.len(), 1);
assert_eq!(
daemon.dispatch_queue[0].task_id, 31,
"blocked task #30 must be filtered out; only the runnable #31 should be queued"
);
}
#[test]
fn enqueue_dispatch_candidates_skips_tasks_assigned_to_non_engineer() {
let tmp = tempfile::tempdir().unwrap();
write_task_with_assignee(tmp.path(), 40, "pm-intake", "mgr");
write_task_with_body(
tmp.path(),
41,
"engineer-candidate",
"todo",
None,
"Touch src/team/telemetry_db.rs only.",
);
let mut daemon = TestDaemonBuilder::new(tmp.path())
.members(vec![
manager_member("mgr", None),
engineer_member("eng-1", Some("mgr"), false),
])
.states(HashMap::from([("eng-1".to_string(), MemberState::Idle)]))
.build();
daemon.enqueue_dispatch_candidates().unwrap();
assert_eq!(daemon.dispatch_queue.len(), 1);
assert_eq!(
daemon.dispatch_queue[0].task_id, 41,
"non-engineer-assigned task #40 must be filtered out; only #41 should dispatch"
);
}
#[test]
fn enqueue_dispatch_candidates_routes_engineer_assigned_task_to_named_engineer() {
let tmp = tempfile::tempdir().unwrap();
write_task_with_assignee(tmp.path(), 50, "for-eng-2", "eng-2");
let mut daemon = TestDaemonBuilder::new(tmp.path())
.members(vec![
manager_member("mgr", None),
engineer_member("eng-1", Some("mgr"), false),
engineer_member("eng-2", Some("mgr"), false),
engineer_member("eng-3", Some("mgr"), false),
])
.states(HashMap::from([
("eng-1".to_string(), MemberState::Idle),
("eng-2".to_string(), MemberState::Idle),
("eng-3".to_string(), MemberState::Idle),
]))
.build();
daemon.enqueue_dispatch_candidates().unwrap();
assert_eq!(daemon.dispatch_queue.len(), 1);
assert_eq!(daemon.dispatch_queue[0].task_id, 50);
assert_eq!(
daemon.dispatch_queue[0].engineer, "eng-2",
"task with `assignee: eng-2` must dispatch to eng-2, not a peer"
);
}
#[test]
fn enqueue_dispatch_candidates_waits_when_assigned_engineer_busy() {
let tmp = tempfile::tempdir().unwrap();
write_task_with_assignee(tmp.path(), 60, "for-eng-2", "eng-2");
let mut daemon = TestDaemonBuilder::new(tmp.path())
.members(vec![
manager_member("mgr", None),
engineer_member("eng-1", Some("mgr"), false),
engineer_member("eng-2", Some("mgr"), false),
])
.states(HashMap::from([("eng-1".to_string(), MemberState::Idle)]))
.build();
daemon.enqueue_dispatch_candidates().unwrap();
assert!(
daemon.dispatch_queue.is_empty(),
"task must remain undispatched while named engineer is unavailable"
);
}
#[test]
fn enqueue_dispatch_candidates_skips_tasks_whose_body_owner_names_non_engineer() {
let tmp = tempfile::tempdir().unwrap();
write_task_with_body(
tmp.path(),
70,
"strategy-star-velocity-gate",
"todo",
None,
"**Owner:** maya-lead (this task), with input from jordan-pm \
and kai-devrel. Mid-window gate decision for launch.\n",
);
write_task_with_body(
tmp.path(),
71,
"engineer-candidate",
"todo",
None,
"Touch src/team/telemetry_db.rs only.\n",
);
let mut daemon = TestDaemonBuilder::new(tmp.path())
.members(vec![
architect_member("maya-lead"),
manager_member("jordan-pm", Some("maya-lead")),
engineer_member("eng-1", Some("jordan-pm"), false),
])
.states(HashMap::from([("eng-1".to_string(), MemberState::Idle)]))
.build();
daemon.enqueue_dispatch_candidates().unwrap();
assert_eq!(daemon.dispatch_queue.len(), 1);
assert_eq!(
daemon.dispatch_queue[0].task_id, 71,
"task with `**Owner:** maya-lead` body must be filtered out \
because maya-lead is a non-engineer member; only #71 should dispatch"
);
}
#[test]
fn enqueue_dispatch_candidates_allows_body_owner_when_it_names_an_engineer() {
let tmp = tempfile::tempdir().unwrap();
write_task_with_body(
tmp.path(),
80,
"engineer-owned",
"todo",
None,
"- Owner: priya-writer drafts; review by peer.\n",
);
let priya = MemberInstance {
name: "priya-writer-1-1".to_string(),
role_name: "priya-writer".to_string(),
role_type: RoleType::Engineer,
agent: Some("claude".to_string()),
reports_to: Some("jordan-pm".to_string()),
use_worktrees: false,
..MemberInstance::default()
};
let mut daemon = TestDaemonBuilder::new(tmp.path())
.members(vec![
architect_member("maya-lead"),
manager_member("jordan-pm", Some("maya-lead")),
priya,
])
.states(HashMap::from([(
"priya-writer-1-1".to_string(),
MemberState::Idle,
)]))
.build();
daemon.enqueue_dispatch_candidates().unwrap();
assert_eq!(daemon.dispatch_queue.len(), 1);
assert_eq!(daemon.dispatch_queue[0].task_id, 80);
assert_eq!(daemon.dispatch_queue[0].engineer, "priya-writer-1-1");
}
#[test]
fn enqueue_dispatch_candidates_waits_when_body_owner_engineer_busy() {
let tmp = tempfile::tempdir().unwrap();
write_task_with_body(
tmp.path(),
90,
"twir-submission",
"todo",
None,
"- Route: dispatch to priya-writer.\n",
);
let priya = MemberInstance {
name: "priya-writer-1-1".to_string(),
role_name: "priya-writer".to_string(),
role_type: RoleType::Engineer,
agent: Some("claude".to_string()),
reports_to: Some("jordan-pm".to_string()),
use_worktrees: false,
..MemberInstance::default()
};
let kai = MemberInstance {
name: "kai-devrel-1-1".to_string(),
role_name: "kai-devrel".to_string(),
role_type: RoleType::Engineer,
agent: Some("claude".to_string()),
reports_to: Some("jordan-pm".to_string()),
use_worktrees: false,
..MemberInstance::default()
};
let mut daemon = TestDaemonBuilder::new(tmp.path())
.members(vec![
architect_member("maya-lead"),
manager_member("jordan-pm", Some("maya-lead")),
priya,
kai,
])
.states(HashMap::from([(
"kai-devrel-1-1".to_string(),
MemberState::Idle,
)]))
.build();
daemon.enqueue_dispatch_candidates().unwrap();
assert!(
daemon.dispatch_queue.is_empty(),
"task must remain undispatched while named body-owner engineer is busy; \
got {:?}",
daemon.dispatch_queue
);
}
#[test]
fn enqueue_dispatch_candidates_dispatches_to_body_owner_when_idle_even_with_other_idle_peers() {
let tmp = tempfile::tempdir().unwrap();
write_task_with_body(
tmp.path(),
91,
"twir-submission",
"todo",
None,
"- Route: dispatch to priya-writer.\n",
);
let priya = MemberInstance {
name: "priya-writer-1-1".to_string(),
role_name: "priya-writer".to_string(),
role_type: RoleType::Engineer,
agent: Some("claude".to_string()),
reports_to: Some("jordan-pm".to_string()),
use_worktrees: false,
..MemberInstance::default()
};
let kai = MemberInstance {
name: "kai-devrel-1-1".to_string(),
role_name: "kai-devrel".to_string(),
role_type: RoleType::Engineer,
agent: Some("claude".to_string()),
reports_to: Some("jordan-pm".to_string()),
use_worktrees: false,
..MemberInstance::default()
};
let mut daemon = TestDaemonBuilder::new(tmp.path())
.members(vec![
architect_member("maya-lead"),
manager_member("jordan-pm", Some("maya-lead")),
priya,
kai,
])
.states(HashMap::from([
("priya-writer-1-1".to_string(), MemberState::Idle),
("kai-devrel-1-1".to_string(), MemberState::Idle),
]))
.build();
daemon.enqueue_dispatch_candidates().unwrap();
assert_eq!(daemon.dispatch_queue.len(), 1);
assert_eq!(daemon.dispatch_queue[0].task_id, 91);
assert_eq!(
daemon.dispatch_queue[0].engineer, "priya-writer-1-1",
"body-owner `priya-writer` must win over idle peer kai-devrel-1-1"
);
}
#[test]
fn dispatch_queue_seeds_role_name_into_domain_tags_for_tag_routing() {
let tmp = tempfile::tempdir().unwrap();
write_task_with_tags(tmp.path(), 90, "role-tagged", &["kai-devrel", "engagement"]);
let member_with_role = |name: &str, role_name: &str| MemberInstance {
name: name.to_string(),
role_name: role_name.to_string(),
role_type: RoleType::Engineer,
agent: Some("codex".to_string()),
reports_to: Some("mgr".to_string()),
use_worktrees: false,
..MemberInstance::default()
};
let mut daemon = TestDaemonBuilder::new(tmp.path())
.members(vec![
manager_member("mgr", None),
member_with_role("alex-dev-1-1", "alex-dev"),
member_with_role("kai-devrel-1-1", "kai-devrel"),
member_with_role("sam-designer-1-1", "sam-designer"),
])
.states(HashMap::from([
("alex-dev-1-1".to_string(), MemberState::Idle),
("kai-devrel-1-1".to_string(), MemberState::Idle),
("sam-designer-1-1".to_string(), MemberState::Idle),
]))
.build();
daemon.enqueue_dispatch_candidates().unwrap();
assert_eq!(daemon.dispatch_queue.len(), 1);
assert_eq!(daemon.dispatch_queue[0].task_id, 90);
assert_eq!(
daemon.dispatch_queue[0].engineer, "kai-devrel-1-1",
"task tagged with a role_name must prefer the engineer whose role_name matches, \
not fall back to the first alphabetical peer"
);
}
#[test]
fn dispatch_queue_seeds_role_name_word_family_variants() {
let tmp = tempfile::tempdir().unwrap();
write_task_with_tags(
tmp.path(),
572,
"Card-1 peak-day hero",
&["pillar-a", "design", "thread-a", "hero", "card-1"],
);
let member_with_role = |name: &str, role_name: &str| MemberInstance {
name: name.to_string(),
role_name: role_name.to_string(),
role_type: RoleType::Engineer,
agent: Some("codex".to_string()),
reports_to: Some("mgr".to_string()),
use_worktrees: false,
..MemberInstance::default()
};
let mut daemon = TestDaemonBuilder::new(tmp.path())
.members(vec![
manager_member("mgr", None),
member_with_role("alex-dev-1-1", "alex-dev"),
member_with_role("kai-devrel-1-1", "kai-devrel"),
member_with_role("sam-designer-1-1", "sam-designer"),
member_with_role("priya-writer-1-1", "priya-writer"),
])
.states(HashMap::from([
("alex-dev-1-1".to_string(), MemberState::Idle),
("kai-devrel-1-1".to_string(), MemberState::Idle),
("sam-designer-1-1".to_string(), MemberState::Idle),
("priya-writer-1-1".to_string(), MemberState::Idle),
]))
.build();
daemon.enqueue_dispatch_candidates().unwrap();
assert_eq!(daemon.dispatch_queue.len(), 1);
assert_eq!(daemon.dispatch_queue[0].task_id, 572);
assert_eq!(
daemon.dispatch_queue[0].engineer, "sam-designer-1-1",
"task tagged `design` must prefer sam-designer over alphabetical \
alex-dev; role_name `sam-designer` seeds `designer` + `design` + \
`designing` into domain_tags"
);
}
#[test]
fn role_name_seed_tags_covers_hyphen_suffix_and_er_variants() {
assert_eq!(
role_name_seed_tags("sam-designer"),
vec![
"sam-designer".to_string(),
"designer".to_string(),
"design".to_string(),
"designing".to_string(),
]
);
assert_eq!(
role_name_seed_tags("priya-writer"),
vec![
"priya-writer".to_string(),
"writer".to_string(),
"writ".to_string(),
"writing".to_string(),
]
);
assert_eq!(
role_name_seed_tags("alex-dev"),
vec!["alex-dev".to_string(), "dev".to_string()]
);
assert_eq!(
role_name_seed_tags("kai-devrel"),
vec!["kai-devrel".to_string(), "devrel".to_string()]
);
assert_eq!(
role_name_seed_tags("architect"),
vec!["architect".to_string()]
);
}
#[test]
fn dispatch_honors_explicit_body_owner_when_tags_do_not_match_role() {
let tmp = tempfile::tempdir().unwrap();
write_task_with_tags_and_body(
tmp.path(),
91,
"body-owner-routing",
&["content", "pillar-b", "x", "writing"],
"- Owner: priya-writer drafts; kai-devrel schedules\n\
- Acceptance: ...\n",
);
let member_with_role = |name: &str, role_name: &str| MemberInstance {
name: name.to_string(),
role_name: role_name.to_string(),
role_type: RoleType::Engineer,
agent: Some("codex".to_string()),
reports_to: Some("mgr".to_string()),
use_worktrees: false,
..MemberInstance::default()
};
let mut daemon = TestDaemonBuilder::new(tmp.path())
.members(vec![
manager_member("mgr", None),
member_with_role("kai-devrel-1-1", "kai-devrel"),
member_with_role("priya-writer-1-1", "priya-writer"),
member_with_role("sam-designer-1-1", "sam-designer"),
])
.states(HashMap::from([
("kai-devrel-1-1".to_string(), MemberState::Idle),
("priya-writer-1-1".to_string(), MemberState::Idle),
("sam-designer-1-1".to_string(), MemberState::Idle),
]))
.build();
daemon.enqueue_dispatch_candidates().unwrap();
assert_eq!(daemon.dispatch_queue.len(), 1);
assert_eq!(daemon.dispatch_queue[0].task_id, 91);
assert_eq!(
daemon.dispatch_queue[0].engineer, "priya-writer-1-1",
"task whose body says `Owner: priya-writer` must route to \
priya-writer-1-1 even when frontmatter tags are thematic \
(content/pillar-b/x/writing) and match no role_name"
);
}
#[test]
fn split_acyclic_blocking_ids_rejects_reverse_edge() {
let tmp = tempfile::tempdir().unwrap();
write_task_with_deps(tmp.path(), 553, "pillar-b-thread", &[554]);
write_task_with_deps(tmp.path(), 554, "dev-to-article", &[]);
write_task_with_deps(tmp.path(), 560, "unrelated", &[]);
let board_dir = tmp.path().join(".batty").join("team_config").join("board");
let (safe, rejected) = split_acyclic_blocking_ids(&board_dir, 554, &[553, 560]).unwrap();
assert_eq!(
rejected,
vec![553],
"edge #554 -> #553 must be rejected — #553 already depends on #554"
);
assert_eq!(
safe,
vec![560],
"unrelated #554 -> #560 edge stays — no cycle"
);
}
#[test]
fn split_acyclic_blocking_ids_rejects_transitive_cycle() {
let tmp = tempfile::tempdir().unwrap();
write_task_with_deps(tmp.path(), 100, "a", &[101]);
write_task_with_deps(tmp.path(), 101, "b", &[102]);
write_task_with_deps(tmp.path(), 102, "c", &[]);
let board_dir = tmp.path().join(".batty").join("team_config").join("board");
let (safe, rejected) = split_acyclic_blocking_ids(&board_dir, 102, &[100]).unwrap();
assert!(safe.is_empty(), "#102 -> #100 reaches #102 via #101");
assert_eq!(rejected, vec![100]);
}
#[test]
fn parse_body_owner_role_extracts_first_role_from_prose() {
assert_eq!(
parse_body_owner_role("- Owner: priya-writer drafts; kai-devrel schedules"),
Some("priya-writer".to_string())
);
assert_eq!(
parse_body_owner_role("Owner: **kai-devrel** (explicit)."),
Some("kai-devrel".to_string())
);
assert_eq!(parse_body_owner_role("Owner: TBD"), None);
assert_eq!(parse_body_owner_role("- Owner: tbd\n"), None);
assert_eq!(parse_body_owner_role("Body with no owner line."), None);
}
#[test]
fn parse_body_owner_role_finds_owner_inside_prose_preamble() {
assert_eq!(
parse_body_owner_role(
"**Round-8 task from maya-lead. Owner: alex-dev. Skeptic-defuser artifact for Article A.**"
),
Some("alex-dev".to_string())
);
assert_eq!(
parse_body_owner_role(
"**Round-11 task from maya-lead. Owner: alex-dev. Daily metrics snapshot.**"
),
Some("alex-dev".to_string())
);
assert_eq!(
parse_body_owner_role("CoOwner: rogue-role. Owner: real-role handles it."),
Some("real-role".to_string())
);
}
#[test]
fn parse_body_owner_role_finds_routing_cues_from_jordan_style_bodies() {
assert_eq!(
parse_body_owner_role(
"**Owner routing**: content/submission task — route to priya-writer for draft, kai-devrel for PR submission (he holds the posting/publishing lane)."
),
Some("priya-writer".to_string())
);
assert_eq!(
parse_body_owner_role(
"**Owner routing**: Research + strategic-analysis task, role-flexible. Primary: priya-writer (narrative audit lens) OR kai-devrel (community/conversion lens). NOT Sam (not visual), NOT Alex."
),
Some("priya-writer".to_string())
);
assert_eq!(
parse_body_owner_role("- Route: dispatch to priya-writer."),
Some("priya-writer".to_string())
);
assert_eq!(
parse_body_owner_role("OWNER: alex-dev-1-1 (explicit per Maya directive)."),
Some("alex-dev".to_string())
);
assert_eq!(
parse_body_owner_role("Please assign to kai-devrel for the release post."),
Some("kai-devrel".to_string())
);
}
#[test]
fn enqueue_dispatch_candidates_skips_recently_orphan_rescued_task() {
let tmp = tempfile::tempdir().unwrap();
write_task_with_body(
tmp.path(),
70,
"rescued-task",
"todo",
None,
"Touch src/team/telemetry_db.rs only.",
);
let mut daemon = TestDaemonBuilder::new(tmp.path())
.members(vec![
manager_member("mgr", None),
engineer_member("eng-1", Some("mgr"), false),
])
.states(HashMap::from([("eng-1".to_string(), MemberState::Idle)]))
.build();
daemon.record_task_rescue(70);
daemon.enqueue_dispatch_candidates().unwrap();
assert!(
daemon.dispatch_queue.is_empty(),
"task under orphan-rescue cooldown must stay off the dispatch queue"
);
}
#[test]
fn enqueue_dispatch_candidates_includes_task_after_orphan_rescue_cooldown_expires() {
let tmp = tempfile::tempdir().unwrap();
write_task_with_body(
tmp.path(),
71,
"expired-rescue",
"todo",
None,
"Touch src/team/telemetry_db.rs only.",
);
let mut daemon = TestDaemonBuilder::new(tmp.path())
.members(vec![
manager_member("mgr", None),
engineer_member("eng-1", Some("mgr"), false),
])
.states(HashMap::from([("eng-1".to_string(), MemberState::Idle)]))
.build();
daemon.config.team_config.board.orphan_rescue_cooldown_secs = 0;
daemon.record_task_rescue(71);
daemon.enqueue_dispatch_candidates().unwrap();
assert_eq!(daemon.dispatch_queue.len(), 1);
assert_eq!(daemon.dispatch_queue[0].task_id, 71);
}
#[test]
fn record_task_rescue_grows_cooldown_exponentially_on_repeat() {
let tmp = tempfile::tempdir().unwrap();
let mut daemon = TestDaemonBuilder::new(tmp.path())
.members(vec![
manager_member("mgr", None),
engineer_member("eng-1", Some("mgr"), false),
])
.build();
daemon.record_task_rescue(99);
let first = daemon.recently_rescued_tasks[&99];
assert_eq!(first.count, 1);
daemon.record_task_rescue(99);
let second = daemon.recently_rescued_tasks[&99];
assert_eq!(second.count, 2);
daemon.record_task_rescue(99);
let third = daemon.recently_rescued_tasks[&99];
assert_eq!(third.count, 3);
let base = std::time::Duration::from_secs(100);
assert_eq!(
third.effective_cooldown(base),
std::time::Duration::from_secs(400)
);
for _ in 0..10 {
daemon.record_task_rescue(99);
}
let capped = daemon.recently_rescued_tasks[&99];
assert_eq!(
capped.effective_cooldown(base),
std::time::Duration::from_secs(1600)
);
}
#[test]
fn record_task_rescue_grows_count_across_dispatch_gate_openings() {
use std::time::Duration;
let tmp = tempfile::tempdir().unwrap();
let mut daemon = TestDaemonBuilder::new(tmp.path())
.members(vec![
manager_member("mgr", None),
engineer_member("eng-1", Some("mgr"), false),
])
.build();
daemon.config.team_config.board.orphan_rescue_cooldown_secs = 100;
daemon.record_task_rescue(42);
let record = daemon.recently_rescued_tasks.get_mut(&42).unwrap();
record.last_rescued_at = std::time::Instant::now() - Duration::from_secs(150);
daemon.record_task_rescue(42);
let grown = daemon.recently_rescued_tasks[&42];
assert_eq!(
grown.count, 2,
"rescue after gate-open but inside cascade window must grow the counter"
);
daemon.record_task_rescue(77);
let record = daemon.recently_rescued_tasks.get_mut(&77).unwrap();
record.last_rescued_at = std::time::Instant::now() - Duration::from_secs(500);
daemon.record_task_rescue(77);
let reset = daemon.recently_rescued_tasks[&77];
assert_eq!(
reset.count, 1,
"rescue past cascade window is a new cascade — counter resets"
);
}
#[test]
fn release_exclusion_blocks_redispatch_to_same_engineer_until_window_expires() {
use std::time::Duration;
let tmp = tempfile::tempdir().unwrap();
let mut daemon = TestDaemonBuilder::new(tmp.path())
.members(vec![
manager_member("mgr", None),
engineer_member("eng-1", Some("mgr"), false),
])
.build();
daemon
.config
.team_config
.board
.dispatch_release_exclusion_secs = 300;
assert!(!daemon.is_release_excluded(42, "eng-1"));
daemon.record_task_release_by(42, "eng-1");
assert!(
daemon.is_release_excluded(42, "eng-1"),
"engineer who just released task must be excluded"
);
assert!(
!daemon.is_release_excluded(42, "eng-2"),
"exclusion is per-(task, engineer), not global"
);
assert!(
!daemon.is_release_excluded(99, "eng-1"),
"exclusion must not leak to other tasks"
);
daemon.recently_released_by.insert(
(42, "eng-1".to_string()),
crate::team::daemon::ReleaseRecord {
last_released_at: std::time::Instant::now() - Duration::from_secs(400),
count: 1,
},
);
assert!(
!daemon.is_release_excluded(42, "eng-1"),
"exclusion must expire after dispatch_release_exclusion_secs"
);
}
#[test]
fn record_task_release_by_grows_exclusion_exponentially_on_repeat() {
use std::time::Duration;
let tmp = tempfile::tempdir().unwrap();
let mut daemon = TestDaemonBuilder::new(tmp.path())
.members(vec![
manager_member("mgr", None),
engineer_member("eng-1", Some("mgr"), false),
])
.build();
daemon
.config
.team_config
.board
.dispatch_release_exclusion_secs = 300;
daemon.record_task_release_by(555, "eng-1");
assert_eq!(
daemon.recently_released_by[&(555, "eng-1".to_string())].count,
1
);
daemon.recently_released_by.insert(
(555, "eng-1".to_string()),
crate::team::daemon::ReleaseRecord {
last_released_at: std::time::Instant::now() - Duration::from_secs(310),
count: 1,
},
);
daemon.record_task_release_by(555, "eng-1");
let after_second = daemon.recently_released_by[&(555, "eng-1".to_string())];
assert_eq!(
after_second.count, 2,
"second release within cascade window must grow the counter"
);
assert!(daemon.is_release_excluded(555, "eng-1"));
daemon.recently_released_by.insert(
(555, "eng-1".to_string()),
crate::team::daemon::ReleaseRecord {
last_released_at: std::time::Instant::now() - Duration::from_secs(3_000),
count: 3,
},
);
daemon.record_task_release_by(555, "eng-1");
let after_reset = daemon.recently_released_by[&(555, "eng-1".to_string())];
assert_eq!(
after_reset.count, 1,
"release past cascade window is a new cascade — counter resets"
);
}
#[test]
fn enqueue_dispatch_candidates_defers_release_excluded_task_and_dispatches_next_priority() {
let tmp = tempfile::tempdir().unwrap();
write_task_with_body(
tmp.path(),
101,
"ci-badge-repair",
"todo",
None,
"- Owner: alex-dev. Prefer CI job repair over badge swap.\n",
);
write_task_with_body(
tmp.path(),
102,
"twir-submission",
"todo",
None,
"- Route: dispatch to priya-writer.\n",
);
let alex = MemberInstance {
name: "alex-dev-1-1".to_string(),
role_name: "alex-dev".to_string(),
role_type: RoleType::Engineer,
agent: Some("claude".to_string()),
reports_to: Some("mgr".to_string()),
use_worktrees: false,
..MemberInstance::default()
};
let priya = MemberInstance {
name: "priya-writer-1-1".to_string(),
role_name: "priya-writer".to_string(),
role_type: RoleType::Engineer,
agent: Some("claude".to_string()),
reports_to: Some("mgr".to_string()),
use_worktrees: false,
..MemberInstance::default()
};
let mut daemon = TestDaemonBuilder::new(tmp.path())
.members(vec![manager_member("mgr", None), alex, priya])
.states(HashMap::from([
("alex-dev-1-1".to_string(), MemberState::Idle),
("priya-writer-1-1".to_string(), MemberState::Idle),
]))
.build();
daemon.record_task_release_by(101, "alex-dev-1-1");
assert!(daemon.is_release_excluded(101, "alex-dev-1-1"));
daemon.enqueue_dispatch_candidates().unwrap();
assert_eq!(
daemon.dispatch_queue.len(),
1,
"release-excluded top-priority task must not starve lower candidates"
);
assert_eq!(daemon.dispatch_queue[0].task_id, 102);
assert_eq!(daemon.dispatch_queue[0].engineer, "priya-writer-1-1");
}
#[test]
fn enqueue_dispatch_candidates_serializes_overlapping_task_and_enqueues_non_overlapping_task() {
let tmp = tempfile::tempdir().unwrap();
write_task_with_body(
tmp.path(),
10,
"active-overlap",
"in-progress",
Some("eng-2"),
"Modify src/team/dispatch/queue.rs and tests.",
);
write_task_with_body(
tmp.path(),
11,
"candidate-overlap",
"todo",
None,
"Update src/team/dispatch/queue.rs overlap logic.",
);
write_task_with_body(
tmp.path(),
12,
"candidate-safe",
"todo",
None,
"Touch src/team/telemetry_db.rs only.",
);
let mut daemon = TestDaemonBuilder::new(tmp.path())
.members(vec![
manager_member("mgr", None),
engineer_member("eng-1", Some("mgr"), false),
engineer_member("eng-2", Some("mgr"), false),
])
.states(HashMap::from([
("eng-1".to_string(), MemberState::Idle),
("eng-2".to_string(), MemberState::Working),
]))
.build();
daemon.enqueue_dispatch_candidates().unwrap();
assert_eq!(daemon.dispatch_queue.len(), 1);
assert_eq!(daemon.dispatch_queue[0].task_id, 12);
let task = crate::task::Task::from_file(
&tmp.path()
.join(".batty")
.join("team_config")
.join("board")
.join("tasks")
.join("011-candidate-overlap.md"),
)
.unwrap();
assert_eq!(task.depends_on, vec![10]);
let events =
crate::team::events::read_events(&crate::team::team_events_path(tmp.path())).unwrap();
assert!(events.iter().any(|event| {
event.event == "dispatch_overlap_skipped" && event.task.as_deref() == Some("11")
}));
}
#[test]
fn enqueue_dispatch_candidates_leaves_serialized_task_unqueued_when_no_safe_alternative_exists()
{
let tmp = tempfile::tempdir().unwrap();
write_task_with_body(
tmp.path(),
20,
"active-overlap",
"in-progress",
Some("eng-2"),
"Modify src/team/dispatch/mod.rs.",
);
write_task_with_body(
tmp.path(),
21,
"candidate-overlap",
"todo",
None,
"Also update src/team/dispatch/mod.rs for prevention logic.",
);
let mut daemon = TestDaemonBuilder::new(tmp.path())
.members(vec![
manager_member("mgr", None),
engineer_member("eng-1", Some("mgr"), false),
engineer_member("eng-2", Some("mgr"), false),
])
.states(HashMap::from([
("eng-1".to_string(), MemberState::Idle),
("eng-2".to_string(), MemberState::Working),
]))
.build();
daemon.enqueue_dispatch_candidates().unwrap();
assert!(daemon.dispatch_queue.is_empty());
let task = crate::task::Task::from_file(
&tmp.path()
.join(".batty")
.join("team_config")
.join("board")
.join("tasks")
.join("021-candidate-overlap.md"),
)
.unwrap();
assert_eq!(task.depends_on, vec![20]);
}
#[test]
fn test_predicted_files_from_body() {
let tmp = tempfile::tempdir().unwrap();
write_task_with_body(
tmp.path(),
30,
"body-path",
"todo",
None,
"Update src/team/daemon.rs to add the new check.",
);
let task = crate::task::Task::from_file(
&tmp.path()
.join(".batty")
.join("team_config")
.join("board")
.join("tasks")
.join("030-body-path.md"),
)
.unwrap();
assert!(predicted_files(&task, tmp.path()).contains(&"src/team/daemon.rs".to_string()));
}
#[test]
fn test_predicted_files_from_frontmatter_files() {
let tmp = tempfile::tempdir().unwrap();
write_task_with_files(
tmp.path(),
31,
"frontmatter-paths",
"todo",
None,
&["src/app.rs", "src/**/*.rs"],
"Use the declared file list.",
);
let task = crate::task::Task::from_file(
&tmp.path()
.join(".batty")
.join("team_config")
.join("board")
.join("tasks")
.join("031-frontmatter-paths.md"),
)
.unwrap();
let predicted = predicted_files(&task, tmp.path());
assert!(predicted.contains(&"src/**/*.rs".to_string()));
assert!(predicted.contains(&"src/app.rs".to_string()));
}
#[test]
fn test_find_overlapping_with_frontmatter_glob() {
let tmp = tempfile::tempdir().unwrap();
write_task_with_files(
tmp.path(),
40,
"active-glob",
"in-progress",
Some("eng-2"),
&["src/**/*.rs"],
"Broad source lock.",
);
write_task_with_body(
tmp.path(),
41,
"candidate-file",
"todo",
None,
"Change src/app.rs only.",
);
let mut active = None;
let mut candidate = None;
for task in crate::task::load_tasks_from_dir(
&tmp.path()
.join(".batty")
.join("team_config")
.join("board")
.join("tasks"),
)
.unwrap()
{
match task.id {
40 => active = Some(task),
41 => candidate = Some(task),
_ => {}
}
}
let conflicts = find_overlapping_tasks(
&candidate.expect("candidate task"),
&[active.expect("active task")],
tmp.path(),
);
assert_eq!(conflicts.len(), 1);
assert_eq!(
conflicts[0].conflicting_files,
vec!["src/app.rs".to_string()]
);
}
#[test]
fn test_predicted_files_from_tags() {
let tmp = tempfile::tempdir().unwrap();
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("001-prior-shim.md"),
"---\nid: 1\ntitle: prior shim\nstatus: done\npriority: high\nclaimed_by: eng-2\ntags:\n - shim\nchanged_paths:\n - src/shim/runtime.rs\nclass: standard\n---\n\nEarlier shim work.\n",
)
.unwrap();
std::fs::write(
tasks_dir.join("031-new-shim.md"),
"---\nid: 31\ntitle: new shim\nstatus: todo\npriority: high\ntags:\n - shim\nclass: standard\n---\n\nNo explicit paths.\n",
)
.unwrap();
let task = crate::task::Task::from_file(&tasks_dir.join("031-new-shim.md")).unwrap();
assert!(predicted_files(&task, tmp.path()).contains(&"src/shim/runtime.rs".to_string()));
}
#[test]
fn test_predicted_files_empty() {
let tmp = tempfile::tempdir().unwrap();
write_task_with_body(
tmp.path(),
32,
"no-hints",
"todo",
None,
"No file hints here.",
);
let task = crate::task::Task::from_file(
&tmp.path()
.join(".batty")
.join("team_config")
.join("board")
.join("tasks")
.join("032-no-hints.md"),
)
.unwrap();
assert!(predicted_files(&task, tmp.path()).is_empty());
}
#[test]
fn test_find_overlapping_no_conflict() {
let tmp = tempfile::tempdir().unwrap();
write_task_with_body(
tmp.path(),
40,
"candidate",
"todo",
None,
"Edit src/team/daemon.rs.",
);
write_task_with_body(
tmp.path(),
41,
"active",
"in-progress",
Some("eng-2"),
"Edit src/team/status.rs.",
);
let tasks_dir = tmp
.path()
.join(".batty")
.join("team_config")
.join("board")
.join("tasks");
let candidate = crate::task::Task::from_file(&tasks_dir.join("040-candidate.md")).unwrap();
let active = crate::task::Task::from_file(&tasks_dir.join("041-active.md")).unwrap();
assert!(find_overlapping_tasks(&candidate, &[active], tmp.path()).is_empty());
}
#[test]
fn test_find_overlapping_with_conflict() {
let tmp = tempfile::tempdir().unwrap();
write_task_with_body(
tmp.path(),
42,
"candidate",
"todo",
None,
"Edit src/team/daemon.rs.",
);
write_task_with_body(
tmp.path(),
43,
"active",
"in-progress",
Some("eng-2"),
"Edit src/team/daemon.rs too.",
);
let tasks_dir = tmp
.path()
.join(".batty")
.join("team_config")
.join("board")
.join("tasks");
let candidate = crate::task::Task::from_file(&tasks_dir.join("042-candidate.md")).unwrap();
let active = crate::task::Task::from_file(&tasks_dir.join("043-active.md")).unwrap();
let conflicts = find_overlapping_tasks(&candidate, &[active], tmp.path());
assert_eq!(conflicts.len(), 1);
assert_eq!(conflicts[0].task_id, "43");
}
#[test]
fn test_find_overlapping_multiple_conflicts() {
let tmp = tempfile::tempdir().unwrap();
write_task_with_body(
tmp.path(),
44,
"candidate",
"todo",
None,
"Edit src/team/daemon.rs.",
);
write_task_with_body(
tmp.path(),
45,
"active-a",
"in-progress",
Some("eng-2"),
"Edit src/team/daemon.rs too.",
);
write_task_with_body(
tmp.path(),
46,
"active-b",
"in-progress",
Some("eng-3"),
"Also touch src/team/daemon.rs.",
);
let tasks_dir = tmp
.path()
.join(".batty")
.join("team_config")
.join("board")
.join("tasks");
let candidate = crate::task::Task::from_file(&tasks_dir.join("044-candidate.md")).unwrap();
let active_a = crate::task::Task::from_file(&tasks_dir.join("045-active-a.md")).unwrap();
let active_b = crate::task::Task::from_file(&tasks_dir.join("046-active-b.md")).unwrap();
let conflicts = find_overlapping_tasks(&candidate, &[active_a, active_b], tmp.path());
assert_eq!(conflicts.len(), 2);
}
#[test]
fn test_dispatch_skips_overlapping_candidate() {
let tmp = tempfile::tempdir().unwrap();
write_task_with_body(
tmp.path(),
50,
"active",
"in-progress",
Some("eng-2"),
"Edit src/team/daemon.rs.",
);
write_task_with_body(
tmp.path(),
51,
"overlap",
"todo",
None,
"Edit src/team/daemon.rs.",
);
write_task_with_body(
tmp.path(),
52,
"safe",
"todo",
None,
"Edit src/team/status.rs.",
);
let mut daemon = TestDaemonBuilder::new(tmp.path())
.members(vec![
manager_member("mgr", None),
engineer_member("eng-1", Some("mgr"), false),
engineer_member("eng-2", Some("mgr"), false),
])
.states(HashMap::from([
("eng-1".to_string(), MemberState::Idle),
("eng-2".to_string(), MemberState::Working),
]))
.build();
daemon.enqueue_dispatch_candidates().unwrap();
assert_eq!(daemon.dispatch_queue.len(), 1);
assert_eq!(daemon.dispatch_queue[0].task_id, 52);
}
#[test]
fn enqueue_dispatch_candidates_skips_benched_engineer() {
let tmp = tempfile::tempdir().unwrap();
write_open_task_file(tmp.path(), 70, "dispatchable", "todo");
write_bench_test_team_config(tmp.path(), 2);
let mut daemon = TestDaemonBuilder::new(tmp.path())
.members(vec![
engineer_member("eng-1", None, false),
engineer_member("eng-2", None, false),
])
.states(HashMap::from([
("eng-1".to_string(), MemberState::Idle),
("eng-2".to_string(), MemberState::Idle),
]))
.build();
crate::team::bench::bench_engineer(tmp.path(), "eng-1", Some("session end")).unwrap();
daemon.enqueue_dispatch_candidates().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, 70);
}
#[test]
fn enqueue_dispatch_candidates_skips_quota_parked_engineer() {
let tmp = tempfile::tempdir().unwrap();
write_open_task_file(tmp.path(), 674, "dispatchable", "todo");
let mut daemon = TestDaemonBuilder::new(tmp.path())
.members(vec![
engineer_member("eng-1", None, false),
engineer_member("eng-2", None, false),
])
.states(HashMap::from([
("eng-1".to_string(), MemberState::Idle),
("eng-2".to_string(), MemberState::Idle),
]))
.build();
let future_deadline = std::time::SystemTime::now()
.duration_since(std::time::UNIX_EPOCH)
.unwrap_or(std::time::Duration::ZERO)
.as_secs()
+ 32 * 3600;
daemon
.backend_quota_retry_at
.insert("eng-1".to_string(), future_deadline);
daemon.enqueue_dispatch_candidates().unwrap();
assert_eq!(daemon.dispatch_queue.len(), 1);
assert_eq!(
daemon.dispatch_queue[0].engineer, "eng-2",
"quota-parked engineer must be skipped even if cached health is Healthy"
);
assert_eq!(daemon.dispatch_queue[0].task_id, 674);
}
#[test]
fn unbench_restores_dispatch_eligibility() {
let tmp = tempfile::tempdir().unwrap();
write_open_task_file(tmp.path(), 71, "dispatchable", "todo");
write_bench_test_team_config(tmp.path(), 2);
let mut daemon = TestDaemonBuilder::new(tmp.path())
.members(vec![
engineer_member("eng-1", None, false),
engineer_member("eng-2", None, false),
])
.states(HashMap::from([
("eng-1".to_string(), MemberState::Idle),
("eng-2".to_string(), MemberState::Idle),
]))
.build();
crate::team::bench::bench_engineer(tmp.path(), "eng-1", Some("pause")).unwrap();
daemon.enqueue_dispatch_candidates().unwrap();
assert_eq!(daemon.dispatch_queue.len(), 1);
assert_eq!(daemon.dispatch_queue[0].engineer, "eng-2");
crate::team::bench::unbench_engineer(tmp.path(), "eng-1").unwrap();
daemon.dispatch_queue.clear();
daemon.enqueue_dispatch_candidates().unwrap();
assert_eq!(daemon.dispatch_queue.len(), 1);
assert_eq!(daemon.dispatch_queue[0].engineer, "eng-1");
assert_eq!(daemon.dispatch_queue[0].task_id, 71);
}
#[test]
fn test_dispatch_all_overlap_picks_least_conflict() {
let tmp = tempfile::tempdir().unwrap();
write_task_with_body(
tmp.path(),
60,
"active-daemon",
"in-progress",
Some("eng-2"),
"Edit src/team/daemon.rs.",
);
write_task_with_body(
tmp.path(),
61,
"active-status",
"in-progress",
Some("eng-3"),
"Edit src/team/status.rs.",
);
write_task_with_body(
tmp.path(),
62,
"overlap-both",
"todo",
None,
"Edit src/team/daemon.rs and src/team/status.rs.",
);
write_task_with_body(
tmp.path(),
63,
"overlap-one",
"todo",
None,
"Edit src/team/daemon.rs only.",
);
let mut daemon = TestDaemonBuilder::new(tmp.path())
.members(vec![
manager_member("mgr", None),
engineer_member("eng-1", Some("mgr"), false),
engineer_member("eng-2", Some("mgr"), false),
engineer_member("eng-3", Some("mgr"), false),
])
.states(HashMap::from([
("eng-1".to_string(), MemberState::Idle),
("eng-2".to_string(), MemberState::Working),
("eng-3".to_string(), MemberState::Working),
]))
.build();
daemon.enqueue_dispatch_candidates().unwrap();
assert!(daemon.dispatch_queue.is_empty());
let task = crate::task::Task::from_file(
&tmp.path()
.join(".batty")
.join("team_config")
.join("board")
.join("tasks")
.join("063-overlap-one.md"),
)
.unwrap();
assert_eq!(task.depends_on, vec![60]);
}
#[test]
fn test_overlap_conflict_struct_fields() {
let conflict = OverlapConflict {
task_id: "42".to_string(),
conflicting_files: vec!["src/team/daemon.rs".to_string()],
in_progress_engineer: "eng-2".to_string(),
};
assert_eq!(conflict.task_id, "42");
assert_eq!(conflict.conflicting_files, vec!["src/team/daemon.rs"]);
assert_eq!(conflict.in_progress_engineer, "eng-2");
}
#[test]
fn file_level_locks_defer_then_release_overlapping_work_for_worktree_teams() {
let tmp = tempfile::tempdir().unwrap();
write_task_with_body(
tmp.path(),
60,
"active-overlap",
"in-progress",
Some("eng-2"),
"Modify src/app.rs only.",
);
write_task_with_files(
tmp.path(),
61,
"waiting-overlap",
"todo",
None,
&["src/*.rs"],
"Lock should wait without rewriting dependencies.",
);
let mut daemon = TestDaemonBuilder::new(tmp.path())
.members(vec![
manager_member("mgr", None),
engineer_member("eng-1", Some("mgr"), true),
engineer_member("eng-2", Some("mgr"), true),
])
.workflow_policy(crate::team::config::WorkflowPolicy {
file_level_locks: true,
..crate::team::config::WorkflowPolicy::default()
})
.states(HashMap::from([
("eng-1".to_string(), MemberState::Idle),
("eng-2".to_string(), MemberState::Working),
]))
.build();
daemon.enqueue_dispatch_candidates().unwrap();
assert!(daemon.dispatch_queue.is_empty());
let waiting_task = crate::task::Task::from_file(
&tmp.path()
.join(".batty")
.join("team_config")
.join("board")
.join("tasks")
.join("061-waiting-overlap.md"),
)
.unwrap();
assert!(
waiting_task.depends_on.is_empty(),
"file-level wait should not rewrite task dependencies"
);
write_task_with_body(
tmp.path(),
60,
"active-overlap",
"done",
Some("eng-2"),
"Modify src/app.rs only.",
);
daemon.enqueue_dispatch_candidates().unwrap();
assert_eq!(daemon.dispatch_queue.len(), 1);
assert_eq!(daemon.dispatch_queue[0].task_id, 61);
}
}