From 4c7866a25ad993b8de271df20916be9d766028fa Mon Sep 17 00:00:00 2001 From: Bruno Charest Date: Tue, 18 Aug 2026 14:17:30 -0400 Subject: [PATCH] =?UTF-8?q?feat:=20phase=202=20=E2=80=94=20am=20lab=20(ben?= =?UTF-8?q?chmark=20N=20agents)=20+=20plugins=20d'=C3=A9v=C3=A9nements=20(?= =?UTF-8?q?contrat=20JSON)?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit - #62: am lab --agents a,b --task — séquentiel ou --parallel, capture durée/exit/sortie/coût estimé, rapport comparatif tableau + --json stable, tâches versionnables dans /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 (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 --- Cargo.lock | 2 +- Cargo.toml | 2 +- README.md | 2 + ROADMAP.md | 4 +- config.yaml | 6 + examples/lab/hello.yaml | 8 + examples/plugins/README.md | 31 ++ examples/plugins/ci-webhook.py | 38 ++ examples/plugins/ci-webhook.yaml | 4 + examples/plugins/notify.py | 24 + examples/plugins/notify.yaml | 7 + man/am-lab.1 | 28 ++ man/am-plugins.1 | 16 + man/am.1 | 10 +- src/app.rs | 4 +- src/cli.rs | 32 ++ src/commands/lab_cmd.rs | 101 ++++ src/commands/mod.rs | 4 + src/commands/plugins_cmd.rs | 102 ++++ src/commands/tip_cmd.rs | 12 + src/config.rs | 19 + src/events.rs | 3 + src/help.rs | 38 ++ src/lab.rs | 629 ++++++++++++++++++++++++ src/lib.rs | 2 + src/plugins.rs | 794 +++++++++++++++++++++++++++++++ src/process.rs | 21 + src/repl.rs | 19 +- tests/lab_plugins_test.rs | 169 +++++++ 29 files changed, 2123 insertions(+), 8 deletions(-) create mode 100644 examples/lab/hello.yaml create mode 100644 examples/plugins/README.md create mode 100644 examples/plugins/ci-webhook.py create mode 100644 examples/plugins/ci-webhook.yaml create mode 100644 examples/plugins/notify.py create mode 100644 examples/plugins/notify.yaml create mode 100644 man/am-lab.1 create mode 100644 man/am-plugins.1 create mode 100644 src/commands/lab_cmd.rs create mode 100644 src/commands/plugins_cmd.rs create mode 100644 src/lab.rs create mode 100644 src/plugins.rs create mode 100644 tests/lab_plugins_test.rs diff --git a/Cargo.lock b/Cargo.lock index 4f518dc..f5fa14d 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -21,7 +21,7 @@ dependencies = [ [[package]] name = "agent-manager" -version = "0.5.3" +version = "0.5.4" dependencies = [ "anyhow", "chrono", diff --git a/Cargo.toml b/Cargo.toml index ca76ab6..673415f 100644 --- a/Cargo.toml +++ b/Cargo.toml @@ -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" diff --git a/README.md b/README.md index 31e8a6d..958c7d6 100644 --- a/README.md +++ b/README.md @@ -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 · am playbook | exporte une plage en playbook YAML et la rejoue pas à pas ({{var}}) | +| am lab --agents a,b --task [--parallel] [--json] | benchmark : même tâche sur plusieurs agents (durée, exit, coût) — tâches versionnables dans /lab/ | +| am plugins [--test ] | 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 --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) | diff --git a/ROADMAP.md b/ROADMAP.md index 6311127..682e969 100644 --- a/ROADMAP.md +++ b/ROADMAP.md @@ -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 | 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). diff --git a/config.yaml b/config.yaml index 22bc63c..f42bd72 100644 --- a/config.yaml +++ b/config.yaml @@ -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 /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 diff --git a/examples/lab/hello.yaml b/examples/lab/hello.yaml new file mode 100644 index 0000000..96f9333 --- /dev/null +++ b/examples/lab/hello.yaml @@ -0,0 +1,8 @@ +# Tâche de benchmark (issue #62) : les mêmes arguments sont passés à la +# commande de chaque agent. Copier dans /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 diff --git a/examples/plugins/README.md b/examples/plugins/README.md new file mode 100644 index 0000000..1e4040e --- /dev/null +++ b/examples/plugins/README.md @@ -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 diff --git a/examples/plugins/ci-webhook.py b/examples/plugins/ci-webhook.py new file mode 100644 index 0000000..c3c36e1 --- /dev/null +++ b/examples/plugins/ci-webhook.py @@ -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()) diff --git a/examples/plugins/ci-webhook.yaml b/examples/plugins/ci-webhook.yaml new file mode 100644 index 0000000..d3b20c0 --- /dev/null +++ b/examples/plugins/ci-webhook.yaml @@ -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 diff --git a/examples/plugins/notify.py b/examples/plugins/notify.py new file mode 100644 index 0000000..53478b8 --- /dev/null +++ b/examples/plugins/notify.py @@ -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 /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()) diff --git a/examples/plugins/notify.yaml b/examples/plugins/notify.yaml new file mode 100644 index 0000000..a83d8d4 --- /dev/null +++ b/examples/plugins/notify.yaml @@ -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 diff --git a/man/am-lab.1 b/man/am-lab.1 new file mode 100644 index 0000000..a680dc2 --- /dev/null +++ b/man/am-lab.1 @@ -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\fR +Comma\-separated agent names (or aliases) to benchmark +.TP +\fB\-\-task\fR \fI\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\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 diff --git a/man/am-plugins.1 b/man/am-plugins.1 new file mode 100644 index 0000000..697572c --- /dev/null +++ b/man/am-plugins.1 @@ -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\fR +Run one plugin against a synthetic event and print its response +.TP +\fB\-h\fR, \fB\-\-help\fR +Print help diff --git a/man/am.1 b/man/am.1 index 74dd335..5d3d4fc 100644 --- a/man/am.1 +++ b/man/am.1 @@ -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 diff --git a/src/app.rs b/src/app.rs index da138c1..a36b68b 100644 --- a/src/app.rs +++ b/src/app.rs @@ -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. diff --git a/src/cli.rs b/src/cli.rs index 229ff39..61b2bbe 100644 --- a/src/cli.rs +++ b/src/cli.rs @@ -423,6 +423,10 @@ pub enum Command { #[arg(long, value_name = "NAME=VALUE")] var: Vec, }, + /// 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, } +/// 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, + /// Task file (a bare name resolves under state_dir/lab/, .yaml appended) + #[arg(long, value_name = "FILE", required_unless_present = "list")] + pub task: Option, + /// 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, + /// 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, +} + /// Remote catalog subcommands (issues #60 #67). #[derive(Subcommand, Debug, Clone)] pub enum CatalogCmd { diff --git a/src/commands/lab_cmd.rs b/src/commands/lab_cmd.rs new file mode 100644 index 0000000..6c9bfcf --- /dev/null +++ b/src/commands/lab_cmd.rs @@ -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 { + if args.list { + let tasks = lab::list_tasks(app); + if app.json() { + let names: Vec = 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 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); +} diff --git a/src/commands/mod.rs b/src/commands/mod.rs index 9498157..2164d9f 100644 --- a/src/commands/mod.rs +++ b/src/commands/mod.rs @@ -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 { 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()) diff --git a/src/commands/plugins_cmd.rs b/src/commands/plugins_cmd.rs new file mode 100644 index 0000000..2fb8575 --- /dev/null +++ b/src/commands/plugins_cmd.rs @@ -0,0 +1,102 @@ +//! plugins: list the event plugins and test one against a synthetic event +//! (issue #75). 'am plugins --test ' 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 { + 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 = 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) +} diff --git a/src/commands/tip_cmd.rs b/src/commands/tip_cmd.rs index 8a2efc6..2f47c43 100644 --- a/src/commands/tip_cmd.rs +++ b/src/commands/tip_cmd.rs @@ -260,6 +260,18 @@ pub static SECTIONS: &[TipSection] = &[ options: &[("--save ", "fichier playbook"), ("playbook ", "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 ", "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é", diff --git a/src/config.rs b/src/config.rs index ec28553..d7a5a44 100644 --- a/src/config.rs +++ b/src/config.rs @@ -136,6 +136,22 @@ pub struct Settings { /// Push automatically when the REPL exits (issue #66, opt-in). #[serde(default)] pub sync_on_exit: Option, + /// Plugin scripts (issue #75): default timeout and enable list. + #[serde(default)] + pub plugins: Option, +} + +/// 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, + /// Only run these plugin names (empty = all the plugins of the dir). + #[serde(default)] + pub enabled: Option>, } /// 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); } diff --git a/src/events.rs b/src/events.rs index 458b46d..f0fa950 100644 --- a/src/events.rs +++ b/src/events.rs @@ -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", } } } diff --git a/src/help.rs b/src/help.rs index 986660b..e12d923 100644 --- a/src/help.rs +++ b/src/help.rs @@ -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 {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", diff --git a/src/lab.rs b/src/lab.rs new file mode 100644 index 0000000..8648c4b --- /dev/null +++ b/src/lab.rs @@ -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 `/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, + /// Arguments passed to every agent's run command. + #[serde(default)] + pub args: Vec, + /// Working directory of the run (default: the current directory). + #[serde(default, skip_serializing_if = "Option::is_none")] + pub cwd: Option, + /// Timeout in seconds (default 600). + #[serde(default, skip_serializing_if = "Option::is_none")] + pub timeout: Option, +} + +/// 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, + 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, + #[serde(skip_serializing_if = "Option::is_none")] + pub tokens_out: Option, + #[serde(skip_serializing_if = "Option::is_none")] + pub cost_usd: Option, + #[serde(skip_serializing_if = "Option::is_none")] + pub error: Option, +} + +/// 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, + 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, + #[serde(skip_serializing_if = "Option::is_none")] + pub cheapest: Option, +} + +/// 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 { + let mut out: Vec = 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 { + 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 { + 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, + env: BTreeMap, + cwd: Option, + timeout: Duration, + cost_model: Option, + sheet: BTreeMap, +} + +/// 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, +) -> Result { + 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 = 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(pipe: Option) -> std::thread::JoinHandle { + 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::::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()); + } +} diff --git a/src/lib.rs b/src/lib.rs index 01beb7b..8cd304c 100644 --- a/src/lib.rs +++ b/src/lib.rs @@ -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; diff --git a/src/plugins.rs b/src/plugins.rs new file mode 100644 index 0000000..23a934f --- /dev/null +++ b/src/plugins.rs @@ -0,0 +1,794 @@ +//! Plugin scripts (issue #75): executable scripts in `/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, + /// Events the plugin subscribes to (sorted, deduplicated). + pub events: Vec, + /// Maximum runtime before the plugin is killed. + pub timeout: Duration, +} + +/// Sidecar manifest of one plugin (`.yaml` next to the script). +#[derive(Debug, Clone, Default, serde::Deserialize)] +#[serde(deny_unknown_fields)] +struct Manifest { + #[serde(default)] + command: Option, + #[serde(default)] + events: Option>, + #[serde(default)] + timeout_secs: Option, +} + +/// 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 { + let enabled_filter: Option> = 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 = 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::>(), + 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 `.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::(&text) { + return manifest; + } + } + } + Manifest::default() +} + +/// A slice is empty (serde helper for `skip_serializing_if` on `&[T]`). +fn slice_is_empty(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, + #[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, + #[serde(skip_serializing_if = "Option::is_none")] + pub duration_s: Option, + #[serde(skip_serializing_if = "Option::is_none")] + pub tokens_in: Option, + #[serde(skip_serializing_if = "Option::is_none")] + pub tokens_out: Option, + #[serde(skip_serializing_if = "Option::is_none")] + pub cost_usd: Option, +} + +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, + #[serde(default)] + pub message: Option, + #[serde(default)] + pub error: Option, + #[serde(default)] + pub data: Option, +} + +/// 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, + pub response: Option, + pub stdout: String, + pub stderr: String, + /// Set when the plugin could not be spawned or the wait failed. + pub error: Option, + pub duration_ms: u64, +} + +impl PluginResult { + fn error(msg: impl Into) -> 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 = 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 { + 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(pipe: Option) -> std::thread::JoinHandle { + 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::>() + ); + 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 + } +} diff --git a/src/process.rs b/src/process.rs index dc1c86a..e28db38 100644 --- a/src/process.rs +++ b/src/process.rs @@ -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( diff --git a/src/repl.rs b/src/repl.rs index f06ec0f..7511274 100644 --- a/src/repl.rs +++ b/src/repl.rs @@ -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 ".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" ) } diff --git a/tests/lab_plugins_test.rs b/tests/lab_plugins_test.rs new file mode 100644 index 0000000..879d7d4 --- /dev/null +++ b/tests/lab_plugins_test.rs @@ -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:?}"), + } +}