use std::collections::{BTreeSet, HashSet};
use std::path::{Path, PathBuf};
use std::time::Duration;
use anyhow::Result;
use team_core::compose::Compose;
use team_core::supervisor::{AgentSpec, AgentState, DrainOutcome, Supervisor, TmuxSupervisor};
use super::agent_filter::AgentSelector;
use super::snapshot::{self, AgentEntry, ReloadPlan, RemovedAgent};
pub fn run(
root: &Path,
dry_run: bool,
project: Option<&str>,
sel: &AgentSelector,
fresh: bool,
force: bool,
) -> Result<()> {
let compose = super::load(root)?;
let errs = team_core::validate::validate(&compose);
if !errs.is_empty() {
for e in &errs {
eprintln!("error: {e}");
}
anyhow::bail!("{} validation error(s) — fix before reload", errs.len());
}
let scoped = project
.map(|name| super::project_filter::resolve(&compose, name))
.transpose()?;
super::up::guard_no_name_collision(&compose, scoped.as_deref())?;
if sel.is_scoped() {
if let Some(id) = scoped.as_deref() {
let targets = super::agent_filter::resolve(&compose, id, sel)?
.expect("scoped selector resolves to a concrete agent set");
return force_restart_scoped(&compose, id, &targets, dry_run, fresh);
}
}
let prev = snapshot::read(&compose.root);
let bin = super::team_mcp_bin().display().to_string();
let next = snapshot::compute(&compose, &bin);
let mut plan = snapshot::plan(prev.as_ref(), &next);
plan.remove
.extend(orphan_removals(prev.is_some(), registry_orphans(&compose)));
if let Some(id) = scoped.as_deref() {
plan = filter_plan_to_project(plan, id);
}
if force {
force_promote_keeps(&mut plan, prev.as_ref());
}
let no_changes = plan.is_empty()
&& prev
.as_ref()
.map(|s| s.compose_digest == next.compose_digest && s.global == next.global)
.unwrap_or(false);
if no_changes {
if dry_run {
println!("no changes (dry run)");
} else {
println!("no changes");
}
return Ok(());
}
if dry_run {
print_plan(&plan, true, fresh);
return Ok(());
}
if let Some(id) = scoped.as_deref() {
super::up::render_project_public(&compose, id)?;
} else {
super::up::ensure_wrapper_and_dirs(&compose)?;
super::up::render_all_public(&compose)?;
super::up::register_all_public(&compose)?;
}
apply_plan(&compose, &plan, fresh)?;
let snap = match scoped.as_deref() {
Some(id) => snapshot::merge_project_into(prev.as_ref(), &next, id),
None => next,
};
snapshot::write(&compose.root, &snap)?;
super::up::record_in_registry(&compose, scoped.as_deref());
Ok(())
}
fn force_restart_scoped(
compose: &Compose,
project_id: &str,
targets: &BTreeSet<String>,
dry_run: bool,
fresh: bool,
) -> Result<()> {
let ids: Vec<String> = compose
.agents()
.filter(|h| h.project == project_id && targets.contains(h.agent))
.map(|h| h.id())
.collect();
if ids.is_empty() {
println!("no agents in scope for project {project_id}.");
return Ok(());
}
if dry_run {
for id in &ids {
println!(
"reloaded · {id} (forced){} (dry run)",
super::up::fresh_suffix(fresh)
);
}
return Ok(());
}
super::up::render_project_public(compose, project_id)?;
let prev = snapshot::read(&compose.root);
let sup = TmuxSupervisor;
let drain_timeout = Duration::from_secs(compose.global.supervisor.drain_timeout_secs);
for id in &ids {
let drain_spec = match prev.as_ref().and_then(|s| s.agents.get(id)) {
Some(e) => spec_from_prior(compose, id, e),
None => match compose.agents().find(|h| &h.id() == id) {
Some(h) => {
AgentSpec::from_handle(h, &compose.root, &compose.global.supervisor.tmux_prefix)
}
None => continue,
},
};
let outcome = sup.drain(&drain_spec, drain_timeout)?;
if let Some(h) = compose.agents().find(|h| &h.id() == id) {
let spec =
AgentSpec::from_handle(h, &compose.root, &compose.global.supervisor.tmux_prefix);
super::up::freshen_for_spec(&compose.root, &spec, &h.spec.runtime, fresh);
sup.up(&spec)?;
}
println!(
"reloaded · {id} (forced){}{}",
super::up::fresh_suffix(fresh),
drain_suffix(outcome)
);
}
let bin = super::team_mcp_bin().display().to_string();
let next = snapshot::compute(compose, &bin);
let snap = snapshot::merge_project_into(prev.as_ref(), &next, project_id);
snapshot::write(&compose.root, &snap)?;
Ok(())
}
fn filter_plan_to_project(plan: ReloadPlan, project_id: &str) -> ReloadPlan {
let prefix = format!("{project_id}:");
let in_project = |id: &str| id.starts_with(&prefix);
ReloadPlan {
add: plan.add.into_iter().filter(|id| in_project(id)).collect(),
change: plan
.change
.into_iter()
.filter(|(id, _)| in_project(id))
.collect(),
remove: plan
.remove
.into_iter()
.filter(|r| in_project(&r.id))
.collect(),
keep: plan.keep.into_iter().filter(|id| in_project(id)).collect(),
change_prior: plan
.change_prior
.into_iter()
.filter(|(id, _)| in_project(id))
.collect(),
}
}
fn force_promote_keeps(plan: &mut ReloadPlan, prev: Option<&snapshot::Snapshot>) {
let keeps = std::mem::take(&mut plan.keep);
for id in keeps {
match prev.and_then(|s| s.agents.get(&id)) {
Some(prior) => {
plan.change_prior.insert(id.clone(), prior.clone());
plan.change.push((id, snapshot::ChangedInputs::forced()));
}
None => plan.keep.push(id),
}
}
}
fn change_line_head(id: &str, inputs: &snapshot::ChangedInputs) -> String {
if inputs.any() {
format!("changed · {id} ({})", inputs.label())
} else {
format!("reloaded · {id} (forced)")
}
}
fn print_plan(plan: &ReloadPlan, dry: bool, fresh: bool) {
let dry_suffix = if dry { " (dry run)" } else { "" };
let fresh_suffix = super::up::fresh_suffix(fresh);
for r in &plan.remove {
println!("removed · {}{dry_suffix}", r.id);
}
for (id, inputs) in &plan.change {
println!("{}{fresh_suffix}{dry_suffix}", change_line_head(id, inputs));
}
for id in &plan.add {
println!("added · {id}{fresh_suffix}{dry_suffix}");
}
}
fn apply_plan(compose: &Compose, plan: &ReloadPlan, fresh: bool) -> Result<()> {
let sup = TmuxSupervisor;
let drain_timeout = Duration::from_secs(compose.global.supervisor.drain_timeout_secs);
for r in &plan.remove {
let outcome = sup.drain(&spec_from_removed(compose, r), drain_timeout)?;
println!("removed · {}{}", r.id, drain_suffix(outcome));
}
for (id, inputs) in &plan.change {
let prior = plan
.change_prior
.get(id)
.expect("change_prior populated alongside every change entry");
let outcome = sup.drain(&spec_from_prior(compose, id, prior), drain_timeout)?;
if let Some(h) = compose.agents().find(|h| &h.id() == id) {
let spec =
AgentSpec::from_handle(h, &compose.root, &compose.global.supervisor.tmux_prefix);
super::up::freshen_for_spec(&compose.root, &spec, &h.spec.runtime, fresh);
sup.up(&spec)?;
}
println!(
"{}{}{}",
change_line_head(id, inputs),
super::up::fresh_suffix(fresh),
drain_suffix(outcome)
);
}
for id in &plan.add {
if let Some(h) = compose.agents().find(|h| &h.id() == id) {
let spec =
AgentSpec::from_handle(h, &compose.root, &compose.global.supervisor.tmux_prefix);
super::up::freshen_for_spec(&compose.root, &spec, &h.spec.runtime, fresh);
sup.up(&spec)?;
println!("added · {id}{}", super::up::fresh_suffix(fresh));
}
}
for id in &plan.keep {
if let Some(h) = compose.agents().find(|h| &h.id() == id) {
let spec =
AgentSpec::from_handle(h, &compose.root, &compose.global.supervisor.tmux_prefix);
if sup.state(&spec)? == AgentState::Stopped {
super::up::freshen_for_spec(&compose.root, &spec, &h.spec.runtime, fresh);
sup.up(&spec)?;
println!("started · {id}{}", super::up::fresh_suffix(fresh));
}
}
}
Ok(())
}
fn drain_suffix(outcome: DrainOutcome) -> &'static str {
match outcome {
DrainOutcome::Graceful => "",
DrainOutcome::TimedOutKilled => " [drain timed out — killed]",
}
}
fn spec_from_removed(compose: &Compose, r: &RemovedAgent) -> AgentSpec {
let (project, agent) = r.id.split_once(':').unwrap_or((r.id.as_str(), ""));
AgentSpec {
project: project.into(),
agent: agent.into(),
tmux_session: r.tmux_session.clone(),
wrapper: super::agent_wrapper(&compose.root),
cwd: compose.root.clone(),
env_file: r.env_file.clone(),
}
}
fn spec_from_prior(compose: &Compose, id: &str, prior: &AgentEntry) -> AgentSpec {
let (project, agent) = id.split_once(':').unwrap_or((id, ""));
AgentSpec {
project: project.into(),
agent: agent.into(),
tmux_session: prior.tmux_session.clone(),
wrapper: super::agent_wrapper(&compose.root),
cwd: compose.root.clone(),
env_file: PathBuf::from(&prior.env_file),
}
}
fn orphan_removals(prev_present: bool, orphans: Vec<RemovedAgent>) -> Vec<RemovedAgent> {
if prev_present {
Vec::new()
} else {
orphans
}
}
fn registry_orphans(compose: &Compose) -> Vec<RemovedAgent> {
let Some(dir) = team_core::registry::config_dir() else {
return Vec::new();
};
let desired: HashSet<String> = compose.agents().map(|h| h.id()).collect();
match team_core::registry::orphans_for_root(&dir, &compose.root, &desired, None, false) {
Ok(orphans) => orphans
.into_iter()
.map(|o| RemovedAgent {
id: o.id(),
tmux_session: o.tmux_session,
env_file: PathBuf::new(),
})
.collect(),
Err(e) => {
eprintln!("warn · teams registry: {e:#}");
Vec::new()
}
}
}
#[cfg(test)]
mod tests {
use super::*;
use crate::cmd::snapshot::{
AgentEntry, ChangedInputs, Fingerprints, PromptFingerprint, RemovedAgent,
};
use std::collections::BTreeMap;
#[test]
fn drain_suffix_empty_on_graceful() {
assert_eq!(drain_suffix(DrainOutcome::Graceful), "");
}
#[test]
fn drain_suffix_annotates_timeout() {
assert!(drain_suffix(DrainOutcome::TimedOutKilled).contains("drain timed out"));
}
#[test]
fn fresh_suffix_annotates_only_when_fresh() {
assert_eq!(super::super::up::fresh_suffix(true), " (fresh)");
assert_eq!(super::super::up::fresh_suffix(false), "");
}
fn entry(env: &str) -> AgentEntry {
AgentEntry {
tmux_session: "a-x".into(),
env_file: env.into(),
fingerprints: Fingerprints {
env: String::new(),
mcp: String::new(),
role_prompt: PromptFingerprint::None,
},
}
}
fn removed(id: &str) -> RemovedAgent {
RemovedAgent {
id: id.into(),
tmux_session: format!("a-{id}"),
env_file: PathBuf::from(""),
}
}
fn changed_inputs() -> ChangedInputs {
ChangedInputs {
env: true,
mcp: false,
role_prompt: false,
}
}
#[test]
fn filter_plan_keeps_only_matching_project_entries() {
let mut change_prior = BTreeMap::new();
change_prior.insert("a:m".into(), entry("/tmp/a-m.env"));
change_prior.insert("b:m".into(), entry("/tmp/b-m.env"));
let plan = ReloadPlan {
add: vec!["a:w".into(), "b:w".into()],
change: vec![
("a:m".into(), changed_inputs()),
("b:m".into(), changed_inputs()),
],
remove: vec![removed("a:gone"), removed("b:gone")],
keep: vec!["a:keep".into(), "b:keep".into()],
change_prior,
};
let filtered = filter_plan_to_project(plan, "a");
assert_eq!(filtered.add, vec!["a:w"]);
assert_eq!(filtered.change.len(), 1);
assert_eq!(filtered.change[0].0, "a:m");
assert_eq!(filtered.remove.len(), 1);
assert_eq!(filtered.remove[0].id, "a:gone");
assert_eq!(filtered.keep, vec!["a:keep"]);
assert_eq!(filtered.change_prior.len(), 1);
assert!(filtered.change_prior.contains_key("a:m"));
}
#[test]
fn filter_plan_does_not_match_prefix_collisions() {
let plan = ReloadPlan {
add: vec!["a:m".into(), "aa:m".into(), "ab:m".into()],
..ReloadPlan::default()
};
let filtered = filter_plan_to_project(plan, "a");
assert_eq!(filtered.add, vec!["a:m"]);
}
#[test]
fn filter_plan_returns_empty_when_no_entries_match() {
let plan = ReloadPlan {
add: vec!["a:m".into(), "b:m".into()],
..ReloadPlan::default()
};
let filtered = filter_plan_to_project(plan, "z");
assert!(filtered.is_empty());
}
fn snapshot_with(ids: &[&str]) -> snapshot::Snapshot {
let mut agents = BTreeMap::new();
for id in ids {
agents.insert((*id).to_string(), entry("/tmp/x.env"));
}
snapshot::Snapshot {
agents,
..Default::default()
}
}
#[test]
fn force_promote_keeps_moves_keeps_into_change_as_forced() {
let prev = snapshot_with(&["a:m", "a:w"]);
let mut plan = ReloadPlan {
keep: vec!["a:m".into(), "a:w".into()],
..ReloadPlan::default()
};
force_promote_keeps(&mut plan, Some(&prev));
assert!(plan.keep.is_empty(), "keeps drained into change");
assert_eq!(plan.change.len(), 2);
for (id, inputs) in &plan.change {
assert!(!inputs.any(), "{id} promoted as forced (all-false)");
assert!(
plan.change_prior.contains_key(id),
"prior carried for {id} so teardown hits the running session"
);
}
assert_eq!(
change_line_head(&plan.change[0].0, &plan.change[0].1),
"reloaded · a:m (forced)"
);
}
#[test]
fn force_promote_keeps_without_prior_entry_leaves_it_kept() {
let prev = snapshot_with(&["a:m"]);
let mut plan = ReloadPlan {
keep: vec!["a:ghost".into()],
..ReloadPlan::default()
};
force_promote_keeps(&mut plan, Some(&prev));
assert_eq!(plan.keep, vec!["a:ghost"]);
assert!(plan.change.is_empty());
}
#[test]
fn change_line_head_distinguishes_diff_from_forced() {
let genuine = ChangedInputs {
env: true,
mcp: true,
role_prompt: false,
};
assert_eq!(change_line_head("a:m", &genuine), "changed · a:m (env+mcp)");
assert_eq!(
change_line_head("a:m", &ChangedInputs::forced()),
"reloaded · a:m (forced)"
);
}
#[test]
fn orphan_removals_only_fires_without_a_prior_snapshot() {
assert!(orphan_removals(true, vec![removed("a:gone")]).is_empty());
let passed = orphan_removals(false, vec![removed("a:gone"), removed("a:zap")]);
assert_eq!(
passed.iter().map(|r| r.id.as_str()).collect::<Vec<_>>(),
vec!["a:gone", "a:zap"]
);
}
}