use anyhow::{Context, Result, bail};
use batty_cli::{
agent,
cli::{
self, ActivityCommand, AutoMergeAction, BoardCommand, Cli, Command, DepsFormatArg,
DiscordCommand, GrafanaCommand, InboxCommand, NudgeCommand, OpenClawCommand,
OpenClawEventTopicArg, OpenClawFollowUpCommand, ProjectCommand, ResearchCommand,
ResearchFormatArg, ResearchKeepPolicyArg, ReviewDispositionArg, TaskCommand, TaskStateArg,
},
env_file, project_registry, release, team,
};
use clap::Parser;
use dialoguer::{Confirm, Input, Select};
use std::collections::HashMap;
use std::path::PathBuf;
use tracing::debug;
fn project_root() -> PathBuf {
let cwd = std::env::current_dir().unwrap_or_else(|_| PathBuf::from("."));
if let Ok(output) = std::process::Command::new("git")
.args(["rev-parse", "--git-common-dir"])
.current_dir(&cwd)
.output()
{
if output.status.success() {
let git_common = String::from_utf8_lossy(&output.stdout).trim().to_string();
let git_path = if std::path::Path::new(&git_common).is_absolute() {
PathBuf::from(&git_common)
} else {
cwd.join(&git_common)
};
if let Some(repo_root) = git_path.parent() {
if let Ok(canonical) = repo_root.canonicalize() {
return canonical;
}
}
}
}
cwd
}
fn setup_tracing(verbose: u8) {
let filter = match verbose {
0 => "warn",
1 => "info",
2 => "debug",
_ => "trace",
};
let env_filter = match std::env::var("BATTY_LOG") {
Ok(value) if !value.trim().is_empty() => tracing_subscriber::EnvFilter::new(value),
_ => tracing_subscriber::EnvFilter::try_from_default_env()
.unwrap_or_else(|_| tracing_subscriber::EnvFilter::new(filter)),
};
tracing_subscriber::fmt()
.with_env_filter(env_filter)
.with_writer(std::io::stderr)
.init();
}
fn format_ts(unix_secs: i64) -> String {
use std::time::{Duration, UNIX_EPOCH};
let dt = UNIX_EPOCH + Duration::from_secs(unix_secs as u64);
let datetime: chrono::DateTime<chrono::Local> = dt.into();
datetime.format("%Y-%m-%d %H:%M:%S").to_string()
}
fn task_state_arg_name(state: TaskStateArg) -> &'static str {
match state {
TaskStateArg::Backlog => "backlog",
TaskStateArg::Todo => "todo",
TaskStateArg::InProgress => "in-progress",
TaskStateArg::Review => "review",
TaskStateArg::Blocked => "blocked",
TaskStateArg::Done => "done",
TaskStateArg::Archived => "archived",
}
}
fn review_disposition_arg_name(disposition: ReviewDispositionArg) -> &'static str {
match disposition {
ReviewDispositionArg::Approved => "approved",
ReviewDispositionArg::ChangesRequested => "changes_requested",
ReviewDispositionArg::Rejected => "rejected",
}
}
fn board_summary_counts(board_dir: &std::path::Path) -> Result<Vec<(&'static str, usize)>> {
const STATUSES: [&str; 7] = [
"backlog",
"todo",
"in-progress",
"review",
"blocked",
"done",
"archived",
];
STATUSES
.into_iter()
.map(|status| {
let output = team::board_cmd::list_tasks(board_dir, Some(status))?;
Ok((status, count_board_list_rows(&output)))
})
.collect()
}
fn count_board_list_rows(output: &str) -> usize {
output
.lines()
.filter(|line| {
line.split_whitespace()
.next()
.is_some_and(|token| token.chars().all(|ch| ch.is_ascii_digit()))
})
.count()
}
fn confirm(prompt: &str, default: bool) -> Result<bool> {
Ok(Confirm::new()
.with_prompt(prompt)
.default(default)
.interact()?)
}
fn input_u64(prompt: &str, default: u64) -> Result<u64> {
let s: String = Input::new()
.with_prompt(prompt)
.default(default.to_string())
.interact_text()?;
Ok(s.parse::<u64>().unwrap_or(default))
}
fn collect_init_overrides() -> Result<team::InitOverrides> {
let mut ov = team::InitOverrides::default();
println!();
println!("── Orchestrator ──────────────────────────────────────────");
println!("The orchestrator is a dedicated tmux pane that runs automated");
println!("triage, review routing, dispatch-gap recovery, and standups.");
ov.orchestrator_pane = Some(confirm("Enable orchestrator pane?", true)?);
println!();
println!("── Board & Dispatch ──────────────────────────────────────");
println!("Auto-dispatch automatically assigns todo tasks from the board");
println!("to idle engineers without manual intervention.");
ov.auto_dispatch = Some(confirm("Enable auto-dispatch of board tasks?", true)?);
println!();
println!("── Engineer Worktrees ────────────────────────────────────");
println!("When enabled, each engineer gets an isolated git worktree.");
println!("New task assignments create a fresh branch in that worktree,");
println!("keeping engineers from stepping on each other's changes.");
ov.use_worktrees = Some(confirm("Enable git worktrees for engineers?", true)?);
println!();
println!("── Automation: Nudges & Standups ─────────────────────────");
println!("Timeout nudges ping agents that appear stuck or idle for too long.");
ov.timeout_nudges = Some(confirm("Enable timeout nudges?", true)?);
println!("Standups periodically ask agents to report their status.");
ov.standups = Some(confirm("Enable periodic standups?", true)?);
println!();
println!("── Automation: Interventions ──────────────────────────────");
println!("These are daemon-driven actions that keep the team moving.");
println!();
println!("Triage: auto-routes new tasks to the right manager/engineer.");
ov.triage_interventions = Some(confirm("Enable triage interventions?", true)?);
println!("Review: nudges reviewers and escalates stale reviews.");
ov.review_interventions = Some(confirm("Enable review interventions?", true)?);
println!("Owned-task recovery: re-dispatches tasks stuck on a dead agent.");
ov.owned_task_interventions = Some(confirm("Enable owned-task recovery?", true)?);
println!("Manager dispatch: nudges managers when todo tasks pile up.");
ov.manager_dispatch_interventions =
Some(confirm("Enable manager dispatch interventions?", true)?);
println!("Architect utilization: nudges the architect when engineers are idle.");
ov.architect_utilization_interventions = Some(confirm(
"Enable architect utilization interventions?",
true,
)?);
println!();
println!("── Auto-Merge ───────────────────────────────────────────");
println!("When enabled, completed engineer branches that pass tests and");
println!("score above a confidence threshold are merged automatically.");
ov.auto_merge_enabled = Some(confirm("Enable auto-merge?", false)?);
println!();
println!("── Timing (seconds) ─────────────────────────────────────");
println!("These control how often the daemon checks on agents and reviews.");
println!();
ov.standup_interval_secs = Some(input_u64(
"Standup interval (secs, how often agents report status)",
600,
)?);
ov.nudge_interval_secs = Some(input_u64(
"Architect nudge interval (secs, idle ping for architect)",
900,
)?);
ov.stall_threshold_secs = Some(input_u64(
"Stall threshold (secs, agent considered stuck after this)",
300,
)?);
ov.review_nudge_threshold_secs = Some(input_u64(
"Review nudge threshold (secs, reviewer gets a reminder)",
1800,
)?);
ov.review_timeout_secs = Some(input_u64(
"Review timeout (secs, stale review gets escalated)",
7200,
)?);
println!();
Ok(ov)
}
fn main() -> Result<()> {
let cli = Cli::parse();
setup_tracing(cli.verbose);
let root = project_root();
env_file::load_project_env(&root)?;
debug!(root = %root.display(), "project root");
match cli.command {
Command::Init {
template,
from,
force,
agent,
} => {
if let Some(ref name) = agent {
if agent::adapter_from_name(name).is_none() {
bail!(
"unknown agent backend '{name}'. Supported: {}",
agent::KNOWN_AGENT_NAMES.join(", ")
);
}
}
let created = if let Some(template_name) = from.as_deref() {
team::init_from_template(&root, template_name)?
} else {
let template_name = match template.unwrap_or(cli::InitTemplate::Simple) {
cli::InitTemplate::Solo => "solo",
cli::InitTemplate::Pair => "pair",
cli::InitTemplate::Simple => "simple",
cli::InitTemplate::Squad => "squad",
cli::InitTemplate::Large => "large",
cli::InitTemplate::Research => "research",
cli::InitTemplate::Software => "software",
cli::InitTemplate::Cleanroom => "cleanroom",
cli::InitTemplate::Batty => "batty",
cli::InitTemplate::Python => "python",
};
let default_name = root
.file_name()
.map(|n| n.to_string_lossy().to_string())
.unwrap_or_else(|| "my-project".to_string());
let project_name: String = Input::new()
.with_prompt("Project name")
.default(default_name)
.interact_text()?;
let selected_agent = if let Some(agent_name) = agent.as_deref() {
agent_name
} else {
let agents = agent::KNOWN_AGENT_NAMES;
let default_idx = agents.iter().position(|&a| a == "claude").unwrap_or(0);
let agent_idx = Select::new()
.with_prompt("Agent backend")
.items(agents)
.default(default_idx)
.interact()?;
agents[agent_idx]
};
let overrides = collect_init_overrides()?;
team::init_team_with_overrides(
&root,
template_name,
Some(&project_name),
Some(selected_agent),
force,
Some(&overrides),
)?
};
println!("Initialized team config ({} files):", created.len());
for path in &created {
println!(" {}", path.display());
}
println!();
println!("Edit .batty/team_config/team.yaml to configure your team.");
println!("Then run: batty start");
}
Command::ExportTemplate { name } => {
let count = team::export_template(&root, &name)?;
println!("Exported template '{name}' ({count} files)");
}
Command::ExportRun => {
let path = team::export_run(&root)?;
println!("Run export written to {}", path.display());
}
Command::Retro { events } => {
let stats = if let Some(events_path) = events {
team::retrospective::analyze_event_log(&events_path)?
} else {
team::retrospective::analyze_project(&root)?
};
match stats {
Some(stats) => {
let path = team::retrospective::generate_retrospective(&root, &stats)?;
println!("Retrospective written to {}", path.display());
}
None => println!("No run data found in event log."),
}
}
Command::Start { attach, quiet } => {
let session = team::start_team(&root, attach)?;
if !attach && !quiet {
println!("Team session started: {session}");
println!("Run `batty attach` to connect.");
}
}
Command::Stop => {
team::stop_team(&root)?;
println!("Team session stopped.");
}
Command::Attach => {
team::attach_team(&root)?;
}
Command::Status {
json,
detail,
health,
} => {
team::team_status(&root, json, detail, health)?;
}
Command::Bench { engineer, reason } => {
let entry = team::bench::bench_engineer(&root, &engineer, reason.as_deref())?;
println!(
"Benched {} at {}{}",
engineer,
entry.timestamp,
entry
.reason
.as_deref()
.map(|reason| format!(" ({reason})"))
.unwrap_or_default()
);
}
Command::Unbench { engineer } => {
if team::bench::unbench_engineer(&root, &engineer)? {
println!("Unbenched {engineer}.");
} else {
println!("{engineer} was not benched.");
}
}
Command::OpenClaw { command } => match command {
OpenClawCommand::Register { force } => {
let path = team::openclaw::register_project(&root, force)?;
println!("OpenClaw project config written to {}", path.display());
}
OpenClawCommand::Status { json } => {
team::openclaw::openclaw_status(&root, json)?;
}
OpenClawCommand::Instruct { role, message } => {
team::openclaw::send_openclaw_instruction(&root, &role, &message)?;
println!("OpenClaw instruction queued for {role}.");
}
OpenClawCommand::Events {
project_id,
all_projects,
json,
topics,
roles,
task_ids,
event_types,
session_names,
since_ts,
limit,
include_archived,
} => {
let subscription = team::openclaw::OpenClawEventSubscription {
topics: topics
.into_iter()
.map(|topic| match topic {
OpenClawEventTopicArg::Completion => {
team::openclaw::OpenClawEventTopic::Completion
}
OpenClawEventTopicArg::Review => {
team::openclaw::OpenClawEventTopic::Review
}
OpenClawEventTopicArg::Stall => {
team::openclaw::OpenClawEventTopic::Stall
}
OpenClawEventTopicArg::Merge => {
team::openclaw::OpenClawEventTopic::Merge
}
OpenClawEventTopicArg::Escalation => {
team::openclaw::OpenClawEventTopic::Escalation
}
OpenClawEventTopicArg::DeliveryFailure => {
team::openclaw::OpenClawEventTopic::DeliveryFailure
}
OpenClawEventTopicArg::Lifecycle => {
team::openclaw::OpenClawEventTopic::Lifecycle
}
})
.collect(),
project_ids: Vec::new(),
session_names,
roles,
task_ids,
event_types,
since_ts,
limit,
include_archived,
};
team::openclaw::openclaw_events(
&root,
&subscription,
project_id.as_deref(),
all_projects,
json,
)?;
}
OpenClawCommand::FollowUp { command } => match command {
OpenClawFollowUpCommand::Run { json } => {
team::openclaw::run_follow_ups(&root, json)?;
}
},
},
Command::Project { command } => match command {
ProjectCommand::Register {
project_id,
name,
aliases,
project_root,
board_dir,
team_name,
session_name,
owner,
tags,
channel_bindings,
thread_bindings,
allow_openclaw_supervision,
allow_cross_project_routing,
allow_shared_service_routing,
archived,
json,
} => {
let mut channel_bindings = channel_bindings
.iter()
.map(|binding| project_registry::parse_channel_binding(binding))
.collect::<Result<Vec<_>>>()?;
channel_bindings.extend(
thread_bindings
.iter()
.map(|binding| project_registry::parse_thread_binding(binding))
.collect::<Result<Vec<_>>>()?,
);
let project =
project_registry::register_project(project_registry::ProjectRegistration {
project_id,
name,
aliases,
project_root,
board_dir,
team_name,
session_name,
channel_bindings,
owner,
tags,
policy_flags: project_registry::ProjectPolicyFlags {
allow_openclaw_supervision,
allow_cross_project_routing,
allow_shared_service_routing,
archived,
},
})?;
if json {
println!("{}", serde_json::to_string_pretty(&project)?);
} else {
println!(
"Registered project {} at {}",
project.project_id,
project.project_root.display()
);
}
}
ProjectCommand::Unregister { project_id, json } => {
let Some(project) = project_registry::unregister_project(&project_id)? else {
bail!("project '{}' is not registered", project_id);
};
if json {
println!("{}", serde_json::to_string_pretty(&project)?);
} else {
println!("Unregistered project {}", project.project_id);
}
}
ProjectCommand::List { json } => {
let projects = project_registry::list_projects()?;
if json {
println!("{}", serde_json::to_string_pretty(&projects)?);
} else if projects.is_empty() {
println!("No projects registered.");
} else {
println!(
"{:<20} {:<24} {:<20} {:<18}",
"PROJECT_ID", "NAME", "TEAM", "SESSION"
);
for project in &projects {
println!(
"{:<20} {:<24} {:<20} {:<18}",
project.project_id,
project.name,
project.team_name,
project.session_name
);
}
}
}
ProjectCommand::Get { project_id, json } => {
let Some(project) = project_registry::get_project(&project_id)? else {
bail!("project '{}' is not registered", project_id);
};
if json {
println!("{}", serde_json::to_string_pretty(&project)?);
} else {
println!("Project: {}", project.project_id);
println!("Name: {}", project.name);
println!("Project root: {}", project.project_root.display());
println!("Board dir: {}", project.board_dir.display());
println!("Team: {}", project.team_name);
println!("Session: {}", project.session_name);
println!("Owner: {}", project.owner.as_deref().unwrap_or("(unowned)"));
println!(
"Tags: {}",
if project.tags.is_empty() {
"(none)".to_string()
} else {
project.tags.join(", ")
}
);
println!(
"Aliases: {}",
if project.aliases.is_empty() {
"(none)".to_string()
} else {
project.aliases.join(", ")
}
);
}
}
ProjectCommand::Start { project_id, json } => {
let result = project_registry::start_project(&project_id)?;
if json {
println!("{}", serde_json::to_string_pretty(&result)?);
} else {
println!("{}", result.audit_message);
println!(
"Lifecycle: {:?} | Running: {}",
result.lifecycle, result.running
);
}
}
ProjectCommand::Stop { project_id, json } => {
let result = project_registry::stop_project(&project_id)?;
if json {
println!("{}", serde_json::to_string_pretty(&result)?);
} else {
println!("{}", result.audit_message);
println!(
"Lifecycle: {:?} | Running: {}",
result.lifecycle, result.running
);
}
}
ProjectCommand::Restart { project_id, json } => {
let result = project_registry::restart_project(&project_id)?;
if json {
println!("{}", serde_json::to_string_pretty(&result)?);
} else {
println!("{}", result.audit_message);
println!(
"Lifecycle: {:?} | Running: {}",
result.lifecycle, result.running
);
}
}
ProjectCommand::Status { project_id, json } => {
let status = project_registry::get_project_status(&project_id)?;
if json {
println!("{}", serde_json::to_string_pretty(&status)?);
} else {
println!("Project: {}", status.project_id);
println!("Name: {}", status.name);
println!("Lifecycle: {:?}", status.lifecycle);
println!("Running: {}", status.running);
println!("Team: {}", status.team_name);
println!("Session: {}", status.session_name);
println!(
"Health: paused={} watchdog={} unhealthy={} triage_backlog={}",
status.health.paused,
status.health.watchdog_state,
status.health.unhealthy_members.len(),
status.health.triage_backlog_count
);
println!(
"Pipeline: active={} review={} runnable={} blocked={} stale_in_progress={} stale_review={}",
status.pipeline.active_task_count,
status.pipeline.review_queue_count,
status.pipeline.runnable_count,
status.pipeline.blocked_count,
status.pipeline.stale_in_progress_count,
status.pipeline.stale_review_count
);
}
}
ProjectCommand::SetActive {
project_id,
channel,
binding,
thread_binding,
json,
} => {
let scope = match (channel, binding, thread_binding) {
(None, None, None) => project_registry::ActiveProjectScope::Global,
(Some(channel), Some(binding), None) => {
project_registry::ActiveProjectScope::Channel { channel, binding }
}
(Some(channel), Some(binding), Some(thread_binding)) => {
project_registry::ActiveProjectScope::Thread {
channel,
binding,
thread_binding,
}
}
_ => bail!(
"set-active requires either no scope, --channel with --binding, or --channel with --binding and --thread-binding"
),
};
let selection = project_registry::set_active_project(&project_id, scope)?;
if json {
println!("{}", serde_json::to_string_pretty(&selection)?);
} else {
println!("Active project set to {}", selection.project_id);
}
}
ProjectCommand::Resolve {
message,
channel,
binding,
thread_binding,
json,
} => {
let decision = project_registry::resolve_project_for_message(
&project_registry::ProjectRoutingRequest {
message,
channel,
binding,
thread_binding,
},
)?;
if json {
println!("{}", serde_json::to_string_pretty(&decision)?);
} else {
println!(
"Selected: {}",
decision
.selected_project_id
.as_deref()
.unwrap_or("(clarification required)")
);
println!("Confidence: {:?}", decision.confidence);
println!("Requires confirmation: {}", decision.requires_confirmation);
println!("Reason: {}", decision.reason);
}
}
},
Command::Watchdog {
project_root,
resume,
} => {
team::run_watchdog(std::path::Path::new(&project_root), resume)?;
}
Command::Send {
from,
role,
message,
} => {
team::send_message_as(&root, from.as_deref(), &role, &message)?;
println!("Message queued for {role}.");
}
Command::Assign { engineer, task } => {
let id = team::assign_task(&root, &engineer, &task)?;
println!(
"Task queued for {engineer}. Inbox message id: {id}. Delivery result will be reported by Batty."
);
match team::wait_for_assignment_result(&root, &id, std::time::Duration::from_secs(8))? {
Some(result) => eprintln!("{}", team::format_assignment_result(&result)),
None => eprintln!(
"Assignment is still queued or pending delivery. No daemon result was available yet for {id}."
),
}
}
Command::Validate { show_checks } => {
team::validate_team(&root, show_checks)?;
}
Command::Config { json } => {
let config_path = team::team_config_path(&root);
if !config_path.exists() {
println!("No team config found. Run `batty init` first.");
return Ok(());
}
let team_config = team::config::TeamConfig::load(&config_path)?;
if json {
let members = team::hierarchy::resolve_hierarchy(&team_config)?;
let output = serde_json::json!({
"config_path": config_path.display().to_string(),
"team": team_config.name,
"roles": team_config.roles.len(),
"members": members.len(),
"board": {
"rotation_threshold": team_config.board.rotation_threshold,
"auto_dispatch": team_config.board.auto_dispatch,
},
"standup": {
"interval_secs": team_config.standup.interval_secs,
"output_lines": team_config.standup.output_lines,
},
"automation": {
"timeout_nudges": team_config.automation.timeout_nudges,
"standups": team_config.automation.standups,
"failure_pattern_detection": team_config.automation.failure_pattern_detection,
"triage_interventions": team_config.automation.triage_interventions,
"review_interventions": team_config.automation.review_interventions,
"owned_task_interventions": team_config.automation.owned_task_interventions,
"manager_dispatch_interventions": team_config.automation.manager_dispatch_interventions,
"architect_utilization_interventions": team_config.automation.architect_utilization_interventions,
},
"workflow": {
"mode": team_config.workflow_mode.as_str(),
"orchestrator_pane": team_config.orchestrator_pane,
},
});
println!("{}", serde_json::to_string_pretty(&output)?);
} else {
println!("Config: {}", config_path.display());
println!("Team: {}", team_config.name);
println!("Roles: {}", team_config.roles.len());
let members = team::hierarchy::resolve_hierarchy(&team_config)?;
println!("Total members: {}", members.len());
println!(
"Board rotation threshold: {}",
team_config.board.rotation_threshold
);
println!("Board auto-dispatch: {}", team_config.board.auto_dispatch);
println!("Standup interval: {}s", team_config.standup.interval_secs);
println!(
"Automation: timeout_nudges={}, standups={}, failure_patterns={}, triage={}, review={}, owned_tasks={}, manager_dispatch={}, architect_utilization={}",
team_config.automation.timeout_nudges,
team_config.automation.standups,
team_config.automation.failure_pattern_detection,
team_config.automation.triage_interventions,
team_config.automation.review_interventions,
team_config.automation.owned_task_interventions,
team_config.automation.manager_dispatch_interventions,
team_config.automation.architect_utilization_interventions,
);
println!(
"Workflow: mode={}, orchestrator_pane={}",
team_config.workflow_mode.as_str(),
team_config.orchestrator_pane
);
}
}
Command::Board { command } => {
let board_dir = root.join(".batty").join("team_config").join("board");
if !board_dir.is_dir() {
bail!(
"no board found at {}; run `batty init` first",
board_dir.display()
);
}
match command {
Some(BoardCommand::List { status }) => {
print!(
"{}",
team::board_cmd::list_tasks(&board_dir, status.as_deref())?
);
}
Some(BoardCommand::Summary) => {
for (status, count) in board_summary_counts(&board_dir)? {
println!("{status:<11} {count}");
}
}
Some(BoardCommand::Deps { format }) => {
let fmt = match format {
DepsFormatArg::Tree => team::deps::DepsFormat::Tree,
DepsFormatArg::Flat => team::deps::DepsFormat::Flat,
DepsFormatArg::Dot => team::deps::DepsFormat::Dot,
};
print!("{}", team::deps::render_deps(&board_dir, fmt)?);
}
Some(BoardCommand::Archive {
older_than,
dry_run,
}) => {
let max_age = team::board::parse_age_threshold(&older_than)?;
let tasks = team::board::done_tasks_older_than(&board_dir, max_age)?;
if dry_run {
println!("[dry-run] Would archive {} task(s):", tasks.len());
}
let summary = team::board::archive_tasks(&board_dir, &tasks, dry_run)?;
if !dry_run {
println!(
"Archived {} task(s) to {}",
summary.archived_count,
summary.archive_dir.display()
);
}
}
Some(BoardCommand::Health) => {
let events_path = root.join(".batty").join("team_config").join("events.jsonl");
let health = team::board_health::compute_health(&board_dir, &events_path)?;
print!("{}", team::board_health::format_health(&health));
}
None => {
let status = std::process::Command::new("kanban-md")
.args(["tui", "--dir", &board_dir.to_string_lossy()])
.status()
.context("failed to run kanban-md — is it installed?")?;
if !status.success() {
bail!("kanban-md tui failed");
}
}
}
}
Command::Inbox {
command,
member,
limit,
all,
raw,
} => match command {
Some(InboxCommand::Purge {
role,
all_roles,
before,
older_than,
all,
}) => {
let before = match (before, older_than) {
(Some(ts), _) => Some(ts),
(_, Some(dur)) => {
let age = team::board::parse_age_threshold(&dur)?;
let now = std::time::SystemTime::now()
.duration_since(std::time::UNIX_EPOCH)
.unwrap()
.as_secs();
Some(now.saturating_sub(age.as_secs()))
}
_ => None,
};
let summary = team::purge_inbox(&root, role.as_deref(), all_roles, before, all)?;
if all_roles {
println!(
"Purged {} delivered message(s) across {} inbox(es).",
summary.messages, summary.roles
);
} else {
let role = role.unwrap_or_default();
println!(
"Purged {} delivered message(s) from {role}.",
summary.messages
);
}
}
None => {
let member =
member.context("member is required unless using `batty inbox purge`")?;
let limit = if all { None } else { Some(limit) };
team::list_inbox(&root, &member, limit, raw)?;
}
},
Command::Read { member, id } => {
team::read_message(&root, &member, &id)?;
}
Command::Ack { member, id } => {
team::ack_message(&root, &member, &id)?;
println!("Message {id} acknowledged for {member}.");
}
Command::Review {
task_id,
disposition,
feedback,
reviewer,
} => {
let board_dir = team::team_config_dir(&root).join("board");
let disposition_str = match disposition {
cli::ReviewAction::Approve => "approve",
cli::ReviewAction::RequestChanges => "request-changes",
cli::ReviewAction::Reject => "reject",
};
team::task_cmd::cmd_review_structured(
&board_dir,
task_id,
disposition_str,
feedback.as_deref(),
&reviewer,
)?;
}
Command::Merge { engineer } => {
team::merge_worktree(&root, &engineer)?;
}
Command::Task { command } => {
let board_dir = team::team_config_dir(&root).join("board");
match command {
TaskCommand::Transition {
task_id,
target_state,
} => team::task_cmd::cmd_transition(
&board_dir,
task_id,
task_state_arg_name(target_state),
)?,
TaskCommand::Assign {
task_id,
execution_owner,
review_owner,
} => team::task_cmd::cmd_assign(
&board_dir,
task_id,
execution_owner.as_deref(),
review_owner.as_deref(),
)?,
TaskCommand::Review {
task_id,
disposition,
feedback,
} => team::task_cmd::cmd_review(
&board_dir,
task_id,
review_disposition_arg_name(disposition),
feedback.as_deref(),
)?,
TaskCommand::Update {
task_id,
branch,
commit,
blocked_on,
clear_blocked,
} => {
let mut fields = HashMap::new();
if let Some(branch) = branch {
fields.insert("branch".to_string(), branch);
}
if let Some(commit) = commit {
fields.insert("commit".to_string(), commit);
}
if let Some(blocked_on) = blocked_on {
fields.insert("blocked_on".to_string(), blocked_on);
}
if clear_blocked {
fields.insert("clear_blocked".to_string(), "true".to_string());
}
team::task_cmd::cmd_update(&board_dir, task_id, fields)?;
}
TaskCommand::AutoMerge { task_id, action } => {
let enabled = match action {
AutoMergeAction::Enable => true,
AutoMergeAction::Disable => false,
};
team::task_cmd::cmd_auto_merge(task_id, enabled, &root)?;
}
TaskCommand::Schedule {
task_id,
at,
cron,
clear,
} => team::task_cmd::cmd_schedule(
&board_dir,
task_id,
at.as_deref(),
cron.as_deref(),
clear,
)?,
}
}
Command::Activity { command } => {
let board_dir = team::team_config_dir(&root).join("board");
match command {
ActivityCommand::AnnotateStatus { task_id, source } => {
team::task_cmd::annotate_latest_status_activity_from_cli(
&board_dir, task_id, &source,
)?;
}
}
}
Command::Metrics => {
team::metrics_cmd::run(&root)?;
}
Command::StressTest {
compact,
duration_hours,
seed,
json_out,
markdown_out,
} => {
let report = team::stress::run(
&root,
team::stress::StressTestOptions {
compact,
duration_hours,
seed,
json_out,
markdown_out,
},
)?;
println!(
"Stress test complete: {} faults, {} failed SLA. JSON: {} Markdown: {}",
report.summary.total_faults,
report.summary.failed_faults,
report.json_report_path.display(),
report.markdown_report_path.display()
);
}
Command::Telemetry { command } => {
let conn =
team::telemetry_db::open(&root).context("failed to open telemetry database")?;
match command {
cli::TelemetryCommand::Summary => {
let rows = team::telemetry_db::query_session_summaries(&conn)?;
if rows.is_empty() {
println!("No session summaries recorded yet.");
} else {
println!(
"{:<24} {:<20} {:<20} {:>10} {:>8} {:>8}",
"SESSION", "STARTED", "ENDED", "COMPLETED", "MERGES", "EVENTS"
);
for row in &rows {
let started = format_ts(row.started_at);
let ended = row
.ended_at
.map(format_ts)
.unwrap_or_else(|| "running".to_string());
println!(
"{:<24} {:<20} {:<20} {:>10} {:>8} {:>8}",
row.session_id,
started,
ended,
row.tasks_completed,
row.total_merges,
row.total_events
);
}
}
}
cli::TelemetryCommand::Agents => {
let rows = team::telemetry_db::query_agent_metrics(&conn)?;
if rows.is_empty() {
println!("No agent metrics recorded yet.");
} else {
println!(
"{:<16} {:>11} {:>8} {:>8} {:>12} {:>8}",
"ROLE", "COMPLETIONS", "FAILURES", "RESTARTS", "CYCLE_SECS", "IDLE_PCT"
);
for row in &rows {
let total_polls = row.idle_polls + row.working_polls;
let idle_pct = if total_polls > 0 {
format!(
"{:.0}%",
row.idle_polls as f64 / total_polls as f64 * 100.0
)
} else {
"-".to_string()
};
println!(
"{:<16} {:>11} {:>8} {:>8} {:>12} {:>8}",
row.role,
row.completions,
row.failures,
row.restarts,
row.total_cycle_secs,
idle_pct
);
}
}
}
cli::TelemetryCommand::Tasks => {
let rows = team::telemetry_db::query_task_metrics(&conn)?;
if rows.is_empty() {
println!("No task metrics recorded yet.");
} else {
println!(
"{:<8} {:<20} {:<20} {:>7} {:>11} {:>9} {:>10} {:>10}",
"TASK",
"STARTED",
"COMPLETED",
"RETRIES",
"ESCALATIONS",
"CTX_RST",
"MERGE_SECS",
"CONFIDENCE"
);
for row in &rows {
let started = row
.started_at
.map(format_ts)
.unwrap_or_else(|| "-".to_string());
let completed = row
.completed_at
.map(format_ts)
.unwrap_or_else(|| "-".to_string());
let merge = row
.merge_time_secs
.map(|s| s.to_string())
.unwrap_or_else(|| "-".to_string());
let confidence = row
.confidence_score
.map(|c| format!("{:.2}", c))
.unwrap_or_else(|| "-".to_string());
println!(
"{:<8} {:<20} {:<20} {:>7} {:>11} {:>9} {:>10} {:>10}",
row.task_id,
started,
completed,
row.retries,
row.escalations,
row.context_restart_count,
merge,
confidence
);
}
}
}
cli::TelemetryCommand::Reviews => {
let row = team::telemetry_db::query_review_metrics(&conn)?;
let total_merges = row.auto_merge_count + row.manual_merge_count;
let auto_rate = if total_merges > 0 {
format!(
"{:.0}%",
row.auto_merge_count as f64 / total_merges as f64 * 100.0
)
} else {
"-".to_string()
};
let total_reviewed = total_merges + row.rework_count;
let rework_rate = if total_reviewed > 0 {
format!(
"{:.0}%",
row.rework_count as f64 / total_reviewed as f64 * 100.0
)
} else {
"-".to_string()
};
let avg_latency = row
.avg_review_latency_secs
.map(|s| format!("{:.0}s", s))
.unwrap_or_else(|| "-".to_string());
println!("Review Pipeline (all sessions)");
println!(
"Auto-merge Rate: {} | Rework Rate: {}",
auto_rate, rework_rate
);
println!(
"Auto: {} | Manual: {} | Rework: {} | Nudges: {} | Escalations: {}",
row.auto_merge_count,
row.manual_merge_count,
row.rework_count,
row.review_nudge_count,
row.review_escalation_count
);
println!("Avg Review Latency: {}", avg_latency);
}
cli::TelemetryCommand::Events { limit } => {
let rows = team::telemetry_db::query_recent_events(&conn, limit)?;
if rows.is_empty() {
println!("No events recorded yet.");
} else {
println!(
"{:<20} {:<24} {:<16} {:<8}",
"TIMESTAMP", "EVENT", "ROLE", "TASK"
);
for row in &rows {
let ts = format_ts(row.timestamp);
let role = row.role.as_deref().unwrap_or("-");
let task = row.task_id.as_deref().unwrap_or("-");
println!("{:<20} {:<24} {:<16} {:<8}", ts, row.event_type, role, task);
}
}
}
}
}
Command::Chat {
agent_type,
cmd,
cwd,
sdk_mode,
} => {
use batty_cli::shim;
let at: shim::classifier::AgentType =
agent_type.parse().map_err(|e: String| anyhow::anyhow!(e))?;
let cmd = if sdk_mode && cmd.is_none() {
match at {
shim::classifier::AgentType::Kiro => {
use batty_cli::agent::kiro::KiroCliAdapter;
KiroCliAdapter::new(None).sdk_launch_command()
}
_ => {
use batty_cli::agent::claude::ClaudeCodeAdapter;
ClaudeCodeAdapter::new(None).sdk_launch_command(None, None)
}
}
} else {
cmd.unwrap_or_else(|| shim::chat::default_cmd(at).to_string())
};
let cwd = PathBuf::from(&cwd)
.canonicalize()
.unwrap_or_else(|_| PathBuf::from(&cwd));
shim::chat::run(at, &cmd, &cwd, sdk_mode)?;
}
Command::Shim {
id,
agent_type,
cmd,
cwd,
rows,
cols,
pty_log_path,
graceful_shutdown_timeout_secs,
auto_commit_on_restart,
sdk_mode,
} => {
use batty_cli::shim;
use std::os::unix::io::FromRawFd;
use std::os::unix::net::UnixStream;
let at: shim::classifier::AgentType =
agent_type.parse().map_err(|e: String| anyhow::anyhow!(e))?;
let stream = unsafe { UnixStream::from_raw_fd(3) };
let channel = shim::protocol::Channel::new(stream);
let args = shim::runtime::ShimArgs {
id,
agent_type: at,
cmd,
cwd: PathBuf::from(cwd),
rows,
cols,
pty_log_path: pty_log_path.map(PathBuf::from),
graceful_shutdown_timeout_secs,
auto_commit_on_restart,
};
if sdk_mode {
match at {
shim::classifier::AgentType::Codex => {
shim::runtime_codex::run_codex_sdk(args, channel)?;
}
shim::classifier::AgentType::Kiro => {
shim::runtime_kiro::run_kiro_acp(args, channel)?;
}
_ => {
shim::runtime_sdk::run_sdk(args, channel)?;
}
}
} else {
shim::runtime::run(args, channel)?;
}
}
Command::ConsolePane {
project_root,
member,
events_log_path,
pty_log_path,
} => {
batty_cli::console_pane::run(
PathBuf::from(project_root).as_path(),
&member,
PathBuf::from(events_log_path).as_path(),
PathBuf::from(pty_log_path).as_path(),
)?;
}
Command::DaemonRestartIfStale { dry_run } => {
let report = team::restart_daemon_if_stale(&root, dry_run)?;
println!("{}", report.render());
}
Command::Daemon {
project_root,
resume,
} => {
let root = std::path::PathBuf::from(project_root);
team::run_daemon(&root, resume)?;
}
Command::Completions { shell } => {
use clap::CommandFactory;
let shell = match shell {
cli::CompletionShell::Bash => clap_complete::Shell::Bash,
cli::CompletionShell::Zsh => clap_complete::Shell::Zsh,
cli::CompletionShell::Fish => clap_complete::Shell::Fish,
};
clap_complete::generate(shell, &mut Cli::command(), "batty", &mut std::io::stdout());
}
Command::Nudge { command } => match command {
NudgeCommand::Disable { name } => {
team::disable_nudge(&root, name.marker_name())?;
println!("Intervention '{}' disabled.", name.marker_name());
}
NudgeCommand::Enable { name } => {
team::enable_nudge(&root, name.marker_name())?;
println!("Intervention '{}' re-enabled.", name.marker_name());
}
NudgeCommand::Status => {
team::nudge_status(&root)?;
}
},
Command::Pause => {
team::pause_team(&root)?;
println!("Nudges and standups paused. Run `batty resume` to resume.");
}
Command::Resume => {
team::resume_team(&root)?;
println!("Nudges and standups resumed.");
}
Command::Load => {
team::show_load(&root)?;
}
Command::Parity { detail, gaps } => {
team::parity::show_parity(&root, detail, gaps)?;
}
Command::Verify => {
team::verification::cmd_verify(&root)?;
}
Command::Release { tag, readiness } => {
if readiness {
release::cmd_release_readiness(&root, tag.as_deref())?;
} else {
release::cmd_release(&root, tag.as_deref())?;
}
}
Command::Queue => {
let entries = team::daemon::load_dispatch_queue_snapshot(&root);
if entries.is_empty() {
println!("Dispatch queue is empty.");
} else {
println!(
"{:<20} {:<8} {:<36} {:>8} LAST FAILURE",
"ENGINEER", "TASK", "TITLE", "FAILURES"
);
println!("{}", "-".repeat(100));
for entry in entries {
println!(
"{:<20} {:<8} {:<36} {:>8} {}",
entry.engineer,
entry.task_id,
entry.task_title,
entry.validation_failures,
entry.last_failure.unwrap_or_else(|| "-".to_string())
);
}
}
}
Command::Dispatch { explain, task } => {
if !explain {
println!("Use `batty dispatch --explain` to inspect the current routing decision.");
} else {
team::allocation::print_dispatch_explanation(&root, task)?;
}
}
Command::Cost => {
team::cost::show_cost(&root)?;
}
Command::Scale { command } => {
team::scale::run(&root, command)?;
}
Command::Reload => {
team::reload::request_topology_reload(&root)?;
println!("Topology reload requested. The daemon will pick it up on the next poll.");
}
Command::Research { command } => match command {
ResearchCommand::Start {
hypothesis,
evaluator,
format,
keep_policy,
max_iterations,
worktree,
} => {
let worktree_dir = if worktree.is_absolute() {
worktree
} else {
std::env::current_dir()?.join(worktree)
};
let mission = team::autoresearch::start_research(
&root,
team::autoresearch::StartResearchOptions {
hypothesis,
evaluator_command: evaluator,
evaluator_format: match format {
ResearchFormatArg::Json => team::autoresearch::EvaluatorFormat::Json,
ResearchFormatArg::ExitCode => {
team::autoresearch::EvaluatorFormat::ExitCode
}
},
keep_policy: match keep_policy {
ResearchKeepPolicyArg::PassOnly => {
team::autoresearch::KeepPolicy::PassOnly
}
ResearchKeepPolicyArg::ScoreImprovement => {
team::autoresearch::KeepPolicy::ScoreImprovement
}
ResearchKeepPolicyArg::ParityImprovement => {
team::autoresearch::KeepPolicy::ParityImprovement
}
},
max_iterations,
worktree_dir,
},
)?;
println!("Started research mission {}", mission.id);
println!("Hypothesis: {}", mission.hypothesis);
println!("Worktree: {}", mission.worktree_dir.display());
}
ResearchCommand::Status => {
team::autoresearch::print_status(&root)?;
}
ResearchCommand::Ledger => {
team::autoresearch::print_ledger(&root)?;
}
ResearchCommand::Stop => match team::autoresearch::stop_current_research(&root)? {
Some(mission) => println!("Stopped research mission {}", mission.id),
None => println!("No active research mission."),
},
},
Command::Doctor { fix, yes } => {
print!("{}", team::doctor::run(&root, fix, yes)?);
}
Command::Worktree { health } => {
if !health {
bail!("worktree command requires --health");
}
print!("{}", team::worktree_health::run(&root)?);
}
Command::Grafana { command } => {
let config_path = team::team_config_path(&root);
let port = if config_path.exists() {
team::config::TeamConfig::load(&config_path)?.grafana.port
} else {
team::grafana::DEFAULT_PORT
};
match command {
GrafanaCommand::Setup => team::grafana::setup(&root, port)?,
GrafanaCommand::Status => team::grafana::status(port)?,
GrafanaCommand::Open => team::grafana::open(port)?,
}
}
Command::GrafanaWebhook { project_root, port } => {
team::grafana::run_alert_webhook(std::path::Path::new(&project_root), port)?;
}
Command::Discord { command } => match command.unwrap_or(DiscordCommand::Setup) {
DiscordCommand::Setup => team::setup_discord(&root)?,
DiscordCommand::Status => team::discord_status(&root)?,
},
Command::Telegram => {
team::setup_telegram(&root)?;
}
}
Ok(())
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn board_list_counts_data_rows_only() {
let output = "\
ID STATUS PRIORITY TITLE CLAIMED TAGS DUE
1 todo high One -- -- --
12 review medium Two @eng-1 dx --
";
assert_eq!(count_board_list_rows(output), 2);
}
#[test]
fn board_list_ignores_headers_and_empty_output() {
let output = "\
ID STATUS PRIORITY TITLE CLAIMED TAGS DUE
";
assert_eq!(count_board_list_rows(output), 0);
assert_eq!(count_board_list_rows(""), 0);
}
}