feat: phase 2 — am lab (benchmark N agents) + plugins d'événements (contrat JSON)

- #62: am lab --agents a,b --task <f> — séquentiel ou --parallel, capture
  durée/exit/sortie/coût estimé, rapport comparatif tableau + --json stable,
  tâches versionnables dans <state>/lab/ (--list), rejouables à l'identique
- #75: plugins/ sous state_dir, déclenchés sur on_install/on_start/on_stop/
  on_update (App::emit), contrat JSON stdin/stdout, timeout configurable,
  échec non bloquant, garde anti-récursion AM_PLUGINS_RUNNING, am plugins
  + am plugins --test <name> (CI), exemples notify + ci-webhook
- intégration REPL (parse, complétion, bannière, is_am_command), help, tips,
  man pages (am-lab.1, am-plugins.1), README, ROADMAP (26/27), v0.5.4
- 334 tests verts (309 + 25)

closes #62
closes #75
This commit is contained in:
2026-08-18 14:17:30 -04:00
parent aa0d8e43eb
commit 4c7866a25a
29 changed files with 2123 additions and 8 deletions
Generated
+1 -1
View File
@@ -21,7 +21,7 @@ dependencies = [
[[package]]
name = "agent-manager"
version = "0.5.3"
version = "0.5.4"
dependencies = [
"anyhow",
"chrono",
+1 -1
View File
@@ -1,6 +1,6 @@
[package]
name = "agent-manager"
version = "0.5.3"
version = "0.5.4"
edition = "2021"
description = "Manage local AI coding agents: list, install, start, stop, update — with automatic dependency handling and a YAML-driven catalog."
license = "MIT"
+2
View File
@@ -136,6 +136,8 @@ sur le PATH) · ⚪ not-installed · 🔴 not-installable (SaaS/desktop)
| am sessions [agent] --status --show --resume | registre des sessions + reprise d'une session arrêtée |
| am history --search --failed --rerun N | historique des commandes (am + shell), filtrable et réexécutable |
| am history 12..25 --save <f> · am playbook <f> | exporte une plage en playbook YAML et la rejoue pas à pas ({{var}}) |
| am lab --agents a,b --task <f> [--parallel] [--json] | benchmark : même tâche sur plusieurs agents (durée, exit, coût) — tâches versionnables dans <state>/lab/ |
| am plugins [--test <nom>] | scripts d'extension sur les événements (on_install/on_start/on_stop/on_update), contrat JSON stdin/stdout, timeout — exemples dans examples/plugins/ |
| am logs <agent> --follow | tail du log d'un agent, en direct |
| am dashboard | tableau de bord TUI temps réel : vue d'ensemble, activité, stats, sessions, projets (Tab/←/→ : onglets, j/k : défilement, q : quitter) |
+2 -2
View File
@@ -443,7 +443,7 @@ moins de 2 secondes.
| [#59](https://git.dracodev.net/Projets/agent-manager/issues/59) | am doctor --watch — vérifications périodiques avec alertes ✅ | M |
| [#60](https://git.dracodev.net/Projets/agent-manager/issues/60) | ✅ Catalogue distant — am catalog update / am catalog add <url> | M |
| [#61](https://git.dracodev.net/Projets/agent-manager/issues/61) | ✅ am suggest — « un agent pour du Python » (tags + usage réel) | M |
| [#62](https://git.dracodev.net/Projets/agent-manager/issues/62) | am lab — benchmark : même tâche sur N agents | L |
| [#62](https://git.dracodev.net/Projets/agent-manager/issues/62) | ✅ am lab — benchmark : même tâche sur N agents | L |
| [#63](https://git.dracodev.net/Projets/agent-manager/issues/63) | ✅ am audit — qui a modifié quoi quand | M |
| [#64](https://git.dracodev.net/Projets/agent-manager/issues/64) | ✅ am update --rollback — backup automatique avant mise à jour | M |
| [#65](https://git.dracodev.net/Projets/agent-manager/issues/65) | ✅ Politiques — pin de version, settings.update_policy | S |
@@ -456,7 +456,7 @@ moins de 2 secondes.
| [#72](https://git.dracodev.net/Projets/agent-manager/issues/72) | ✅ doctor vérifie ollama / llama-server comme n'importe quel outil | S |
| [#73](https://git.dracodev.net/Projets/agent-manager/issues/73) | ✅ am models prune — purge des modèles inutilisés | S |
| [#74](https://git.dracodev.net/Projets/agent-manager/issues/74) | i18n — messages EN/FR | L |
| [#75](https://git.dracodev.net/Projets/agent-manager/issues/75) | Plugin scripts (hooks avancés, intégration CI) | M |
| [#75](https://git.dracodev.net/Projets/agent-manager/issues/75) | ✅ Plugin scripts (hooks avancés, intégration CI) | M |
Critère de sortie du jalon : parc auto-supervisé (crash = redémarrage +
alerte).
+6
View File
@@ -33,6 +33,12 @@ settings:
# monitor_thresholds:
# cpu_pct: 80.0
# mem_mb: 2048
# Plugins d'événements (#75) : scripts exécutés sur on_install/on_start/
# on_stop/on_update, contrat JSON stdin/stdout, timeout, échec non bloquant.
# Le répertoire <state>/plugins/ contient les scripts (+ manifeste .yaml).
# plugins:
# timeout_secs: 10 # temps max d'un plugin (défaut 10 s)
# enabled: [notify] # vide = tous les plugins du répertoire
# Sauvegarde git de l'état (#66) : dépôt où am sync pousse state.json,
# le journal et l'historique (logs/ et backups/ exclus automatiquement).
# sync_repo: https://git.dracodev.net/bruno/am-state.git
+8
View File
@@ -0,0 +1,8 @@
# Tâche de benchmark (issue #62) : les mêmes arguments sont passés à la
# commande de chaque agent. Copier dans <state_dir>/lab/ ou passer un
# chemin explicite :
# am lab --agents claude-code,aider --task hello.yaml --json
name: hello
description: "Répondre à la question la plus simple (sanity check du lab)."
args: ["--prompt", "Réponds uniquement 'coucou' et rien d'autre."]
# timeout: 300 # secondes par agent (défaut 600) — --timeout l'emporte
+31
View File
@@ -0,0 +1,31 @@
# Plugin scripts — agent-manager (issue #75)
#
# Les plugins sont des points d'extension par script, exécutés sur les
# événements du cycle de vie des agents (complément des hooks de l'axe 6) :
# on_install · on_start · on_stop · on_update
#
# Contrat :
# - entrée : l'événement JSON sur stdin
# {"event":"on_start","kind":"start","ts":"…","agent":"claude-code",
# "pid":1234,"cwd":"…","project":"…","args":[…],"env_keys":[…],…}
# - sortie : un objet JSON sur stdout
# {"ok": true, "message": "…", "data": {…}}
# ou {"ok": false, "error": "…"} — signalé en warning, jamais fatal
# - une sortie invalide (pas de JSON) est signalée en warning sans casser am
# - un plugin qui dépasse son timeout est tué (défaut 10 s)
# - les plugins n'ont pas accès aux valeurs des secrets (env_keys seulement)
# - AM_PLUGINS_RUNNING=1 est exporté pendant un plugin : un am imbriqué
# saute les plugins (anti-récursion)
#
# Installation :
# cp notify.py notify.yaml ~/.local/state/agent-manager/plugins/ (Linux)
# cp notify.py notify.yaml %LOCALAPPDATA%\agent-manager\plugins\ (Windows)
#
# Vérification :
# am plugins # liste les plugins découverts
# am plugins --test notify # exécute le contrat sur un événement de test
#
# Configuration (config.yaml) :
# plugins:
# timeout_secs: 10 # temps max d'un plugin (défaut 10 s)
# enabled: [notify] # vide = tous les plugins du répertoire
+38
View File
@@ -0,0 +1,38 @@
#!/usr/bin/env python3
"""Plugin exemple (issue #75) : poste les événements install/update vers un
webhook CI. Si la variable d'environnement CI_WEBHOOK_URL n'est pas définie,
le plugin répond ok sans rien faire (échec non bloquant, jamais fatal).
Exemple : CI_WEBHOOK_URL=https://hooks.example.com/am-ci am update --all
"""
import json
import os
import sys
import urllib.request
def main() -> int:
ev = json.load(sys.stdin)
url = os.environ.get("CI_WEBHOOK_URL")
if not url:
print(json.dumps({"ok": True, "message": "CI_WEBHOOK_URL non définie — rien à faire"}))
return 0
req = urllib.request.Request(
url,
data=json.dumps(ev).encode(),
headers={"Content-Type": "application/json"},
method="POST",
)
try:
with urllib.request.urlopen(req, timeout=5) as resp:
status = resp.status
except Exception as exc: # échec non bloquant : warning seulement
print(json.dumps({"ok": False, "error": f"webhook: {exc}"}))
return 0
print(json.dumps({"ok": True, "message": f"webhook {status}"}))
return 0
if __name__ == "__main__":
sys.exit(main())
+4
View File
@@ -0,0 +1,4 @@
# Manifest du plugin ci-webhook (issue #75).
command: "python ci-webhook.py" # Linux : "python3 ci-webhook.py"
events: [on_install, on_update]
timeout_secs: 10
+24
View File
@@ -0,0 +1,24 @@
#!/usr/bin/env python3
"""Plugin exemple (issue #75) : notifie le démarrage / l'arrêt d'un agent.
Contrat : l'événement JSON arrive sur stdin, la réponse JSON part sur stdout.
Installation : copier notify.py + notify.yaml dans <state_dir>/plugins/ puis
relancer am. Test : am plugins --test notify
"""
import json
import sys
def main() -> int:
ev = json.load(sys.stdin)
agent = ev.get("agent") or "inconnu"
event = ev.get("event", "?")
pid = ev.get("pid")
detail = f" (pid {pid})" if pid else ""
print(json.dumps({"ok": True, "message": f"{event} de {agent}{detail}"}))
return 0
if __name__ == "__main__":
sys.exit(main())
+7
View File
@@ -0,0 +1,7 @@
# Manifest du plugin notify (issue #75).
# command: interpréteur + script (indispensable sur Windows pour .py).
command: "python notify.py" # Linux : "python3 notify.py"
# Événements écoutés (défaut : les 4 événements du cycle de vie).
events: [on_start, on_stop]
# Temps maximal d'exécution avant kill (défaut : settings.plugins.timeout_secs, 10).
timeout_secs: 5
+28
View File
@@ -0,0 +1,28 @@
.ie \n(.g .ds Aq \(aq
.el .ds Aq '
.TH am-lab 1 "lab "
.SH NAME
lab \- Benchmark: run the same task on several agents and compare (issue #62)
.SH SYNOPSIS
\fBlab\fR [\fB\-\-agents\fR] [\fB\-\-task\fR] [\fB\-\-parallel\fR] [\fB\-\-timeout\fR] [\fB\-\-list\fR] [\fB\-h\fR|\fB\-\-help\fR]
.SH DESCRIPTION
Benchmark: run the same task on several agents and compare (issue #62)
.SH OPTIONS
.TP
\fB\-\-agents\fR \fI<AGENTS>\fR
Comma\-separated agent names (or aliases) to benchmark
.TP
\fB\-\-task\fR \fI<FILE>\fR
Task file (a bare name resolves under state_dir/lab/, .yaml appended)
.TP
\fB\-\-parallel\fR
Run every agent concurrently instead of one after another
.TP
\fB\-\-timeout\fR \fI<SECS>\fR
Override the task timeout in seconds
.TP
\fB\-\-list\fR
List the task files of the lab directory
.TP
\fB\-h\fR, \fB\-\-help\fR
Print help
+16
View File
@@ -0,0 +1,16 @@
.ie \n(.g .ds Aq \(aq
.el .ds Aq '
.TH am-plugins 1 "plugins "
.SH NAME
plugins \- List the event plugins and test one (issue #75)
.SH SYNOPSIS
\fBplugins\fR [\fB\-\-test\fR] [\fB\-h\fR|\fB\-\-help\fR]
.SH DESCRIPTION
List the event plugins and test one (issue #75)
.SH OPTIONS
.TP
\fB\-\-test\fR \fI<NAME>\fR
Run one plugin against a synthetic event and print its response
.TP
\fB\-h\fR, \fB\-\-help\fR
Print help
+8 -2
View File
@@ -1,6 +1,6 @@
.ie \n(.g .ds Aq \(aq
.el .ds Aq '
.TH am 1 "am 0.5.3"
.TH am 1 "am 0.5.4"
.SH NAME
am \- agent\-manager (am) — manage local AI coding agents
.SH SYNOPSIS
@@ -171,6 +171,12 @@ Search the command history (am and shell commands)
am\-playbook(1)
Replay a saved playbook step by step with confirmation (issue #51)
.TP
am\-lab(1)
Benchmark: run the same task on several agents and compare (issue #62)
.TP
am\-plugins(1)
List the event plugins and test one (issue #75)
.TP
am\-info(1)
Show detailed information about one agent
.TP
@@ -210,4 +216,4 @@ Export the configuration and installation state (backup)
am\-import(1)
Import a previously exported configuration and state
.SH VERSION
v0.5.3
v0.5.4
+3 -1
View File
@@ -131,11 +131,13 @@ impl App {
}
/// Append an event to the journal; failures are never fatal (logged at
/// verbose level only).
/// verbose level only). The plugins subscribed to the event kind are
/// then fired (issue #75) — non-blocking, never fatal either.
pub fn emit(&self, event: &crate::events::Event) {
if let Err(e) = crate::events::append(&self.events_dir(), event) {
self.log.verbose(&format!("cannot write event journal: {e:#}"));
}
crate::plugins::dispatch(self, event);
}
/// Ask the user for confirmation on stderr. Accepts y/yes/o/oui.
+32
View File
@@ -423,6 +423,10 @@ pub enum Command {
#[arg(long, value_name = "NAME=VALUE")]
var: Vec<String>,
},
/// Benchmark: run the same task on several agents and compare (issue #62)
Lab(LabArgs),
/// List the event plugins and test one (issue #75)
Plugins(PluginsArgs),
/// Show detailed information about one agent
Info {
/// Agent name or alias
@@ -570,6 +574,34 @@ pub struct ModelsArgs {
pub days: Option<u64>,
}
/// Arguments of the lab command (issue #62).
#[derive(Args, Debug, Clone, Default)]
pub struct LabArgs {
/// Comma-separated agent names (or aliases) to benchmark
#[arg(long, value_name = "AGENTS", required_unless_present = "list")]
pub agents: Option<String>,
/// Task file (a bare name resolves under state_dir/lab/, .yaml appended)
#[arg(long, value_name = "FILE", required_unless_present = "list")]
pub task: Option<PathBuf>,
/// Run every agent concurrently instead of one after another
#[arg(long, action = ArgAction::SetTrue)]
pub parallel: bool,
/// Override the task timeout in seconds
#[arg(long, value_name = "SECS")]
pub timeout: Option<u64>,
/// List the task files of the lab directory
#[arg(long, action = ArgAction::SetTrue)]
pub list: bool,
}
/// Arguments of the plugins command (issue #75).
#[derive(Args, Debug, Clone, Default)]
pub struct PluginsArgs {
/// Run one plugin against a synthetic event and print its response
#[arg(long, value_name = "NAME")]
pub test: Option<String>,
}
/// Remote catalog subcommands (issues #60 #67).
#[derive(Subcommand, Debug, Clone)]
pub enum CatalogCmd {
+101
View File
@@ -0,0 +1,101 @@
//! lab: benchmark the same task on several agents (issue #62).
use super::*;
use crate::cli::LabArgs;
use crate::lab::{self, LabReport};
use anyhow::Result;
pub fn run(app: &App, args: &LabArgs) -> Result<i32> {
if args.list {
let tasks = lab::list_tasks(app);
if app.json() {
let names: Vec<String> = tasks.iter().map(|p| p.display().to_string()).collect();
crate::output::print_json(&names);
return Ok(0);
}
if tasks.is_empty() {
app.log.info(&format!(
"no tasks in {} — write a task.yaml (name, args, timeout) or run 'lab --agents a,b --task hello'",
lab::dir(app).display()
));
return Ok(0);
}
let mut table = crate::output::Table::new(vec!["TASK"]);
for p in tasks {
table.row(vec![p.display().to_string()]);
}
print!("{}", table.render());
return Ok(0);
}
let path = match args.task.as_deref() {
Some(p) => lab::resolve_task(app, p),
None => anyhow::bail!("--task <file> is required"),
};
if !path.exists() {
anyhow::bail!(
"task not found: {} — write it or use 'lab --list' to see the tasks of {}",
path.display(),
lab::dir(app).display()
);
}
let task = lab::load_task(&path)?;
let agents = lab::split_agents(args.agents.as_deref().unwrap_or(""));
if agents.is_empty() {
anyhow::bail!(
"--agents expects at least one agent name, got '{}'",
args.agents.as_deref().unwrap_or("")
);
}
app.log.info(&format!(
"lab '{}' — {} agent(s), {}",
task.name,
agents.len(),
if args.parallel { "parallèle" } else { "séquentiel" }
));
let report = lab::run_benchmark(app, &agents, &task, args.parallel, args.timeout)?;
if app.json() {
crate::output::print_json(&report);
} else {
print_report(app, &report);
}
// Exit code: 0 when every run succeeded, 1 otherwise (CI-friendly).
Ok(if report.runs.iter().all(|r| r.status == "ok") {
0
} else {
1
})
}
fn print_report(app: &App, report: &LabReport) {
let mut table = crate::output::Table::new(vec![
"AGENT", "STATUS", "EXIT", "DURATION", "COST", "OUTPUT",
]);
for r in &report.runs {
table.row(vec![
r.agent.clone(),
r.status.clone(),
r.exit_code
.map(|c| c.to_string())
.unwrap_or_else(|| "-".into()),
format!("{:.2}s", r.duration_ms as f64 / 1000.0),
r.cost_usd
.map(|c| format!("${c:.4}"))
.unwrap_or_else(|| "-".into()),
crate::output::truncate(r.output.lines().next().unwrap_or(""), 40),
]);
}
print!("{}", table.render());
let s = &report.summary;
let mut line = format!(
"{} run(s) — {} ok, {} failed, {} timeout",
s.runs, s.ok, s.failed, s.timeout
);
if let Some(f) = &s.fastest {
line.push_str(&format!(" — fastest: {f}"));
}
if let Some(c) = &s.cheapest {
line.push_str(&format!(" — cheapest: {c}"));
}
app.log.info(&line);
}
+4
View File
@@ -15,6 +15,7 @@ pub mod history_cmd;
pub mod info_cmd;
pub mod init_cmd;
pub mod install_cmd;
pub mod lab_cmd;
pub mod list_cmd;
pub mod log_cmd;
pub mod logs_cmd;
@@ -38,6 +39,7 @@ pub mod status_cmd;
pub mod suggest_cmd;
pub mod sync_cmd;
pub mod playbook_cmd;
pub mod plugins_cmd;
pub mod theme_cmd;
pub mod timeline_cmd;
pub mod tip_cmd;
@@ -195,6 +197,8 @@ pub fn execute_command(app: &App, cmd: &Command) -> Result<i32> {
save.as_deref(),
),
Command::Playbook { path, var } => playbook_cmd::run(app, path, var),
Command::Lab(args) => lab_cmd::run(app, args),
Command::Plugins(args) => plugins_cmd::run(app, args),
Command::Info { agent } => info_cmd::run(app, agent),
Command::Init { force, template } => {
init_cmd::run(app, *force, template.as_deref())
+102
View File
@@ -0,0 +1,102 @@
//! plugins: list the event plugins and test one against a synthetic event
//! (issue #75). 'am plugins --test <name>' is the CI hook: it exercises the
//! JSON contract end to end and exits non-zero on failure.
use super::*;
use crate::cli::PluginsArgs;
use crate::plugins::{self, PluginStatus};
use anyhow::Result;
pub fn run(app: &App, args: &PluginsArgs) -> Result<i32> {
let dir = plugins::dir(app);
let all = plugins::discover(app, &dir);
if let Some(name) = &args.test {
let Some(plugin) = all.iter().find(|p| &p.name == name) else {
anyhow::bail!(
"unknown plugin '{name}' — 'am plugins' lists the plugins of {}",
dir.display()
);
};
// Synthetic event (a start, agent 'test') to exercise the contract.
let event = crate::events::Event::now(crate::events::EventKind::Start)
.with_agent("test")
.with_args(vec!["--probe".to_string()]);
let r = plugins::run_plugin(app, plugin, "on_start", &event);
if app.json() {
crate::output::print_json(&serde_json::json!({
"name": plugin.name,
"status": r.status.as_str(),
"exit_code": r.exit_code,
"duration_ms": r.duration_ms,
"response": r.response,
"stdout": r.stdout,
"stderr": r.stderr,
"error": r.error,
}));
} else {
app.log.info(&format!(
"plugin '{}' — {} in {}ms",
plugin.name,
r.status.as_str(),
r.duration_ms
));
if let Some(code) = r.exit_code {
app.log.info(&format!("exit code: {code}"));
}
if let Some(resp) = &r.response {
if let Some(msg) = &resp.message {
app.log.info(&format!("message: {msg}"));
}
if let Some(err) = &resp.error {
app.log.warn(&format!("error: {err}"));
}
}
if !r.stdout.trim().is_empty() {
app.log.info(&format!("stdout: {}", r.stdout.trim()));
}
if !r.stderr.trim().is_empty() {
app.log.warn(&format!("stderr: {}", r.stderr.trim()));
}
if let Some(e) = &r.error {
app.log.warn(&format!("{e}"));
}
}
return Ok(if r.status == PluginStatus::Ok { 0 } else { 1 });
}
if app.json() {
let items: Vec<serde_json::Value> = all
.iter()
.map(|p| {
serde_json::json!({
"name": p.name,
"events": p.events,
"timeout_secs": p.timeout.as_secs(),
"script": p.file.display().to_string(),
})
})
.collect();
crate::output::print_json(&items);
return Ok(0);
}
if all.is_empty() {
app.log.info(&format!(
"no plugins in {} — copy one from examples/plugins/ (notify, ci-webhook)",
dir.display()
));
return Ok(0);
}
let mut table = crate::output::Table::new(vec!["PLUGIN", "EVENTS", "TIMEOUT", "SCRIPT"]);
for p in &all {
table.row(vec![
p.name.clone(),
p.events.join(", "),
format!("{}s", p.timeout.as_secs()),
p.file.display().to_string(),
]);
}
print!("{}", table.render());
Ok(0)
}
+12
View File
@@ -260,6 +260,18 @@ pub static SECTIONS: &[TipSection] = &[
options: &[("--save <f>", "fichier playbook"), ("playbook <f>", "rejoue pas à pas avec confirmation")],
example: "history 12..25 --save deploy.yaml && playbook deploy.yaml",
},
TipEntry {
usage: "lab --agents claude-code,aider --task hello.yaml",
about: "benchmark : même tâche sur plusieurs agents (durée, exit, coût)",
options: &[("--parallel", "exécution concurrente"), ("--json", "rapport comparatif stable (scripts/CI)")],
example: "lab --agents claude-code,aider --task bench/hello.yaml --json",
},
TipEntry {
usage: "plugins --test notify",
about: "scripts d'extension sur les événements (contrat JSON stdin/stdout, timeout)",
options: &[("--test <name>", "exécute un plugin sur un événement de test (CI)")],
example: "plugins && plugins --test notify",
},
TipEntry {
usage: "timeline",
about: "une vue chronologique de toute l'activité",
+19
View File
@@ -136,6 +136,22 @@ pub struct Settings {
/// Push automatically when the REPL exits (issue #66, opt-in).
#[serde(default)]
pub sync_on_exit: Option<bool>,
/// Plugin scripts (issue #75): default timeout and enable list.
#[serde(default)]
pub plugins: Option<PluginSettings>,
}
/// Plugin scripts settings (issue #75): how long one plugin may run and
/// which plugins of the directory are enabled (empty = every plugin).
#[derive(Debug, Clone, Serialize, Deserialize, Default)]
#[serde(deny_unknown_fields)]
pub struct PluginSettings {
/// Default timeout in seconds for every plugin (default 10).
#[serde(default)]
pub timeout_secs: Option<u64>,
/// Only run these plugin names (empty = all the plugins of the dir).
#[serde(default)]
pub enabled: Option<Vec<String>>,
}
/// Token price model: USD per million tokens (issue #49).
@@ -750,6 +766,9 @@ pub fn merge(base: &mut Config, overlay: Config) {
if o.sync_on_exit.is_some() {
s.sync_on_exit = o.sync_on_exit;
}
if o.plugins.is_some() {
s.plugins = o.plugins;
}
for (k, v) in o.hooks {
s.hooks.insert(k, v);
}
+3
View File
@@ -46,6 +46,8 @@ pub enum EventKind {
Cost,
/// Threshold alert from 'am monitor' (issue #50).
Alert,
/// One benchmark run of 'am lab' (issue #62).
Lab,
}
impl EventKind {
@@ -70,6 +72,7 @@ impl EventKind {
EventKind::Prune => "prune",
EventKind::Cost => "cost",
EventKind::Alert => "alert",
EventKind::Lab => "lab",
}
}
}
+38
View File
@@ -339,6 +339,44 @@ pub static HELP_SPECS: &[HelpSpec] = &[
HelpExample { desc: "Replay without asking, with a variable.", code: "playbook deploy.yaml --yes --var agent=claude-code" },
],
},
HelpSpec {
name: "lab",
category: "Commands",
usage: "lab --agents a,b --task <file> {flags}",
about: "Benchmark: run the same task on several agents and compare duration, exit code, output and estimated cost (issue #62).",
search_terms: &["benchmark", "compare", "comparer", "bench", "perf"],
flags: &[
HelpFlag { short: "", long: "--agents", value: "LIST", desc: "Comma-separated agent names (or aliases) to benchmark (required)" },
HelpFlag { short: "", long: "--task", value: "FILE", desc: "Task file (bare name resolves under state_dir/lab/, .yaml appended)" },
HelpFlag { short: "", long: "--parallel", value: "", desc: "Run every agent concurrently instead of one after another" },
HelpFlag { short: "", long: "--timeout", value: "SECS", desc: "Override the task timeout in seconds" },
HelpFlag { short: "", long: "--list", value: "", desc: "List the task files of the lab directory" },
],
subcommands: &[],
parameters: &[],
io: None,
examples: &[
HelpExample { desc: "Run the task sequentially on two agents.", code: "lab --agents claude-code,aider --task hello.yaml" },
HelpExample { desc: "Machine-readable comparative report.", code: "lab --agents claude-code,codex --task bench.yaml --json --parallel" },
],
},
HelpSpec {
name: "plugins",
category: "Commands",
usage: "plugins {flags}",
about: "List the event plugins of the state directory and test one against a synthetic event (issue #75).",
search_terms: &["plugin", "extension", "hooks", "notify", "webhook"],
flags: &[
HelpFlag { short: "", long: "--test", value: "NAME", desc: "Run one plugin on a synthetic event and print its response (CI hook)" },
],
subcommands: &[],
parameters: &[],
io: Some(("JSON (stdin)", "JSON object (stdout)")),
examples: &[
HelpExample { desc: "List the installed plugins.", code: "plugins" },
HelpExample { desc: "Exercise the contract of one plugin.", code: "plugins --test notify" },
],
},
HelpSpec {
name: "top",
category: "Commands",
+629
View File
@@ -0,0 +1,629 @@
//! am lab (issue #62): run the same task on several agents and compare
//! duration, exit code, output and estimated cost. Tasks are plain YAML
//! files living in `<state_dir>/lab/` — versionable and replayable:
//! running the same file again reproduces the same benchmark.
//!
//! A task is the list of arguments passed to each agent's run command
//! (plus an optional working directory and timeout).
use crate::app::App;
use anyhow::{bail, Context, Result};
use serde::{Deserialize, Serialize};
use std::collections::BTreeMap;
use std::path::{Path, PathBuf};
use std::process::{Command, Stdio};
use std::time::{Duration, Instant};
use wait_timeout::ChildExt;
/// Default timeout for one task run (seconds).
pub const DEFAULT_TASK_TIMEOUT_SECS: u64 = 600;
/// Name of the tasks directory inside the state directory.
pub const DIR_NAME: &str = "lab";
/// A benchmark task: a versionable YAML file run on every agent.
#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct LabTask {
/// Short name of the task (required).
pub name: String,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub description: Option<String>,
/// Arguments passed to every agent's run command.
#[serde(default)]
pub args: Vec<String>,
/// Working directory of the run (default: the current directory).
#[serde(default, skip_serializing_if = "Option::is_none")]
pub cwd: Option<String>,
/// Timeout in seconds (default 600).
#[serde(default, skip_serializing_if = "Option::is_none")]
pub timeout: Option<u64>,
}
/// Outcome of one agent on the task.
#[derive(Debug, Clone, Default, Serialize)]
pub struct LabRun {
pub agent: String,
/// ok | failed | timeout | error
pub status: String,
#[serde(skip_serializing_if = "Option::is_none")]
pub exit_code: Option<i32>,
pub duration_ms: u64,
/// Captured stdout of the run.
pub output: String,
#[serde(skip_serializing_if = "String::is_empty")]
pub stderr: String,
#[serde(skip_serializing_if = "Option::is_none")]
pub tokens_in: Option<u64>,
#[serde(skip_serializing_if = "Option::is_none")]
pub tokens_out: Option<u64>,
#[serde(skip_serializing_if = "Option::is_none")]
pub cost_usd: Option<f64>,
#[serde(skip_serializing_if = "Option::is_none")]
pub error: Option<String>,
}
/// The comparative report (also the --json contract: stable fields,
/// runs sorted by agent name).
#[derive(Debug, Clone, Serialize)]
pub struct LabReport {
pub task: String,
pub runs: Vec<LabRun>,
pub summary: LabSummary,
}
#[derive(Debug, Clone, Default, Serialize)]
pub struct LabSummary {
pub runs: usize,
pub ok: usize,
pub failed: usize,
pub timeout: usize,
#[serde(skip_serializing_if = "Option::is_none")]
pub fastest: Option<String>,
#[serde(skip_serializing_if = "Option::is_none")]
pub cheapest: Option<String>,
}
/// The tasks directory of the state directory.
pub fn dir(app: &App) -> PathBuf {
app.events_dir().join(DIR_NAME)
}
/// Resolve a task path: a bare name goes into the lab directory (with
/// .yaml/.yml appended when the exact name does not exist), an explicit
/// path is honored as-is.
pub fn resolve_task(app: &App, dest: &Path) -> PathBuf {
let has_sep = dest.components().count() > 1
|| dest.components().any(|c| {
matches!(
c,
std::path::Component::ParentDir | std::path::Component::CurDir
)
});
if has_sep || dest.is_absolute() {
return dest.to_path_buf();
}
let base = dir(app).join(dest);
if base.exists() {
return base;
}
for ext in ["yaml", "yml"] {
let candidate = dir(app).join(format!("{}.{}", dest.display(), ext));
if candidate.exists() {
return candidate;
}
}
base
}
/// Task files of the lab directory (yaml/yml), sorted by name.
pub fn list_tasks(app: &App) -> Vec<PathBuf> {
let mut out: Vec<PathBuf> = std::fs::read_dir(dir(app))
.map(|rd| {
rd.flatten()
.map(|e| e.path())
.filter(|p| p.is_file())
.filter(|p| {
matches!(
p.extension().and_then(|e| e.to_str()),
Some("yaml" | "yml")
)
})
.collect()
})
.unwrap_or_default();
out.sort();
out
}
/// Load and validate a task file.
pub fn load_task(path: &Path) -> Result<LabTask> {
let text = std::fs::read_to_string(path)
.with_context(|| format!("cannot read {}", path.display()))?;
let task: LabTask = serde_yaml::from_str(&text)
.with_context(|| format!("{}: invalid task YAML", path.display()))?;
if task.name.trim().is_empty() {
bail!("{}: 'name' is required", path.display());
}
Ok(task)
}
/// Split the --agents list ("a,b, c") into trimmed names.
pub fn split_agents(list: &str) -> Vec<String> {
list.split(',')
.map(|s| s.trim().to_string())
.filter(|s| !s.is_empty())
.collect()
}
/// A pre-resolved run context: owned data, safe to move across threads.
struct RunCtx {
agent: String,
program: String,
args: Vec<String>,
env: BTreeMap<String, String>,
cwd: Option<PathBuf>,
timeout: Duration,
cost_model: Option<String>,
sheet: BTreeMap<String, crate::config::CostModel>,
}
/// Run the task on every agent and build the comparative report.
/// Sequential by default; `parallel` spawns one thread per agent.
pub fn run_benchmark(
app: &App,
agents: &[String],
task: &LabTask,
parallel: bool,
timeout_override: Option<u64>,
) -> Result<LabReport> {
let timeout = Duration::from_secs(
timeout_override.unwrap_or(task.timeout.unwrap_or(DEFAULT_TASK_TIMEOUT_SECS)),
);
let task_cwd = task.cwd.as_deref().map(crate::config::expand_path);
if let Some(cwd) = &task_cwd {
if !cwd.is_dir() {
bail!("task cwd does not exist: {}", cwd.display());
}
}
let sheet = app
.config
.settings
.cost_models
.clone()
.unwrap_or_else(crate::config::default_cost_models);
let mut ctxs = Vec::new();
for name in agents {
let agent = crate::commands::require_agent(app, name)?;
let exec = crate::commands::resolve_exec(app, agent, &task.args, &BTreeMap::new())?;
ctxs.push(RunCtx {
agent: agent.name.clone(),
program: exec.program,
args: exec.args,
env: exec.env,
cwd: task_cwd.clone(),
timeout,
cost_model: agent.cost_model.clone(),
sheet: sheet.clone(),
});
}
let mut runs: Vec<LabRun> = if parallel {
let handles: Vec<_> = ctxs
.into_iter()
.map(|ctx| std::thread::spawn(move || run_one(&ctx)))
.collect();
handles
.into_iter()
.map(|h| {
h.join().unwrap_or_else(|_| LabRun {
agent: "?".to_string(),
status: "error".to_string(),
error: Some("runner thread panicked".to_string()),
..Default::default()
})
})
.collect()
} else {
ctxs.iter().map(run_one).collect()
};
// Stable order for --json: always sorted by agent name.
runs.sort_by(|a, b| a.agent.cmp(&b.agent));
// Journal one Lab event per run (axe 1: everything is traceable).
for run in &runs {
let mut ev = crate::events::Event::now(crate::events::EventKind::Lab)
.with_agent(run.agent.clone())
.with_reason(run.status.clone())
.with_args(task.args.clone())
.with_duration(run.duration_ms / 1000);
if let Some(code) = run.exit_code {
ev = ev.with_exit_code(code);
}
if let Some(cost) = run.cost_usd {
ev = ev
.with_usage(run.tokens_in.unwrap_or(0), run.tokens_out.unwrap_or(0))
.with_cost(cost);
}
app.emit(&ev);
}
Ok(LabReport {
task: task.name.clone(),
summary: summarize(&runs),
runs,
})
}
/// Run one agent on the task with a timeout; captures stdout/stderr.
fn run_one(ctx: &RunCtx) -> LabRun {
let started = Instant::now();
let mut cmd = Command::new(&ctx.program);
cmd.args(&ctx.args)
.envs(&ctx.env)
.stdin(Stdio::null())
.stdout(Stdio::piped())
.stderr(Stdio::piped());
if let Some(cwd) = &ctx.cwd {
cmd.current_dir(cwd);
}
set_process_group(&mut cmd);
let mut child = match cmd.spawn() {
Ok(c) => c,
Err(e) => {
return LabRun {
agent: ctx.agent.clone(),
status: "error".to_string(),
error: Some(format!("cannot spawn: {e:#}")),
duration_ms: started.elapsed().as_millis() as u64,
..Default::default()
}
}
};
let stdout = read_pipe(child.stdout.take());
let stderr = read_pipe(child.stderr.take());
let deadline = started + ctx.timeout;
let outcome = loop {
match child.wait_timeout(Duration::from_millis(200)) {
Ok(Some(status)) => break Ok(status),
Ok(None) => {
if Instant::now() >= deadline {
crate::process::kill_tree(child.id());
let _ = child.wait();
break Err("timeout".to_string());
}
}
Err(e) => {
crate::process::kill_tree(child.id());
let _ = child.wait();
break Err(format!("cannot wait: {e:#}"));
}
}
};
let out = stdout.join().unwrap_or_default();
let err = stderr.join().unwrap_or_default();
let (status, exit_code, error) = match outcome {
Ok(st) => {
let code = st.code().unwrap_or(-1);
(
if code == 0 { "ok" } else { "failed" }.to_string(),
Some(code),
None,
)
}
Err(e) => {
let kind = if e == "timeout" { "timeout" } else { "error" };
(kind.to_string(), None, Some(e))
}
};
let mut run = LabRun {
agent: ctx.agent.clone(),
status,
exit_code,
duration_ms: started.elapsed().as_millis() as u64,
output: out,
stderr: err,
error,
..Default::default()
};
// Estimated cost when the agent exposes its usage (axe 1, issue #49).
let cwd_str = ctx.cwd.as_deref().map(|p| p.display().to_string());
if let Some(usage) = crate::costs::collect_usage(&ctx.agent, cwd_str.as_deref()) {
let cost = crate::costs::estimate_cost(ctx.cost_model.as_deref(), usage, &ctx.sheet);
run.tokens_in = Some(usage.tokens_in);
run.tokens_out = Some(usage.tokens_out);
run.cost_usd = Some(cost);
}
run
}
/// Compute the summary rows of the report.
fn summarize(runs: &[LabRun]) -> LabSummary {
let mut s = LabSummary {
runs: runs.len(),
..Default::default()
};
let mut fastest: Option<(&str, u64)> = None;
let mut cheapest: Option<(&str, f64)> = None;
for r in runs {
match r.status.as_str() {
"ok" => s.ok += 1,
"failed" => s.failed += 1,
"timeout" => s.timeout += 1,
_ => {}
}
if r.status == "ok" && fastest.map(|(_, d)| r.duration_ms < d).unwrap_or(true) {
fastest = Some((&r.agent, r.duration_ms));
}
if let Some(cost) = r.cost_usd {
if cheapest.map(|(_, c)| cost < c).unwrap_or(true) {
cheapest = Some((&r.agent, cost));
}
}
}
s.fastest = fastest.map(|(a, _)| a.to_string());
s.cheapest = cheapest.map(|(a, _)| a.to_string());
s
}
/// Read a child pipe to the end on a dedicated thread (avoids pipe
/// deadlocks while the parent waits with a timeout).
fn read_pipe<R: std::io::Read + Send + 'static>(pipe: Option<R>) -> std::thread::JoinHandle<String> {
std::thread::spawn(move || {
use std::io::Read;
let mut buf = String::new();
if let Some(mut p) = pipe {
let _ = p.read_to_string(&mut buf);
}
buf
})
}
/// Put the child in its own process group so `kill_tree` can stop it with
/// its whole subtree on timeout.
fn set_process_group(cmd: &mut Command) {
#[cfg(unix)]
{
use std::os::unix::process::CommandExt;
cmd.process_group(0);
}
#[cfg(windows)]
{
use std::os::windows::process::CommandExt;
const CREATE_NEW_PROCESS_GROUP: u32 = 0x0000_0200;
cmd.creation_flags(CREATE_NEW_PROCESS_GROUP);
}
}
// ---------------------------------------------------------------------------
// Tests
// ---------------------------------------------------------------------------
#[cfg(test)]
mod tests {
use super::*;
use crate::cli::Cli;
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 = Cli::parse_from(["am", "--config", cfg.to_str().unwrap(), "list"]);
let mut app = App::from_cli(cli).unwrap();
app.paths.state_file = dir.path().join(format!("state-{tag}.json"));
app.paths.log_dir = dir.path().join("logs");
app
}
/// An app whose catalog defines three shell-based fake agents:
/// bench-ok (echoes its first argument), bench-fail (exit 3),
/// bench-slow (infinite loop, for the timeout test).
fn bench_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",
" - name: bench-ok\n",
" installable: false\n",
" run: sh\n",
" args: [\"-c\", \"echo out: $0\"]\n",
" - name: bench-fail\n",
" installable: false\n",
" run: sh\n",
" args: [\"-c\", \"echo boom >&2; exit 3\"]\n",
" - name: bench-slow\n",
" installable: false\n",
" run: sh\n",
" args: [\"-c\", \"while true; do :; done\"]\n",
),
)
.unwrap();
let cli = Cli::parse_from(["am", "--config", cfg.to_str().unwrap(), "list"]);
let mut app = App::from_cli(cli).unwrap();
app.paths.state_file = dir.path().join(format!("state-{tag}.json"));
app.paths.log_dir = dir.path().join("logs");
app
}
#[test]
fn split_agents_trims_and_filters() {
assert_eq!(split_agents("a,b"), vec!["a".to_string(), "b".to_string()]);
assert_eq!(split_agents(" a , b, c "), vec!["a", "b", "c"]);
assert_eq!(split_agents(" , "), Vec::<String>::new());
}
#[test]
fn task_loads_with_defaults() {
let dir = tempfile::tempdir().unwrap();
let f = dir.path().join("hello.yaml");
std::fs::write(
&f,
"name: hello\ndescription: test\nargs: [\"--prompt\", \"coucou\"]\n",
)
.unwrap();
let t = load_task(&f).unwrap();
assert_eq!(t.name, "hello");
assert_eq!(t.args, vec!["--prompt".to_string(), "coucou".to_string()]);
assert_eq!(t.timeout, None);
assert!(load_task(&dir.path().join("missing.yaml")).is_err());
std::fs::write(&f, "args: [x]\n").unwrap();
assert!(load_task(&f).is_err(), "name is required");
}
#[test]
fn bare_name_resolves_into_the_lab_dir() {
let app = test_app("dest");
let resolved = resolve_task(&app, Path::new("hello"));
assert_eq!(resolved, app.events_dir().join(DIR_NAME).join("hello"));
let abs = PathBuf::from("C:/tmp/hello.yaml");
assert_eq!(resolve_task(&app, &abs), abs);
}
#[test]
fn list_tasks_filters_yaml_files() {
let app = test_app("list");
let d = dir(&app);
std::fs::create_dir_all(&d).unwrap();
std::fs::write(d.join("a.yaml"), "name: a\n").unwrap();
std::fs::write(d.join("b.yml"), "name: b\n").unwrap();
std::fs::write(d.join("notes.md"), "nope").unwrap();
let tasks = list_tasks(&app);
assert_eq!(tasks.len(), 2);
assert!(tasks[0].display().to_string().ends_with("a.yaml"));
assert!(tasks[1].display().to_string().ends_with("b.yml"));
}
#[test]
fn benchmark_runs_sequentially_and_reports() {
let app = bench_app("seq");
let task = LabTask {
name: "hello".into(),
description: None,
args: vec!["hello".to_string()],
cwd: None,
timeout: Some(30),
};
let agents = vec!["bench-ok".to_string(), "bench-fail".to_string()];
let report = run_benchmark(&app, &agents, &task, false, None).unwrap();
assert_eq!(report.task, "hello");
assert_eq!(report.runs.len(), 2);
let ok = report.runs.iter().find(|r| r.agent == "bench-ok").unwrap();
assert_eq!(ok.status, "ok");
assert_eq!(ok.exit_code, Some(0));
assert!(ok.output.contains("out: hello"), "got: {}", ok.output);
let fail = report.runs.iter().find(|r| r.agent == "bench-fail").unwrap();
assert_eq!(fail.status, "failed");
assert_eq!(fail.exit_code, Some(3));
assert!(fail.stderr.contains("boom"));
// Summary counts.
assert_eq!(report.summary.ok, 1);
assert_eq!(report.summary.failed, 1);
assert_eq!(report.summary.fastest.as_deref(), Some("bench-ok"));
}
#[test]
fn benchmark_is_replayable_identically() {
let app = bench_app("replay");
let task = LabTask {
name: "hello".into(),
description: None,
args: vec!["hello".to_string()],
cwd: None,
timeout: Some(30),
};
let agents = vec!["bench-ok".to_string()];
let r1 = run_benchmark(&app, &agents, &task, false, None).unwrap();
let r2 = run_benchmark(&app, &agents, &task, false, None).unwrap();
assert_eq!(r1.runs.len(), r2.runs.len());
assert_eq!(r1.runs[0].status, r2.runs[0].status);
assert_eq!(r1.runs[0].exit_code, r2.runs[0].exit_code);
assert_eq!(r1.runs[0].output, r2.runs[0].output);
// The --json contract is stable: same fields on both serializations.
let j1 = serde_json::to_string(&r1).unwrap();
let j2 = serde_json::to_string(&r2).unwrap();
assert!(j1.contains("\"task\":\"hello\""));
assert!(j1.contains("\"status\":\"ok\""));
assert!(j2.contains("\"status\":\"ok\""));
}
#[test]
fn timeout_kills_the_run() {
let app = bench_app("to");
let task = LabTask {
name: "slow".into(),
description: None,
args: vec![],
cwd: None,
timeout: Some(1),
};
let agents = vec!["bench-slow".to_string()];
let report = run_benchmark(&app, &agents, &task, false, None).unwrap();
assert_eq!(report.runs[0].status, "timeout");
assert_eq!(report.summary.timeout, 1);
}
#[test]
fn parallel_mode_runs_every_agent() {
let app = bench_app("par");
let task = LabTask {
name: "hello".into(),
description: None,
args: vec!["hello".to_string()],
cwd: None,
timeout: Some(30),
};
let agents = vec!["bench-ok".to_string(), "bench-fail".to_string()];
let report = run_benchmark(&app, &agents, &task, true, None).unwrap();
assert_eq!(report.runs.len(), 2);
// Runs are sorted by agent name: stable --json.
assert_eq!(report.runs[0].agent, "bench-fail");
assert_eq!(report.runs[1].agent, "bench-ok");
assert_eq!(report.summary.ok, 1);
}
#[test]
fn missing_cwd_is_an_error() {
let app = bench_app("cwd");
let task = LabTask {
name: "bad".into(),
description: None,
args: vec![],
cwd: Some("C:/does/not/exist-anywhere".to_string()),
timeout: None,
};
assert!(run_benchmark(&app, &["bench-ok".to_string()], &task, false, None).is_err());
}
#[test]
fn unknown_agent_is_an_error() {
let app = bench_app("unk");
let task = LabTask {
name: "t".into(),
description: None,
args: vec![],
cwd: None,
timeout: None,
};
assert!(run_benchmark(&app, &["ghost".to_string()], &task, false, None).is_err());
}
}
+2
View File
@@ -31,6 +31,7 @@ pub mod help;
pub mod history;
pub mod hooks;
pub mod installers;
pub mod lab;
pub mod models;
pub mod nav;
pub mod output;
@@ -46,6 +47,7 @@ pub mod shell;
pub mod state;
pub mod sync;
pub mod playbook;
pub mod plugins;
pub mod tables;
pub mod theme;
pub mod toolchain;
+794
View File
@@ -0,0 +1,794 @@
//! Plugin scripts (issue #75): executable scripts in `<state_dir>/plugins/`
//! subscribed to lifecycle events, complementing the axe-6 hooks.
//!
//! Contract: the event JSON arrives on stdin, the plugin answers with JSON
//! on stdout. Failures are non-blocking (a warning at most) and a per-plugin
//! timeout kills runaway scripts. While a plugin runs, the `AM_PLUGINS_RUNNING`
//! guard is exported so a nested `am` process skips plugins (no recursion).
//!
//! Layout of the directory:
//! ```text
//! state_dir/plugins/
//! notify.sh # the plugin itself (any executable, or a marker when
//! # the manifest declares a `command`)
//! notify.yaml # optional sidecar manifest: events, timeout_secs, command
//! ```
use crate::app::App;
use crate::config::PluginSettings;
use crate::events::{Event, EventKind};
use serde::Serialize;
use std::collections::BTreeMap;
use std::io::Write;
use std::path::{Path, PathBuf};
use std::process::{Command, Stdio};
use std::time::{Duration, Instant};
use wait_timeout::ChildExt;
/// Default timeout for one plugin run (overridable in settings or manifest).
pub const DEFAULT_TIMEOUT_SECS: u64 = 10;
/// Name of the plugins directory inside the state directory.
pub const DIR_NAME: &str = "plugins";
/// Recursion guard: while a plugin runs, nested am processes skip plugins.
pub const RECURSION_GUARD: &str = "AM_PLUGINS_RUNNING";
/// Events plugins can subscribe to (lifecycle, same names as the axe-6 hooks).
pub const EVENT_NAMES: &[&str] = &["on_install", "on_start", "on_stop", "on_update"];
/// Map a journal event kind to the plugin event name it fires.
fn event_name(kind: EventKind) -> Option<&'static str> {
match kind {
EventKind::Install => Some("on_install"),
EventKind::Start => Some("on_start"),
EventKind::Stop => Some("on_stop"),
EventKind::Update => Some("on_update"),
_ => None,
}
}
/// One discovered plugin.
#[derive(Debug, Clone)]
pub struct Plugin {
pub name: String,
/// Script file (also the discovery anchor when a manifest is present).
pub file: PathBuf,
/// Optional command prefix from the manifest (e.g. "python notify.py");
/// when absent the script file itself is executed.
pub command: Option<String>,
/// Events the plugin subscribes to (sorted, deduplicated).
pub events: Vec<String>,
/// Maximum runtime before the plugin is killed.
pub timeout: Duration,
}
/// Sidecar manifest of one plugin (`<name>.yaml` next to the script).
#[derive(Debug, Clone, Default, serde::Deserialize)]
#[serde(deny_unknown_fields)]
struct Manifest {
#[serde(default)]
command: Option<String>,
#[serde(default)]
events: Option<Vec<String>>,
#[serde(default)]
timeout_secs: Option<u64>,
}
/// The plugins directory of the state directory.
pub fn dir(app: &App) -> PathBuf {
app.events_dir().join(DIR_NAME)
}
fn is_manifest_ext(path: &Path) -> bool {
matches!(
path.extension().and_then(|e| e.to_str()),
Some("yaml" | "yml" | "json")
)
}
/// Discover the enabled plugins of a directory, sorted by name. Manifest
/// files (yaml/yml/json) are consumed as sidecars of their script and never
/// listed as plugins themselves.
pub fn discover(app: &App, dir: &Path) -> Vec<Plugin> {
let enabled_filter: Option<Vec<String>> = app
.config
.settings
.plugins
.as_ref()
.and_then(|p: &PluginSettings| p.enabled.clone())
.filter(|l| !l.is_empty());
let default_timeout = app
.config
.settings
.plugins
.as_ref()
.and_then(|p: &PluginSettings| p.timeout_secs)
.unwrap_or(DEFAULT_TIMEOUT_SECS);
let mut entries: Vec<PathBuf> = std::fs::read_dir(dir)
.map(|rd| {
rd.flatten()
.map(|e| e.path())
.filter(|p| {
let name = p.file_name().and_then(|n| n.to_str()).unwrap_or("");
!name.starts_with('.')
})
.collect()
})
.unwrap_or_default();
entries.sort();
let mut out = Vec::new();
for path in entries {
if !path.is_file() || is_manifest_ext(&path) {
continue;
}
let Some(stem) = path.file_stem().and_then(|s| s.to_str()) else {
continue;
};
let stem = stem.to_string();
let manifest = load_manifest(&path, &stem);
let mut events = match &manifest.events {
Some(list) => list
.iter()
.filter(|e| EVENT_NAMES.contains(&e.as_str()))
.cloned()
.collect::<Vec<String>>(),
None => EVENT_NAMES.iter().map(|s| s.to_string()).collect(),
};
events.sort();
events.dedup();
if events.is_empty() {
app.log.warn(&format!(
"plugin '{stem}': no valid event in its manifest (valid: {}) — ignored",
EVENT_NAMES.join(", ")
));
continue;
}
if let Some(filter) = &enabled_filter {
if !filter.iter().any(|n| n == &stem) {
continue;
}
}
let timeout = manifest.timeout_secs.unwrap_or(default_timeout).max(1);
out.push(Plugin {
name: stem,
file: path,
command: manifest.command,
events,
timeout: Duration::from_secs(timeout),
});
}
out
}
/// Load the sidecar manifest `<stem>.yaml|yml|json` next to the script.
/// An unreadable or invalid manifest silently yields the defaults.
fn load_manifest(script: &Path, stem: &str) -> Manifest {
let parent = script.parent().unwrap_or(Path::new("."));
for ext in ["yaml", "yml", "json"] {
let m = parent.join(format!("{stem}.{ext}"));
if let Ok(text) = std::fs::read_to_string(&m) {
if let Ok(manifest) = serde_yaml::from_str::<Manifest>(&text) {
return manifest;
}
}
}
Manifest::default()
}
/// A slice is empty (serde helper for `skip_serializing_if` on `&[T]`).
fn slice_is_empty<T>(s: &[T]) -> bool {
s.is_empty()
}
/// Payload sent to a plugin on stdin (mirrors the journal event fields).
#[derive(Debug, Serialize)]
pub struct PluginInput<'a> {
/// Plugin event name, e.g. "on_start".
pub event: &'a str,
/// Journal kind, e.g. "start".
pub kind: &'a str,
pub ts: &'a str,
#[serde(skip_serializing_if = "Option::is_none")]
pub agent: Option<&'a str>,
#[serde(skip_serializing_if = "Option::is_none")]
pub pid: Option<u32>,
#[serde(skip_serializing_if = "Option::is_none")]
pub cwd: Option<&'a str>,
#[serde(skip_serializing_if = "Option::is_none")]
pub project: Option<&'a str>,
#[serde(skip_serializing_if = "Option::is_none")]
pub session: Option<&'a str>,
#[serde(skip_serializing_if = "Option::is_none")]
pub reason: Option<&'a str>,
#[serde(skip_serializing_if = "slice_is_empty")]
pub args: &'a [String],
#[serde(skip_serializing_if = "slice_is_empty")]
pub env_keys: &'a [String],
#[serde(skip_serializing_if = "Option::is_none")]
pub exit_code: Option<i32>,
#[serde(skip_serializing_if = "Option::is_none")]
pub duration_s: Option<u64>,
#[serde(skip_serializing_if = "Option::is_none")]
pub tokens_in: Option<u64>,
#[serde(skip_serializing_if = "Option::is_none")]
pub tokens_out: Option<u64>,
#[serde(skip_serializing_if = "Option::is_none")]
pub cost_usd: Option<f64>,
}
impl<'a> PluginInput<'a> {
fn from_event(event_name: &'a str, ev: &'a Event) -> Self {
Self {
event: event_name,
kind: ev.kind.as_str(),
ts: &ev.ts,
agent: ev.agent.as_deref(),
pid: ev.pid,
cwd: ev.cwd.as_deref(),
project: ev.project.as_deref(),
session: ev.session.as_deref(),
reason: ev.reason.as_deref(),
args: &ev.args,
env_keys: &ev.env_keys,
exit_code: ev.exit_code,
duration_s: ev.duration_s,
tokens_in: ev.tokens_in,
tokens_out: ev.tokens_out,
cost_usd: ev.cost_usd,
}
}
}
/// The JSON response contract a plugin answers with on stdout.
#[derive(Debug, Clone, Default, serde::Deserialize, serde::Serialize)]
pub struct PluginResponse {
#[serde(default)]
pub ok: Option<bool>,
#[serde(default)]
pub message: Option<String>,
#[serde(default)]
pub error: Option<String>,
#[serde(default)]
pub data: Option<serde_json::Value>,
}
/// Outcome of one plugin run.
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum PluginStatus {
Ok,
Failed,
Timeout,
}
impl PluginStatus {
pub fn as_str(&self) -> &'static str {
match self {
PluginStatus::Ok => "ok",
PluginStatus::Failed => "failed",
PluginStatus::Timeout => "timeout",
}
}
}
/// Full result of one plugin run (used by `am plugins --test` and the logs).
#[derive(Debug)]
pub struct PluginResult {
pub status: PluginStatus,
pub exit_code: Option<i32>,
pub response: Option<PluginResponse>,
pub stdout: String,
pub stderr: String,
/// Set when the plugin could not be spawned or the wait failed.
pub error: Option<String>,
pub duration_ms: u64,
}
impl PluginResult {
fn error(msg: impl Into<String>) -> Self {
Self {
status: PluginStatus::Failed,
exit_code: None,
response: None,
stdout: String::new(),
stderr: String::new(),
error: Some(msg.into()),
duration_ms: 0,
}
}
}
/// Fire the plugins subscribed to this event (called from `App::emit`).
/// Non-blocking: every failure is logged as a warning, never fatal.
pub fn dispatch(app: &App, event: &Event) {
let Some(name) = event_name(event.kind) else {
return;
};
if std::env::var_os(RECURSION_GUARD).is_some() {
return; // a plugin is already running: never recurse
}
if app.dry_run() {
return;
}
let dir = dir(app);
if !dir.is_dir() {
return;
}
for plugin in discover(app, &dir) {
if !plugin.events.iter().any(|e| e == name) {
continue;
}
let result = run_plugin(app, &plugin, name, event);
report(app, &plugin, name, &result);
}
}
/// Run one plugin against one event with the JSON contract. Never fails:
/// every outcome is returned and the caller decides how to report it.
pub fn run_plugin(app: &App, plugin: &Plugin, event_name: &str, event: &Event) -> PluginResult {
let started = Instant::now();
let input = PluginInput::from_event(event_name, event);
let input_json = match serde_json::to_string(&input) {
Ok(j) => j,
Err(e) => {
return PluginResult::error(format!("cannot serialize the event payload: {e:#}"));
}
};
let (prog, args) = match &plugin.command {
Some(cmd) => {
let tokens = crate::playbook::split_words(cmd);
if tokens.is_empty() {
return PluginResult::error("empty command in the plugin manifest");
}
let (p, pre) = crate::runner::resolve_program(&tokens[0]);
let mut a = pre;
a.extend(tokens[1..].iter().cloned());
(p, a)
}
None => {
let (p, pre) = crate::runner::resolve_program(&plugin.file.display().to_string());
(p, pre)
}
};
let mut env = BTreeMap::new();
for (k, v) in std::env::vars_os() {
env.insert(
k.to_string_lossy().to_string(),
v.to_string_lossy().to_string(),
);
}
env.insert(RECURSION_GUARD.to_string(), "1".to_string());
let mut cmd = Command::new(&prog);
cmd.args(&args)
.envs(&env)
.stdin(Stdio::piped())
.stdout(Stdio::piped())
.stderr(Stdio::piped());
// The plugin runs in its own directory: relative paths in the manifest
// command (e.g. "python3 notify.py") resolve next to the script.
if let Some(parent) = plugin.file.parent() {
cmd.current_dir(parent);
}
set_process_group(&mut cmd);
app.log
.verbose(&format!("plugin '{}': {} {event_name}", plugin.name, prog));
let mut child = match cmd.spawn() {
Ok(c) => c,
Err(e) => {
return PluginResult::error(format!(
"cannot run {}: {e:#}",
plugin.file.display()
));
}
};
// Feed the event JSON, then close stdin (EOF).
if let Some(mut stdin) = child.stdin.take() {
let _ = stdin.write_all(input_json.as_bytes());
}
let stdout = read_pipe(child.stdout.take());
let stderr = read_pipe(child.stderr.take());
let deadline = started + plugin.timeout;
let mut timed_out = false;
let mut exit_code: Option<i32> = None;
loop {
match child.wait_timeout(Duration::from_millis(200)) {
Ok(Some(status)) => {
exit_code = status.code();
break;
}
Ok(None) => {
if Instant::now() >= deadline {
crate::process::kill_tree(child.id());
let _ = child.wait();
timed_out = true;
break;
}
}
Err(e) => {
crate::process::kill_tree(child.id());
let _ = child.wait();
return PluginResult::error(format!("cannot wait for the plugin: {e:#}"));
}
}
}
let out = stdout.join().unwrap_or_default();
let err_out = stderr.join().unwrap_or_default();
let duration_ms = started.elapsed().as_millis() as u64;
if timed_out {
return PluginResult {
status: PluginStatus::Timeout,
exit_code: None,
response: None,
stdout: out,
stderr: err_out,
error: Some(format!(
"timeout after {}s — killed",
plugin.timeout.as_secs()
)),
duration_ms,
};
}
let code = exit_code.unwrap_or(-1);
let response = parse_response(&out);
PluginResult {
status: if code == 0 {
PluginStatus::Ok
} else {
PluginStatus::Failed
},
exit_code: Some(code),
response,
stdout: out,
stderr: err_out,
error: None,
duration_ms,
}
}
/// Log the outcome of a plugin run as warnings (never fatal).
fn report(app: &App, plugin: &Plugin, event_name: &str, r: &PluginResult) {
let who = format!("plugin '{}' on {event_name}", plugin.name);
match r.status {
PluginStatus::Ok => {
let resp = r.response.as_ref();
let message = resp.and_then(|x| x.message.clone());
let error = resp.and_then(|x| x.error.clone());
let ok_flag = resp.and_then(|x| x.ok).unwrap_or(true);
if let Some(err) = error {
app.log.warn(&format!("{who}: {err}"));
} else if !ok_flag {
app.log.warn(&format!(
"{who}: returned ok=false{}",
message
.map(|m| format!(" ({m})"))
.unwrap_or_default()
));
} else if let Some(m) = message {
app.log
.verbose(&format!("{who}: ok in {}ms — {m}", r.duration_ms));
} else {
app.log
.verbose(&format!("{who}: ok in {}ms", r.duration_ms));
}
}
PluginStatus::Failed => {
let detail = r
.response
.as_ref()
.and_then(|x| x.error.clone())
.or_else(|| {
let t = r.stderr.trim();
if t.is_empty() {
None
} else {
Some(t.to_string())
}
})
.unwrap_or_default();
app.log.warn(&format!(
"{who} failed (exit {}){}",
r.exit_code.unwrap_or(-1),
if detail.is_empty() {
String::new()
} else {
format!(": {detail}")
}
));
}
PluginStatus::Timeout => {
app.log.warn(&format!(
"{who}: {}",
r.error.as_deref().unwrap_or("timed out")
));
}
}
if let Some(e) = &r.error {
if r.status != PluginStatus::Timeout {
app.log.warn(&format!("{who}: {e}"));
}
}
if r.response.is_none() && !r.stdout.trim().is_empty() {
app.log.warn(&format!(
"{who}: invalid plugin output (expected a JSON object on stdout): {}",
crate::output::truncate(r.stdout.trim(), 80)
));
}
}
/// Parse the plugin stdout as the response contract (a JSON object).
fn parse_response(out: &str) -> Option<PluginResponse> {
let trimmed = out.trim();
if trimmed.is_empty() {
return None;
}
let value: serde_json::Value = serde_json::from_str(trimmed).ok()?;
if !value.is_object() {
return None;
}
serde_json::from_value(value).ok()
}
/// Read a child pipe to the end on a dedicated thread (avoids pipe deadlocks
/// while the parent waits with a timeout).
fn read_pipe<R: std::io::Read + Send + 'static>(pipe: Option<R>) -> std::thread::JoinHandle<String> {
std::thread::spawn(move || {
use std::io::Read;
let mut buf = String::new();
if let Some(mut p) = pipe {
let _ = p.read_to_string(&mut buf);
}
buf
})
}
/// Put the child in its own process group so `kill_tree` can stop it with
/// its whole subtree.
fn set_process_group(cmd: &mut Command) {
#[cfg(unix)]
{
use std::os::unix::process::CommandExt;
cmd.process_group(0);
}
#[cfg(windows)]
{
use std::os::windows::process::CommandExt;
const CREATE_NEW_PROCESS_GROUP: u32 = 0x0000_0200;
cmd.creation_flags(CREATE_NEW_PROCESS_GROUP);
}
}
// ---------------------------------------------------------------------------
// Tests
// ---------------------------------------------------------------------------
#[cfg(test)]
mod tests {
use super::*;
use crate::cli::Cli;
use clap::Parser;
fn app_with_plugins(tag: &str) -> (App, PathBuf) {
let dir = tempfile::tempdir().unwrap();
let root = dir.path().to_path_buf();
std::mem::forget(dir);
let cfg = root.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 = Cli::parse_from(["am", "--config", cfg.to_str().unwrap(), "list"]);
let mut app = App::from_cli(cli).unwrap();
app.paths.state_file = root.join(format!("state-{tag}.json"));
app.paths.log_dir = root.join("logs");
let pdir = app.events_dir().join(DIR_NAME);
std::fs::create_dir_all(&pdir).unwrap();
(app, pdir)
}
/// Write a script plus a manifest running it through the system shell
/// with an absolute path (a bare .sh file cannot be spawned directly on
/// Windows; the `command` override is the portable way).
fn write_shell_plugin(pdir: &Path, name: &str, script: &str, manifest_extra: &str) {
let script_path = pdir.join(format!("{name}.sh"));
std::fs::write(&script_path, script).unwrap();
let abs = script_path.display().to_string().replace('\\', "/");
std::fs::write(
pdir.join(format!("{name}.yaml")),
format!("command: \"sh {abs}\"\n{manifest_extra}"),
)
.unwrap();
}
#[test]
fn manifest_events_and_timeout_are_parsed() {
let (app, pdir) = app_with_plugins("man");
std::fs::write(pdir.join("notify.sh"), "#!/bin/sh\ncat\n").unwrap();
std::fs::write(
pdir.join("notify.yaml"),
"events: [on_start, on_stop]\ntimeout_secs: 3\n",
)
.unwrap();
let plugins = discover(&app, &pdir);
assert_eq!(plugins.len(), 1);
assert_eq!(plugins[0].name, "notify");
assert_eq!(
plugins[0].events,
vec!["on_start".to_string(), "on_stop".to_string()]
);
assert_eq!(plugins[0].timeout.as_secs(), 3);
// The manifest itself is never a plugin.
assert!(!plugins.iter().any(|p| p.name == "notify.yaml"));
}
#[test]
fn defaults_are_all_events_and_ten_seconds() {
let (app, pdir) = app_with_plugins("def");
std::fs::write(pdir.join("all.sh"), "#!/bin/sh\ncat\n").unwrap();
let plugins = discover(&app, &pdir);
assert_eq!(plugins.len(), 1);
assert_eq!(
plugins[0].events,
EVENT_NAMES.iter().map(|s| s.to_string()).collect::<Vec<_>>()
);
assert_eq!(plugins[0].timeout.as_secs(), DEFAULT_TIMEOUT_SECS);
}
#[test]
fn enabled_filter_and_settings_timeout_apply() {
let (mut app, pdir) = app_with_plugins("en");
std::fs::write(pdir.join("a.sh"), "#!/bin/sh\ncat\n").unwrap();
std::fs::write(pdir.join("b.sh"), "#!/bin/sh\ncat\n").unwrap();
app.config.settings.plugins = Some(PluginSettings {
timeout_secs: Some(7),
enabled: Some(vec!["b".to_string()]),
});
let plugins = discover(&app, &pdir);
assert_eq!(plugins.len(), 1);
assert_eq!(plugins[0].name, "b");
assert_eq!(plugins[0].timeout.as_secs(), 7);
}
#[test]
fn unknown_events_are_filtered_out() {
let (app, pdir) = app_with_plugins("unk");
std::fs::write(pdir.join("x.sh"), "#!/bin/sh\ncat\n").unwrap();
std::fs::write(pdir.join("x.yaml"), "events: [on_boom, on_start]\n").unwrap();
let plugins = discover(&app, &pdir);
assert_eq!(plugins.len(), 1);
assert_eq!(plugins[0].events, vec!["on_start".to_string()]);
}
#[test]
fn plugin_runs_with_the_json_contract() {
let (app, pdir) = app_with_plugins("run");
write_shell_plugin(
&pdir,
"echo",
"#!/bin/sh\ncat >/dev/null\necho '{\"ok\": true, \"message\": \"pong\"}'\n",
"",
);
let plugins = discover(&app, &pdir);
assert_eq!(plugins.len(), 1);
let ev = Event::now(EventKind::Start)
.with_agent("claude-code")
.with_pid(42);
let r = run_plugin(&app, &plugins[0], "on_start", &ev);
assert_eq!(r.status, PluginStatus::Ok, "stderr: {}", r.stderr);
assert_eq!(r.exit_code, Some(0));
let resp = r.response.expect("a JSON response");
assert_eq!(resp.ok, Some(true));
assert_eq!(resp.message.as_deref(), Some("pong"));
}
#[test]
fn failing_plugin_is_reported_without_panicking() {
let (app, pdir) = app_with_plugins("fail");
write_shell_plugin(
&pdir,
"fail",
"#!/bin/sh\necho '{\"ok\": false, \"error\": \"boom\"}' >&2\nexit 3\n",
"",
);
let plugins = discover(&app, &pdir);
let ev = Event::now(EventKind::Start);
let r = run_plugin(&app, &plugins[0], "on_start", &ev);
assert_eq!(r.status, PluginStatus::Failed);
assert_eq!(r.exit_code, Some(3));
// dispatch only warns — it never returns an error or panics.
dispatch(&app, &ev);
}
#[test]
fn invalid_output_is_signaled_without_breaking_am() {
let (app, pdir) = app_with_plugins("bad");
write_shell_plugin(
&pdir,
"bad",
"#!/bin/sh\ncat >/dev/null\necho 'pas du json'\n",
"",
);
let plugins = discover(&app, &pdir);
let ev = Event::now(EventKind::Start);
let r = run_plugin(&app, &plugins[0], "on_start", &ev);
assert_eq!(r.status, PluginStatus::Ok); // exit 0: the run itself is fine
assert!(r.response.is_none(), "non-JSON output has no response");
dispatch(&app, &ev); // warning logged, nothing panics
}
#[test]
fn timeout_is_enforced_and_kills_the_plugin() {
let (app, pdir) = app_with_plugins("slow");
write_shell_plugin(
&pdir,
"slow",
"#!/bin/sh\nwhile true; do :; done\n",
"timeout_secs: 1\n",
);
let plugins = discover(&app, &pdir);
assert_eq!(plugins[0].timeout.as_secs(), 1);
let ev = Event::now(EventKind::Start);
let r = run_plugin(&app, &plugins[0], "on_start", &ev);
assert_eq!(r.status, PluginStatus::Timeout);
assert!(r.error.as_deref().unwrap_or("").contains("timeout"));
}
#[test]
fn command_override_runs_through_the_shell() {
let (app, pdir) = app_with_plugins("cmd");
// The manifest command wins over the file: the file is only a
// marker, the command does the work (portable on every OS).
std::fs::write(pdir.join("echo.sh"), "this is not a script\n").unwrap();
let worker = pdir.parent().unwrap().join("worker.sh");
std::fs::write(
&worker,
"#!/bin/sh\ncat >/dev/null\necho '{\"ok\": true, \"message\": \"via-command\"}'\n",
)
.unwrap();
let abs = worker.display().to_string().replace('\\', "/");
std::fs::write(pdir.join("echo.yaml"), format!("command: \"sh {abs}\"\n")).unwrap();
let plugins = discover(&app, &pdir);
let ev = Event::now(EventKind::Update).with_agent("aider");
let r = run_plugin(&app, &plugins[0], "on_update", &ev);
assert_eq!(r.status, PluginStatus::Ok, "stderr: {}", r.stderr);
let resp = r.response.expect("a JSON response");
assert_eq!(resp.message.as_deref(), Some("via-command"));
}
#[test]
fn dispatch_skips_other_kinds_and_missing_dir() {
let (app, pdir) = app_with_plugins("skip");
// No plugins dir: dispatch is a no-op.
std::fs::remove_dir_all(&pdir).unwrap();
dispatch(&app, &Event::now(EventKind::Start));
// A kind with no plugin event (run) is a no-op too.
dispatch(&app, &Event::now(EventKind::Run));
}
#[test]
fn recursion_guard_skips_dispatch() {
let (app, pdir) = app_with_plugins("guard");
write_shell_plugin(
&pdir,
"echo",
"#!/bin/sh\ncat >/dev/null\necho '{\"ok\": true}'\n",
"",
);
std::env::set_var(RECURSION_GUARD, "1");
dispatch(&app, &Event::now(EventKind::Start)); // must not run anything
std::env::remove_var(RECURSION_GUARD);
dispatch(&app, &Event::now(EventKind::Start)); // normal path: no panic
}
}
+21
View File
@@ -90,6 +90,27 @@ pub fn stop_pid(pid: u32, force: bool, timeout_secs: u64) -> bool {
!is_running(pid)
}
/// Kill a process and its whole tree (process group on unix, taskkill /T on
/// Windows). Used by the plugin runner and the lab benchmark to stop
/// runaway scripts and timed-out task runs.
pub fn kill_tree(pid: u32) {
#[cfg(unix)]
{
use std::os::unix::process::CommandExt;
let _ = Command::new("kill")
.arg("-KILL")
.arg(format!("-{pid}")) // negative pid = the process group
.status();
let _ = Command::new("kill").arg("-KILL").arg(pid.to_string()).status();
}
#[cfg(windows)]
{
let _ = Command::new("taskkill")
.args(["/PID", &pid.to_string(), "/T", "/F"])
.status();
}
}
/// Start a program detached: stdout/stderr are redirected to the agent log
/// file. Returns the PID. Fails when the process exits immediately.
pub fn spawn_background(
+18 -1
View File
@@ -116,6 +116,8 @@ const COMMAND_DESCRIPTIONS: &[(&str, &str)] = &[
("sync", "push the state (journal, sessions, config) into a git repo"),
("migrate", "export/import a machine transfer bundle (config + state)"),
("playbook", "replay a saved history sequence step by step (am history --save)"),
("lab", "benchmark the same task on several agents (duration, cost, --json)"),
("plugins", "list the event plugins and test one (JSON contract on stdin/stdout)"),
("shell", "show or switch the system shell"),
("theme", "show or switch the color theme"),
("tip", "cheat sheet of the most useful commands"),
@@ -230,7 +232,7 @@ impl AmCompleter {
"self-update", "self-uninstall", "export", "import", "shell", "theme",
"tip", "dashboard", "favorite", "unfavorite", "note", "tag", "untag", "tags",
"profile", "man", "models", "catalog", "suggest", "audit",
"service", "schedule", "monitor", "sync", "migrate", "playbook",
"service", "schedule", "monitor", "sync", "migrate", "playbook", "lab", "plugins",
"ls", "dir", "cd", "ps", "where", "get", "help", "version", "exit",
],
config_sub: vec!["show", "path", "edit", "validate", "add"],
@@ -782,6 +784,9 @@ pub fn banner_box(
rows.push(inner(
" automate service install · schedule add · doctor --watch · monitor · sync · migrate · playbook".to_string(),
));
rows.push(inner(
" bench lab --agents a,b --task t.yaml · plugins · plugins --test <name>".to_string(),
));
rows.push(inner(
" system self-update · self-uninstall · export · import".to_string(),
));
@@ -1617,6 +1622,16 @@ fn handle_line(
.collect(),
}
},
"lab" => Command::Lab(crate::cli::LabArgs {
agents: opt_value("--agents"),
task: opt_value("--task").map(std::path::PathBuf::from),
parallel: flag("--parallel"),
timeout: opt_value("--timeout").and_then(|v| v.parse().ok()),
list: flag("--list"),
}),
"plugins" => Command::Plugins(crate::cli::PluginsArgs {
test: opt_value("--test"),
}),
"init" => Command::Init {
force: flag("--force"),
template: opt_value("--template"),
@@ -1864,6 +1879,8 @@ pub fn is_am_command(word: &str) -> bool {
| "sync"
| "migrate"
| "playbook"
| "lab"
| "plugins"
)
}
+169
View File
@@ -0,0 +1,169 @@
//! Integration tests: am lab (issue #62) and am plugins (issue #75).
mod common;
use agent_manager::cli::{Cli, LabArgs, PluginsArgs};
use agent_manager::commands::{lab_cmd, plugins_cmd};
use clap::Parser;
use std::path::PathBuf;
/// A config with two shell-based fake agents for the lab benchmark.
fn lab_config() -> String {
concat!(
"agents:\n",
" - name: bench-ok\n",
" installable: false\n",
" run: sh\n",
" args: [\"-c\", \"echo out: $0\"]\n",
" - name: bench-fail\n",
" installable: false\n",
" run: sh\n",
" args: [\"-c\", \"echo boom >&2; exit 3\"]\n",
)
.to_string()
}
#[test]
fn lab_runs_a_task_on_every_agent_and_reports() {
let dir = std::env::temp_dir().join(format!("am-lab-e2e-{}", std::process::id()));
std::fs::create_dir_all(&dir).unwrap();
let cfg = common::write_config(&dir, &lab_config());
let task = dir.join("hello.yaml");
std::fs::write(
&task,
"name: hello\ndescription: e2e\nargs: [\"hello\"]\ntimeout: 30\n",
)
.unwrap();
let cli = Cli::parse_from([
"am",
"--config",
cfg.to_str().unwrap(),
"lab",
"--agents",
"bench-ok",
"--task",
"hello.yaml",
]);
let mut app = agent_manager::app::App::from_cli(cli).expect("app should build");
common::isolate(&mut app, &dir);
let code = lab_cmd::run(
&app,
&LabArgs {
agents: Some("bench-ok, bench-fail".to_string()),
task: Some(task.clone()),
parallel: false,
timeout: None,
list: false,
},
)
.unwrap();
// One agent fails → the exit code is 1 (CI-friendly).
assert_eq!(code, 1);
// The benchmark is replayable with the same file.
let code2 = lab_cmd::run(
&app,
&LabArgs {
agents: Some("bench-ok".to_string()),
task: Some(task),
parallel: false,
timeout: None,
list: false,
},
)
.unwrap();
assert_eq!(code2, 0);
// Each run was journaled as a lab event.
let events = agent_manager::events::read_events(&app.events_dir(), 0);
let lab_events = events.iter().filter(|e| e.kind == agent_manager::events::EventKind::Lab);
assert!(lab_events.count() >= 3, "expected >= 3 lab events");
let _ = std::fs::remove_dir_all(&dir);
}
#[test]
fn lab_list_shows_the_tasks_of_the_lab_dir() {
let dir = std::env::temp_dir().join(format!("am-lab-list-{}", std::process::id()));
std::fs::create_dir_all(&dir).unwrap();
let cfg = common::write_config(&dir, "agents: []\n");
let cli = Cli::parse_from(["am", "--config", cfg.to_str().unwrap(), "lab", "--list"]);
let mut app = agent_manager::app::App::from_cli(cli).expect("app should build");
common::isolate(&mut app, &dir);
std::fs::create_dir_all(app.events_dir().join("lab")).unwrap();
std::fs::write(app.events_dir().join("lab/hello.yaml"), "name: hello\n").unwrap();
let code = lab_cmd::run(
&app,
&LabArgs {
agents: None,
task: None,
parallel: false,
timeout: None,
list: true,
},
)
.unwrap();
assert_eq!(code, 0);
let _ = std::fs::remove_dir_all(&dir);
}
#[test]
fn plugin_runs_on_an_emitted_event() {
let dir = std::env::temp_dir().join(format!("am-plugins-e2e-{}", std::process::id()));
std::fs::create_dir_all(&dir).unwrap();
let cfg = common::write_config(&dir, "agents: []\n");
let cli = Cli::parse_from(["am", "--config", cfg.to_str().unwrap(), "plugins"]);
let mut app = agent_manager::app::App::from_cli(cli).expect("app should build");
common::isolate(&mut app, &dir);
// Install an echo plugin: valid JSON on stdout (contract) + a marker
// file proving it ran. The absolute path keeps it independent of cwd.
let pdir = app.events_dir().join("plugins");
std::fs::create_dir_all(&pdir).unwrap();
let marker = pdir.join("pong.log");
let script = format!(
"#!/bin/sh\ncat >/dev/null\necho '{{\"ok\": true, \"message\": \"pong\"}}' | tee '{}'\n",
marker.display().to_string().replace('\\', "/")
);
std::fs::write(pdir.join("pong.sh"), script).unwrap();
let script_path = pdir.join("pong.sh").display().to_string().replace('\\', "/");
std::fs::write(pdir.join("pong.yaml"), format!("command: \"sh {script_path}\"\n")).unwrap();
// List: the plugin is discovered.
let code = plugins_cmd::run(&app, &PluginsArgs { test: None }).unwrap();
assert_eq!(code, 0);
// --test: the JSON contract is exercised.
let code = plugins_cmd::run(&app, &PluginsArgs { test: Some("pong".to_string()) }).unwrap();
assert_eq!(code, 0);
// A real emitted event fires the plugin (acceptance criterion).
app.emit(&agent_manager::events::Event::now(
agent_manager::events::EventKind::Start,
));
assert!(marker.exists(), "the plugin must have run on the emitted event");
let _ = std::fs::remove_dir_all(&dir);
}
#[test]
fn plugins_list_and_test_are_available_in_cli() {
let cli = Cli::parse_from(["am", "plugins", "--test", "notify"]);
match cli.command {
Some(agent_manager::cli::Command::Plugins(args)) => {
assert_eq!(args.test.as_deref(), Some("notify"));
}
other => panic!("unexpected: {other:?}"),
}
let cli = Cli::parse_from(["am", "lab", "--agents", "a,b", "--task", "t.yaml"]);
match cli.command {
Some(agent_manager::cli::Command::Lab(args)) => {
assert_eq!(args.agents.as_deref(), Some("a,b"));
assert_eq!(args.task, Some(PathBuf::from("t.yaml")));
}
other => panic!("unexpected: {other:?}"),
}
}