Files
agent-manager/src/sessions.rs
T

340 lines
10 KiB
Rust

//! Session registry: one record per agent session (background start/stop
//! cycles and interactive REPL sessions). sessions.json is an index derived
//! from the event journal; it can be rebuilt with 'am doctor --fix'.
use crate::app::App;
use anyhow::{anyhow, Context, Result};
use serde::{Deserialize, Serialize};
use std::collections::BTreeMap;
use std::path::PathBuf;
pub const SESSIONS_VERSION: u32 = 1;
#[derive(Debug, Clone, Serialize, Deserialize, PartialEq)]
pub struct SessionRecord {
pub id: String,
/// "agent" or "repl".
#[serde(default = "kind_agent")]
pub kind: String,
#[serde(default)]
pub agent: Option<String>,
#[serde(default)]
pub pid: Option<u32>,
#[serde(default)]
pub cwd: Option<String>,
#[serde(default)]
pub project: Option<String>,
/// RFC 3339 start timestamp.
pub started_at: String,
/// RFC 3339 end timestamp (None while running).
#[serde(default)]
pub ended_at: Option<String>,
#[serde(default)]
pub exit_code: Option<i32>,
#[serde(default)]
pub log: Option<String>,
#[serde(default)]
pub args: Vec<String>,
/// Variable names only — never values.
#[serde(default)]
pub env_keys: Vec<String>,
/// running | stopped | failed | interrupted
#[serde(default = "status_running")]
pub status: String,
}
fn kind_agent() -> String {
"agent".to_string()
}
fn status_running() -> String {
"running".to_string()
}
#[derive(Debug, Clone, Serialize, Deserialize, Default)]
pub struct SessionsFile {
#[serde(default)]
pub version: u32,
#[serde(default)]
pub sessions: BTreeMap<String, SessionRecord>,
}
/// Path of sessions.json (inside the state directory).
pub fn file_path(app: &App) -> PathBuf {
app.events_dir().join("sessions.json")
}
pub fn load(app: &App) -> Result<SessionsFile> {
let p = file_path(app);
if !p.exists() {
return Ok(SessionsFile {
version: SESSIONS_VERSION,
sessions: BTreeMap::new(),
});
}
let text = std::fs::read_to_string(&p)
.with_context(|| format!("cannot read sessions file {}", p.display()))?;
let mut sf: SessionsFile = serde_json::from_str(&text).map_err(|e| {
anyhow!(
"sessions file {} is corrupt ({e}); run 'am doctor --fix' to repair it",
p.display()
)
})?;
if sf.version < SESSIONS_VERSION {
sf.version = SESSIONS_VERSION;
save(app, &sf)?;
}
Ok(sf)
}
pub fn save(app: &App, sf: &SessionsFile) -> Result<()> {
let p = file_path(app);
if let Some(parent) = p.parent() {
std::fs::create_dir_all(parent)?;
}
let tmp = p.with_extension("json.tmp");
let text = serde_json::to_string_pretty(sf)?;
std::fs::write(&tmp, text)?;
std::fs::rename(&tmp, &p)
.with_context(|| format!("cannot write sessions file {}", p.display()))?;
Ok(())
}
/// Open a new agent session; returns its id.
#[allow(clippy::too_many_arguments)]
pub fn start_agent(
app: &App,
agent: &str,
pid: u32,
cwd: Option<String>,
project: Option<String>,
args: Vec<String>,
env_keys: Vec<String>,
log: Option<String>,
) -> Result<String> {
let mut sf = load(app)?;
let id = crate::repl::session_id();
sf.sessions.insert(
id.clone(),
SessionRecord {
id: id.clone(),
kind: "agent".to_string(),
agent: Some(agent.to_string()),
pid: Some(pid),
cwd,
project,
started_at: crate::installers::now_rfc3339(),
ended_at: None,
exit_code: None,
log,
args,
env_keys,
status: "running".to_string(),
},
);
save(app, &sf)?;
let _ = app.state.bump_sessions();
Ok(id)
}
/// Open a REPL session.
pub fn start_repl(app: &App, id: &str, cwd: Option<String>) -> Result<()> {
let mut sf = load(app)?;
sf.sessions.insert(
id.to_string(),
SessionRecord {
id: id.to_string(),
kind: "repl".to_string(),
agent: None,
pid: None,
cwd,
project: None,
started_at: crate::installers::now_rfc3339(),
ended_at: None,
exit_code: None,
log: None,
args: Vec::new(),
env_keys: Vec::new(),
status: "running".to_string(),
},
);
save(app, &sf)?;
let _ = app.state.bump_sessions();
Ok(())
}
/// Close the running session of an agent: the one matching the pid when
/// given, otherwise the most recent one.
pub fn finish_agent(app: &App, agent: &str, pid: Option<u32>, exit_code: i32, ok: bool) -> Result<()> {
let mut sf = load(app)?;
let mut candidates: Vec<String> = sf
.sessions
.iter()
.filter(|(_, r)| r.kind == "agent" && r.status == "running" && r.agent.as_deref() == Some(agent))
.map(|(k, _)| k.clone())
.collect();
candidates.sort_by(|a, b| sf.sessions[b].started_at.cmp(&sf.sessions[a].started_at));
let key = if let Some(p) = pid {
candidates
.iter()
.find(|k| sf.sessions[*k].pid == Some(p))
.cloned()
.or_else(|| candidates.first().cloned())
} else {
candidates.first().cloned()
};
if let Some(k) = key {
if let Some(r) = sf.sessions.get_mut(&k) {
r.ended_at = Some(crate::installers::now_rfc3339());
r.exit_code = Some(exit_code);
r.status = if ok { "stopped" } else { "failed" }.to_string();
}
save(app, &sf)?;
}
Ok(())
}
/// Close a REPL session.
pub fn finish_repl(app: &App, id: &str) -> Result<()> {
let mut sf = load(app)?;
if let Some(r) = sf.sessions.get_mut(id) {
r.ended_at = Some(crate::installers::now_rfc3339());
r.exit_code = Some(0);
r.status = "stopped".to_string();
save(app, &sf)?;
}
Ok(())
}
/// Mark running sessions whose process is dead as interrupted. Returns the
/// number of sessions reconciled. REPL sessions (no pid) are left untouched:
/// the current REPL has no pid, and only 'am sessions' triggers this.
pub fn reconcile(app: &App) -> Result<usize> {
let mut sf = load(app)?;
let now = crate::installers::now_rfc3339();
let mut changed = 0;
for r in sf.sessions.values_mut() {
if r.status != "running" {
continue;
}
let dead = match r.pid {
Some(p) => !crate::process::is_running(p),
None => false,
};
if dead {
r.status = "interrupted".to_string();
r.ended_at = Some(now.clone());
changed += 1;
}
}
if changed > 0 {
save(app, &sf)?;
}
Ok(changed)
}
/// Duration of a session in seconds (elapsed when still running).
pub fn duration_s(r: &SessionRecord) -> Option<u64> {
let start = chrono::DateTime::parse_from_rfc3339(&r.started_at).ok()?;
let end = match &r.ended_at {
Some(e) => chrono::DateTime::parse_from_rfc3339(e).ok()?,
None => chrono::Utc::now().into(),
};
let secs = end
.with_timezone(&chrono::Utc)
.signed_duration_since(start.with_timezone(&chrono::Utc))
.num_seconds();
Some(secs.max(0) as u64)
}
#[cfg(test)]
mod tests {
use super::*;
use clap::Parser;
fn test_app(tag: &str) -> App {
let dir = tempfile::tempdir().unwrap();
let cfg = dir.path().join("config.yaml");
std::fs::write(
&cfg,
concat!(
"version: \"1.0\"\n",
"settings:\n",
" auto_install_deps: false\n",
" confirm_before_run: false\n",
"agents: []\n",
),
)
.unwrap();
let cli = crate::cli::Cli::parse_from(["am", "--config", cfg.to_str().unwrap()]);
let mut app = App::from_cli(cli).unwrap();
let mut p = app.paths.clone();
p.install_dir = dir.path().join("agents");
p.bin_dir = dir.path().join("agents").join("bin");
p.log_dir = dir.path().join("logs");
p.state_file = dir.path().join("state.json");
p.probe_cache_file = dir.path().join("probe-cache.json");
p.config_dir = Some(dir.path().join("config"));
app.paths = p;
app.state = crate::state::StateStore::new(dir.path().join("state.json"));
let _ = tag;
app
}
#[test]
fn agent_session_lifecycle() {
let app = test_app("sess");
let id = start_agent(&app, "claude-code", 4242, Some("/tmp".into()), None, vec![], vec![], None).unwrap();
let sf = load(&app).unwrap();
let r = sf.sessions.get(&id).unwrap();
assert_eq!(r.status, "running");
assert_eq!(r.pid, Some(4242));
assert_eq!(app.state.load().unwrap().sessions_count, 1);
finish_agent(&app, "claude-code", Some(4242), 0, true).unwrap();
let sf = load(&app).unwrap();
let r = sf.sessions.get(&id).unwrap();
assert_eq!(r.status, "stopped");
assert_eq!(r.exit_code, Some(0));
assert!(r.ended_at.is_some());
assert!(duration_s(r).is_some());
}
#[test]
fn reconcile_marks_dead_pid_interrupted() {
let app = test_app("recon");
let id = start_agent(&app, "aider", u32::MAX - 1, None, None, vec![], vec![], None).unwrap();
let changed = reconcile(&app).unwrap();
assert_eq!(changed, 1);
let sf = load(&app).unwrap();
assert_eq!(sf.sessions.get(&id).unwrap().status, "interrupted");
}
#[test]
fn repl_session_lifecycle() {
let app = test_app("repl");
start_repl(&app, "s1", None).unwrap();
finish_repl(&app, "s1").unwrap();
let sf = load(&app).unwrap();
assert_eq!(sf.sessions.get("s1").unwrap().status, "stopped");
}
#[test]
fn finish_matches_by_pid_first() {
let app = test_app("multi");
start_agent(&app, "jcode", 100, None, None, vec![], vec![], None).unwrap();
start_agent(&app, "jcode", 200, None, None, vec![], vec![], None).unwrap();
finish_agent(&app, "jcode", Some(200), 1, false).unwrap();
let sf = load(&app).unwrap();
let stopped: Vec<_> = sf
.sessions
.values()
.filter(|r| r.status == "stopped" || r.status == "failed")
.collect();
assert_eq!(stopped.len(), 1);
assert_eq!(stopped[0].pid, Some(200));
assert_eq!(stopped[0].exit_code, Some(1));
assert_eq!(stopped[0].status, "failed");
}
}