304 lines
11 KiB
Rust
304 lines
11 KiB
Rust
//! start, stop, restart and run: process management commands.
|
|
|
|
use super::*;
|
|
use crate::cli::StartArgs;
|
|
use crate::events::{self, Event, EventKind};
|
|
use crate::process;
|
|
use crate::sessions;
|
|
use anyhow::{bail, Context, Result};
|
|
use std::collections::BTreeMap;
|
|
use std::ffi::OsString;
|
|
use std::path::PathBuf;
|
|
use std::process::Command;
|
|
|
|
/// Resolve the start target: the explicit agent/alias/group when given,
|
|
/// otherwise the default agent of the project containing cwd (issue #30).
|
|
pub fn resolve_start_target(app: &App, cwd: &std::path::Path, agent: Option<&str>) -> Result<String> {
|
|
if let Some(a) = agent.filter(|s| !s.is_empty()) {
|
|
return Ok(a.to_string());
|
|
}
|
|
crate::projects::default_agent_for_dir(app, cwd).ok_or_else(|| {
|
|
anyhow!(
|
|
"no agent given — set default_agent in this project's profile ('am init' does it) or pass an agent name"
|
|
)
|
|
})
|
|
}
|
|
|
|
pub fn start(app: &App, opts: &StartArgs) -> Result<i32> {
|
|
start_with_cwd(app, opts, None)
|
|
}
|
|
|
|
/// Start with an explicit working directory (used by sessions resume).
|
|
pub fn start_with_cwd(app: &App, opts: &StartArgs, cwd: Option<&std::path::Path>) -> Result<i32> {
|
|
let current = std::env::current_dir().unwrap_or_default();
|
|
let target = resolve_start_target(app, ¤t, opts.agent.as_deref())?;
|
|
let extra_env = parse_env_list(&opts.env)?;
|
|
let extra_args = parse_extra_args(&opts.args);
|
|
if let Some(group) = crate::catalog::Catalog::parse_group_selector(&target) {
|
|
let members = app.catalog.group_members(group);
|
|
if members.is_empty() {
|
|
bail!("unknown group '{group}'");
|
|
}
|
|
for m in members {
|
|
start_one(app, m, &extra_args, &extra_env, opts.notify, true, cwd)?;
|
|
}
|
|
return Ok(0);
|
|
}
|
|
let agent = require_agent(app, &target)?;
|
|
start_one(app, agent, &extra_args, &extra_env, opts.notify, opts.background, cwd)
|
|
}
|
|
|
|
fn start_one(
|
|
app: &App,
|
|
agent: &AgentDef,
|
|
extra_args: &[String],
|
|
extra_env: &BTreeMap<String, String>,
|
|
notify: bool,
|
|
background: bool,
|
|
cwd: Option<&std::path::Path>,
|
|
) -> Result<i32> {
|
|
let exec = resolve_exec(app, agent, extra_args, extra_env)?;
|
|
let bin = PathBuf::from(&exec.program);
|
|
let is_managed = app.state.get(&agent.name).ok().flatten().is_some();
|
|
let ctx = crate::context::detect(&std::env::current_dir().unwrap_or_default());
|
|
let current = std::env::current_dir().unwrap_or_default();
|
|
crate::hooks::run_hooks(app, "on_start", cwd.unwrap_or(¤t));
|
|
if background {
|
|
let pid = process::spawn_background(app, &agent.name, &bin, &exec.args, &exec.env, cwd)?;
|
|
if app.dry_run() {
|
|
return Ok(0);
|
|
}
|
|
if is_managed {
|
|
app.state
|
|
.update_pid(&agent.name, Some(pid), Some(crate::installers::now_rfc3339()))?;
|
|
} else {
|
|
app.log.warn(&format!(
|
|
"agent '{}' is external (not installed by agent-manager); PID {pid} is not recorded",
|
|
agent.name
|
|
));
|
|
}
|
|
app.emit(
|
|
&Event::now(EventKind::Start)
|
|
.with_agent(agent.name.clone())
|
|
.with_pid(pid)
|
|
.with_session(crate::repl::session_id())
|
|
.with_cwd(cwd_string())
|
|
.with_project(ctx.root.clone())
|
|
.with_args(exec.args.clone())
|
|
.with_env_keys(events::env_keys(&exec.env)),
|
|
);
|
|
if is_managed {
|
|
let _ = app.state.touch(&agent.name);
|
|
}
|
|
let log = process::agent_log_path(app, &agent.name);
|
|
let _ = sessions::start_agent(
|
|
app,
|
|
&agent.name,
|
|
pid,
|
|
cwd_string(),
|
|
ctx.root.clone(),
|
|
exec.args.clone(),
|
|
events::env_keys(&exec.env),
|
|
Some(log.display().to_string()),
|
|
);
|
|
app.log.success(&format!(
|
|
"{} started (pid {pid}); log: {}",
|
|
agent.title(),
|
|
log.display()
|
|
));
|
|
if notify {
|
|
process::desktop_notify("agent-manager", &format!("{} started", agent.title()));
|
|
}
|
|
Ok(0)
|
|
} else {
|
|
if app.dry_run() {
|
|
app.log.dry(format!(
|
|
"would run {} {} (foreground)",
|
|
exec.program,
|
|
exec.args.join(" ")
|
|
));
|
|
return Ok(0);
|
|
}
|
|
let started = std::time::Instant::now();
|
|
let (prog, prefix) = crate::runner::resolve_program(&exec.program);
|
|
let mut full_args = prefix;
|
|
full_args.extend(exec.args.iter().cloned());
|
|
let mut cmd = Command::new(&prog);
|
|
cmd.args(&full_args).envs(&exec.env);
|
|
if let Some(dir) = cwd {
|
|
cmd.current_dir(dir);
|
|
}
|
|
let status = cmd
|
|
.status()
|
|
.with_context(|| format!("failed to run {}", exec.program))?;
|
|
let code = status.code().unwrap_or(1);
|
|
app.emit(
|
|
&Event::now(EventKind::Run)
|
|
.with_agent(agent.name.clone())
|
|
.with_cwd(cwd_string())
|
|
.with_project(ctx.root.clone())
|
|
.with_args(exec.args.clone())
|
|
.with_env_keys(events::env_keys(&exec.env))
|
|
.with_exit_code(code)
|
|
.with_duration(started.elapsed().as_secs()),
|
|
);
|
|
let _ = app.state.touch(&agent.name);
|
|
Ok(code)
|
|
}
|
|
}
|
|
|
|
pub fn stop(app: &App, target: &str, force: bool, timeout: Option<u64>) -> Result<i32> {
|
|
let timeout = timeout.unwrap_or_else(|| app.config.settings.stop_timeout());
|
|
if let Some(group) = crate::catalog::Catalog::parse_group_selector(target) {
|
|
let members = app.catalog.group_members(group);
|
|
if members.is_empty() {
|
|
bail!("unknown group '{group}'");
|
|
}
|
|
for m in members.iter().rev() {
|
|
stop_one(app, m, force, timeout)?;
|
|
}
|
|
return Ok(0);
|
|
}
|
|
let agent = require_agent(app, target)?;
|
|
stop_one(app, agent, force, timeout)
|
|
}
|
|
|
|
fn stop_one(app: &App, agent: &AgentDef, force: bool, timeout: u64) -> Result<i32> {
|
|
let Some(entry) = app.state.get(&agent.name).ok().flatten() else {
|
|
app.log.info(&format!(
|
|
"agent '{}' is not managed by agent-manager; nothing to stop",
|
|
agent.name
|
|
));
|
|
return Ok(0);
|
|
};
|
|
let Some(pid) = entry.pid else {
|
|
app.log.info(&format!("agent '{}' is not running", agent.name));
|
|
return Ok(0);
|
|
};
|
|
if app.dry_run() {
|
|
app.log.dry(format!(
|
|
"would stop pid {pid} ({}{})",
|
|
agent.name,
|
|
if force { "force" } else { "graceful" }
|
|
));
|
|
return Ok(0);
|
|
}
|
|
if !process::is_running(pid) {
|
|
app.log.warn(&format!(
|
|
"stale PID {pid} for '{}' (process no longer running); cleaning up",
|
|
agent.name
|
|
));
|
|
app.emit(
|
|
&Event::now(EventKind::Stop)
|
|
.with_agent(agent.name.clone())
|
|
.with_pid(pid)
|
|
.with_reason("stale"),
|
|
);
|
|
let _ = sessions::finish_agent(app, &agent.name, Some(pid), 0, true);
|
|
app.state.update_pid(&agent.name, None, None)?;
|
|
let current = std::env::current_dir().unwrap_or_default();
|
|
crate::hooks::run_hooks(app, "on_stop", ¤t);
|
|
return Ok(0);
|
|
}
|
|
app.log.info(&format!("stopping '{}' (pid {pid})", agent.name));
|
|
let stopped = process::stop_pid(pid, force, timeout);
|
|
let duration = seconds_since(entry.started_at.as_deref());
|
|
let mut ev = Event::now(EventKind::Stop)
|
|
.with_agent(agent.name.clone())
|
|
.with_pid(pid)
|
|
.with_reason(if force { "force" } else { "signal" })
|
|
.with_exit_code(if stopped { 0 } else { 1 });
|
|
if let Some(d) = duration {
|
|
ev = ev.with_duration(d);
|
|
}
|
|
app.emit(&ev);
|
|
let _ = sessions::finish_agent(app, &agent.name, Some(pid), if stopped { 0 } else { 1 }, stopped);
|
|
app.state.update_pid(&agent.name, None, None)?;
|
|
let current = std::env::current_dir().unwrap_or_default();
|
|
crate::hooks::run_hooks(app, "on_stop", ¤t);
|
|
if stopped {
|
|
app.log.success(&format!("{} stopped", agent.title()));
|
|
Ok(0)
|
|
} else {
|
|
app.log.error(&format!("failed to stop pid {pid}"));
|
|
Ok(1)
|
|
}
|
|
}
|
|
|
|
pub fn restart(app: &App, opts: &StartArgs, force: bool, timeout: Option<u64>) -> Result<i32> {
|
|
let timeout = timeout.unwrap_or_else(|| app.config.settings.stop_timeout());
|
|
let cwd = std::env::current_dir().unwrap_or_default();
|
|
let target = resolve_start_target(app, &cwd, opts.agent.as_deref())?;
|
|
if let Some(group) = crate::catalog::Catalog::parse_group_selector(&target) {
|
|
let members = app.catalog.group_members(group);
|
|
if members.is_empty() {
|
|
bail!("unknown group '{group}'");
|
|
}
|
|
for m in members.iter().rev() {
|
|
stop_one(app, m, force, timeout)?;
|
|
}
|
|
let extra_env = parse_env_list(&opts.env)?;
|
|
let extra_args = parse_extra_args(&opts.args);
|
|
for m in members {
|
|
start_one(app, m, &extra_args, &extra_env, opts.notify, true, None)?;
|
|
}
|
|
return Ok(0);
|
|
}
|
|
let agent = require_agent(app, &target)?;
|
|
stop_one(app, agent, force, timeout)?;
|
|
let extra_env = parse_env_list(&opts.env)?;
|
|
let extra_args = parse_extra_args(&opts.args);
|
|
start_one(app, agent, &extra_args, &extra_env, opts.notify, opts.background, None)
|
|
}
|
|
|
|
/// run: execute the agent command directly, no process management.
|
|
pub fn run(app: &App, target: &str, extra: &[OsString]) -> Result<i32> {
|
|
let agent = require_agent(app, target)?;
|
|
let extra_args: Vec<String> = extra.iter().map(|o| o.to_string_lossy().to_string()).collect();
|
|
let exec = resolve_exec(app, agent, &extra_args, &BTreeMap::new())?;
|
|
if app.dry_run() {
|
|
app.log.dry(format!(
|
|
"would run {} {}",
|
|
exec.program,
|
|
exec.args.join(" ")
|
|
));
|
|
return Ok(0);
|
|
}
|
|
let started = std::time::Instant::now();
|
|
let (prog, prefix) = crate::runner::resolve_program(&exec.program);
|
|
let mut full_args = prefix;
|
|
full_args.extend(exec.args.iter().cloned());
|
|
let status = Command::new(&prog)
|
|
.args(&full_args)
|
|
.envs(&exec.env)
|
|
.status()
|
|
.with_context(|| format!("failed to run {}", exec.program))?;
|
|
let code = status.code().unwrap_or(1);
|
|
app.emit(
|
|
&Event::now(EventKind::Run)
|
|
.with_agent(agent.name.clone())
|
|
.with_cwd(cwd_string())
|
|
.with_args(exec.args.clone())
|
|
.with_env_keys(events::env_keys(&exec.env))
|
|
.with_exit_code(code)
|
|
.with_duration(started.elapsed().as_secs()),
|
|
);
|
|
let _ = app.state.touch(&agent.name);
|
|
Ok(code)
|
|
}
|
|
|
|
/// Current directory as a string (None when unavailable).
|
|
fn cwd_string() -> Option<String> {
|
|
std::env::current_dir().ok().map(|p| p.display().to_string())
|
|
}
|
|
|
|
/// Seconds elapsed since an RFC 3339 timestamp (None when unparsable).
|
|
fn seconds_since(ts: Option<&str>) -> Option<u64> {
|
|
let t = chrono::DateTime::parse_from_rfc3339(ts?).ok()?;
|
|
let now = chrono::Utc::now();
|
|
let secs = now
|
|
.signed_duration_since(t.with_timezone(&chrono::Utc))
|
|
.num_seconds();
|
|
Some(secs.max(0) as u64)
|
|
}
|