use std::path::{Path, PathBuf};
use std::time::Duration;
use async_trait::async_trait;
use serde_json::Value;
use crate::error::{Error, Result};
use crate::tools::{Tool, ToolContext};
pub const DEFAULT_PLUGIN_TOOL_TIMEOUT_SECS: u64 = 30;
pub const PLUGIN_TOOL_MAX_OUTPUT_BYTES: usize = 1024 * 1024;
#[derive(Debug, Clone, Copy, PartialEq, Eq, Default)]
pub enum TrustDecision {
#[default]
Ask,
Always,
Never,
}
impl TrustDecision {
pub fn parse(s: &str) -> Option<TrustDecision> {
match s {
"ask" => Some(TrustDecision::Ask),
"always" => Some(TrustDecision::Always),
"never" => Some(TrustDecision::Never),
_ => None,
}
}
}
pub fn is_trusted(config: &crate::Config) -> bool {
config.trust_enabled && config.trust_default == TrustDecision::Always
}
#[derive(Debug, Clone, PartialEq)]
pub struct PluginToolSpec {
pub command: String,
pub args: Vec<String>,
pub description: String,
pub params: Value,
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct PluginHookSpec {
pub event: String,
pub command: String,
}
#[derive(Debug, Clone, PartialEq)]
pub struct PluginManifest {
pub name: String,
pub version: String,
pub tools: Vec<(String, PluginToolSpec)>,
pub hooks: Vec<PluginHookSpec>,
}
#[derive(Debug, Default, serde::Deserialize)]
struct RawManifest {
name: Option<String>,
#[serde(default)]
version: String,
#[serde(default)]
tools: Vec<RawTool>,
#[serde(default)]
hooks: Vec<RawHook>,
}
#[derive(Debug, Default, serde::Deserialize)]
struct RawTool {
#[serde(default)]
name: String,
#[serde(default)]
command: String,
#[serde(default)]
args: Vec<String>,
#[serde(default)]
description: String,
params: Option<Value>,
}
#[derive(Debug, Default, serde::Deserialize)]
struct RawHook {
#[serde(default)]
event: String,
#[serde(default)]
command: String,
}
pub fn parse_manifest_str(text: &str, fallback_name: &str) -> Result<PluginManifest> {
let raw: RawManifest = toml::from_str(text)
.map_err(|e| Error::tool("plugins", format!("parsing manifest: {e}")))?;
let name = raw
.name
.filter(|n| !n.trim().is_empty())
.unwrap_or_else(|| fallback_name.to_string());
let version = if raw.version.trim().is_empty() {
"0.0.0".to_string()
} else {
raw.version
};
let tools = raw
.tools
.into_iter()
.filter(|t| !t.name.trim().is_empty() && !t.command.trim().is_empty())
.map(|t| {
(
t.name,
PluginToolSpec {
command: t.command,
args: t.args,
description: t.description,
params: t
.params
.unwrap_or_else(|| serde_json::json!({"type": "object"})),
},
)
})
.collect();
let hooks = raw
.hooks
.into_iter()
.filter(|h| !h.event.trim().is_empty() && !h.command.trim().is_empty())
.map(|h| PluginHookSpec {
event: h.event,
command: h.command,
})
.collect();
Ok(PluginManifest {
name,
version,
tools,
hooks,
})
}
pub fn parse_manifest(path: &Path, fallback_name: &str) -> Result<PluginManifest> {
let text = std::fs::read_to_string(path)
.map_err(|e| Error::tool("plugins", format!("reading {}: {e}", path.display())))?;
parse_manifest_str(&text, fallback_name)
}
pub fn default_plugins_dir() -> PathBuf {
crate::agent::global_instructions_dir().join("plugins")
}
pub fn discover_manifests(dirs: &[PathBuf]) -> Vec<(String, PathBuf)> {
let mut found: std::collections::BTreeMap<String, PathBuf> = std::collections::BTreeMap::new();
for dir in dirs {
let Ok(entries) = std::fs::read_dir(dir) else {
continue;
};
for entry in entries.flatten() {
let path = entry.path();
if !path.is_dir() {
continue;
}
let manifest = path.join("plugin.toml");
if !manifest.is_file() {
continue;
}
let Some(name) = path.file_name().and_then(|n| n.to_str()) else {
continue;
};
found.insert(name.to_string(), manifest);
}
}
found.into_iter().collect()
}
#[cfg(unix)]
fn kill_group(pid: u32) {
crate::lsp::kill_process_group(pid);
}
#[cfg(not(unix))]
fn kill_group(_pid: u32) {}
fn shell_quote(s: &str) -> String {
format!("'{}'", s.replace('\'', "'\\''"))
}
async fn drain_capped<R>(mut reader: R, cap: usize) -> (String, bool)
where
R: tokio::io::AsyncRead + Unpin,
{
use tokio::io::AsyncReadExt;
let mut buf: Vec<u8> = Vec::new();
let mut truncated = false;
let mut chunk = [0u8; 8192];
loop {
match reader.read(&mut chunk).await {
Ok(0) => break,
Ok(n) => {
if buf.len() < cap {
let room = cap - buf.len();
let take = room.min(n);
buf.extend_from_slice(&chunk[..take]);
if take < n {
truncated = true;
}
} else {
truncated = true;
}
}
Err(_) => break,
}
}
(String::from_utf8_lossy(&buf).into_owned(), truncated)
}
#[derive(Debug, Clone)]
pub struct PluginTool {
name: String,
description: String,
params: Value,
command: String,
args: Vec<String>,
timeout: Duration,
}
impl PluginTool {
pub fn new(plugin_name: &str, tool_name: &str, spec: &PluginToolSpec) -> Self {
PluginTool {
name: format!("plugin__{plugin_name}__{tool_name}"),
description: spec.description.clone(),
params: spec.params.clone(),
command: spec.command.clone(),
args: spec.args.clone(),
timeout: Duration::from_secs(DEFAULT_PLUGIN_TOOL_TIMEOUT_SECS),
}
}
#[cfg(test)]
fn with_timeout(mut self, timeout: Duration) -> Self {
self.timeout = timeout;
self
}
}
#[async_trait]
impl Tool for PluginTool {
fn name(&self) -> &str {
&self.name
}
fn description(&self) -> &str {
&self.description
}
fn parameters(&self) -> Value {
self.params.clone()
}
async fn execute(&self, args: Value, ctx: &ToolContext) -> Result<String> {
let quoted = format!(
"{} {}",
shell_quote(&self.command),
self.args
.iter()
.map(|a| shell_quote(a))
.collect::<Vec<_>>()
.join(" ")
);
let mut cmd = crate::tools::build_sandboxed_sh("ed, ctx)?;
cmd.current_dir(&ctx.cwd)
.stdin(std::process::Stdio::piped())
.stdout(std::process::Stdio::piped())
.stderr(std::process::Stdio::piped())
.kill_on_drop(true);
#[cfg(unix)]
cmd.process_group(0);
let mut child = cmd
.spawn()
.map_err(|e| Error::tool("plugins", format!("spawn `{}`: {e}", self.command)))?;
let pid = child.id();
let args_json = serde_json::to_vec(&args)
.map_err(|e| Error::tool("plugins", format!("encoding tool args: {e}")))?;
let mut stdin = child
.stdin
.take()
.ok_or_else(|| Error::tool("plugins", "no stdin"))?;
let stdout = child
.stdout
.take()
.ok_or_else(|| Error::tool("plugins", "no stdout"))?;
let stderr = child
.stderr
.take()
.ok_or_else(|| Error::tool("plugins", "no stderr"))?;
let run = async {
use tokio::io::AsyncWriteExt;
let _ = stdin.write_all(&args_json).await;
let _ = stdin.flush().await;
drop(stdin);
let stdout_task = tokio::spawn(drain_capped(stdout, PLUGIN_TOOL_MAX_OUTPUT_BYTES));
let stderr_task = tokio::spawn(drain_capped(stderr, PLUGIN_TOOL_MAX_OUTPUT_BYTES));
let status = child.wait().await;
let (out, out_truncated) = stdout_task.await.unwrap_or_default();
let (err, err_truncated) = stderr_task.await.unwrap_or_default();
(status, out, out_truncated, err, err_truncated)
};
let outcome = tokio::time::timeout(self.timeout, run).await;
if let Some(pid) = pid {
kill_group(pid);
}
let (status, out, out_truncated, err, err_truncated) = match outcome {
Ok(result) => result,
Err(_) => {
return Err(Error::tool(
"plugins",
format!(
"plugin tool `{}` timed out after {:?}",
self.name, self.timeout
),
));
}
};
let mut result = out;
if out_truncated {
result.push_str(&format!(
"\n[plugin output truncated at {PLUGIN_TOOL_MAX_OUTPUT_BYTES} bytes]"
));
}
if !err.trim().is_empty() {
result.push_str("\n[stderr]\n");
result.push_str(&err);
if err_truncated {
result.push_str(&format!(
"\n[plugin stderr truncated at {PLUGIN_TOOL_MAX_OUTPUT_BYTES} bytes]"
));
}
}
match status {
Ok(s) if !s.success() => {
result.push_str(&format!(
"\n[plugin tool `{}` exited {}]",
self.name,
s.code().map(|c| c.to_string()).unwrap_or_default()
));
}
Err(e) => {
return Err(Error::tool(
"plugins",
format!("plugin tool `{}` wait failed: {e}", self.name),
));
}
_ => {}
}
Ok(result)
}
}
#[derive(Debug, Default)]
pub struct LoadedPlugins {
pub tools: Vec<PluginTool>,
pub hooks: Vec<(String, PluginHookSpec)>,
pub loaded_plugin_names: Vec<String>,
pub warnings: Vec<String>,
}
#[derive(Debug)]
pub enum PluginLoadOutcome {
Disabled,
BlockedPendingTrust,
Loaded(LoadedPlugins),
}
pub fn discover_and_load(config: &crate::Config) -> PluginLoadOutcome {
if !config.plugins_enabled {
return PluginLoadOutcome::Disabled;
}
if !is_trusted(config) {
return PluginLoadOutcome::BlockedPendingTrust;
}
let mut dirs = vec![default_plugins_dir()];
dirs.extend(config.plugins_dirs.iter().cloned());
let manifests = discover_manifests(&dirs);
let mut loaded = LoadedPlugins::default();
for (name, path) in manifests {
match parse_manifest(&path, &name) {
Ok(manifest) => {
for (tool_name, spec) in &manifest.tools {
loaded
.tools
.push(PluginTool::new(&manifest.name, tool_name, spec));
}
for hook in &manifest.hooks {
loaded.hooks.push((manifest.name.clone(), hook.clone()));
}
loaded.loaded_plugin_names.push(manifest.name);
}
Err(e) => {
loaded
.warnings
.push(format!("plugin `{name}` ({}): {e}", path.display()));
}
}
}
PluginLoadOutcome::Loaded(loaded)
}
pub fn register_into(config: &crate::Config, registry: &mut crate::tools::ToolRegistry) {
match discover_and_load(config) {
PluginLoadOutcome::Disabled => {}
PluginLoadOutcome::BlockedPendingTrust => {
eprintln!(
"warning: [capabilities.plugins] is enabled but this workspace is not trusted \
([capabilities.trust] default must be \"always\") — no plugin was loaded"
);
}
PluginLoadOutcome::Loaded(loaded) => {
for warning in &loaded.warnings {
eprintln!("warning: {warning}");
}
for (plugin, hook) in &loaded.hooks {
eprintln!(
"warning: plugin `{plugin}`'s `{}` hook is registered but is not yet \
emitted in this build (no-op) — see crate::plugins's module doc comment",
hook.event
);
}
for tool in loaded.tools {
registry.register(tool);
}
}
}
}
#[cfg(test)]
mod tests {
use super::*;
fn tmp(tag: &str) -> PathBuf {
let dir = std::env::temp_dir().join(format!(
"supercode-plugins-test-{tag}-{}-{}",
std::process::id(),
std::time::SystemTime::now()
.duration_since(std::time::UNIX_EPOCH)
.map(|d| d.as_nanos())
.unwrap_or(0)
));
std::fs::create_dir_all(&dir).unwrap();
dir
}
fn write_manifest(dir: &Path, name: &str, toml: &str) -> PathBuf {
let plugin_dir = dir.join(name);
std::fs::create_dir_all(&plugin_dir).unwrap();
let manifest = plugin_dir.join("plugin.toml");
std::fs::write(&manifest, toml).unwrap();
manifest
}
#[test]
fn trust_decision_parse_round_trips_known_values() {
assert_eq!(TrustDecision::parse("ask"), Some(TrustDecision::Ask));
assert_eq!(TrustDecision::parse("always"), Some(TrustDecision::Always));
assert_eq!(TrustDecision::parse("never"), Some(TrustDecision::Never));
assert_eq!(TrustDecision::parse("bogus"), None);
}
#[test]
fn is_trusted_requires_both_trust_enabled_and_default_always() {
let base = crate::Config::builder().model("m").build();
assert!(!is_trusted(&base), "trust disabled by default");
let enabled_ask = crate::Config::builder()
.model("m")
.trust_enabled(true)
.trust_default(TrustDecision::Ask)
.build();
assert!(
!is_trusted(&enabled_ask),
"trust enabled but default=ask must NOT be trusted (no interactive upgrade wired)"
);
let enabled_never = crate::Config::builder()
.model("m")
.trust_enabled(true)
.trust_default(TrustDecision::Never)
.build();
assert!(!is_trusted(&enabled_never));
let enabled_always = crate::Config::builder()
.model("m")
.trust_enabled(true)
.trust_default(TrustDecision::Always)
.build();
assert!(is_trusted(&enabled_always));
let disabled_always = crate::Config::builder()
.model("m")
.trust_enabled(false)
.trust_default(TrustDecision::Always)
.build();
assert!(
!is_trusted(&disabled_always),
"trust_enabled=false must gate regardless of trust_default"
);
}
#[test]
fn parse_manifest_str_parses_tools_and_hooks() {
let toml = r#"
name = "demo"
version = "1.2.3"
[[tools]]
name = "greet"
command = "echo"
args = ["hi"]
description = "says hi"
params = { type = "object" }
[[hooks]]
event = "post_tool"
command = "notify.sh"
"#;
let m = parse_manifest_str(toml, "fallback").unwrap();
assert_eq!(m.name, "demo");
assert_eq!(m.version, "1.2.3");
assert_eq!(m.tools.len(), 1);
assert_eq!(m.tools[0].0, "greet");
assert_eq!(m.tools[0].1.command, "echo");
assert_eq!(m.tools[0].1.args, vec!["hi".to_string()]);
assert_eq!(m.hooks.len(), 1);
assert_eq!(m.hooks[0].event, "post_tool");
assert_eq!(m.hooks[0].command, "notify.sh");
}
#[test]
fn parse_manifest_str_falls_back_to_directory_name_when_name_absent() {
let m = parse_manifest_str("version = \"0.1.0\"", "my-dir-name").unwrap();
assert_eq!(m.name, "my-dir-name");
}
#[test]
fn parse_manifest_str_defaults_version_and_params_when_absent() {
let toml = r#"
[[tools]]
name = "t"
command = "echo"
"#;
let m = parse_manifest_str(toml, "p").unwrap();
assert_eq!(m.version, "0.0.0");
assert_eq!(m.tools[0].1.params, serde_json::json!({"type": "object"}));
}
#[test]
fn parse_manifest_str_skips_malformed_entries_without_failing_the_manifest() {
let toml = r#"
[[tools]]
name = ""
command = "echo"
[[tools]]
name = "ok"
command = ""
[[tools]]
name = "good"
command = "echo"
[[hooks]]
event = ""
command = "x"
"#;
let m = parse_manifest_str(toml, "p").unwrap();
assert_eq!(m.tools.len(), 1, "only the fully-valid tool survives");
assert_eq!(m.tools[0].0, "good");
assert!(m.hooks.is_empty());
}
#[test]
fn parse_manifest_str_rejects_malformed_toml() {
assert!(parse_manifest_str("not valid toml [[[", "p").is_err());
}
#[test]
fn discover_manifests_finds_plugin_toml_under_immediate_subdirs() {
let dir = tmp("discover");
write_manifest(&dir, "alpha", "name = \"alpha\"\n");
write_manifest(&dir, "beta", "name = \"beta\"\n");
std::fs::create_dir_all(dir.join("not-a-plugin")).unwrap();
let found = discover_manifests(std::slice::from_ref(&dir));
let names: Vec<&str> = found.iter().map(|(n, _)| n.as_str()).collect();
assert_eq!(names, vec!["alpha", "beta"], "sorted by name");
std::fs::remove_dir_all(&dir).ok();
}
#[test]
fn discover_manifests_missing_dir_is_silently_skipped() {
let missing = tmp("missing-parent").join("does-not-exist");
let found = discover_manifests(&[missing]);
assert!(found.is_empty());
}
#[test]
fn discover_manifests_later_dir_wins_on_name_collision() {
let dir_a = tmp("collide-a");
let dir_b = tmp("collide-b");
write_manifest(&dir_a, "dup", "version = \"1.0.0\"\n");
write_manifest(&dir_b, "dup", "version = \"2.0.0\"\n");
let found = discover_manifests(&[dir_a.clone(), dir_b.clone()]);
assert_eq!(found.len(), 1);
let (_, path) = &found[0];
assert!(path.starts_with(&dir_b), "later dir must win");
std::fs::remove_dir_all(&dir_a).ok();
std::fs::remove_dir_all(&dir_b).ok();
}
#[test]
fn discover_and_load_is_disabled_when_plugins_off_default_off_byte_identity() {
let config = crate::Config::builder().model("m").build();
assert!(!config.plugins_enabled);
assert!(matches!(
discover_and_load(&config),
PluginLoadOutcome::Disabled
));
}
#[test]
fn discover_and_load_is_blocked_pending_trust_when_untrusted() {
let config = crate::Config::builder()
.model("m")
.plugins_enabled(true)
.trust_enabled(true)
.trust_default(TrustDecision::Ask)
.build();
assert!(matches!(
discover_and_load(&config),
PluginLoadOutcome::BlockedPendingTrust
));
}
#[test]
fn discover_and_load_is_blocked_when_trust_module_itself_is_off() {
let config = crate::Config::builder()
.model("m")
.plugins_enabled(true)
.trust_enabled(false)
.build();
assert!(matches!(
discover_and_load(&config),
PluginLoadOutcome::BlockedPendingTrust
));
}
#[test]
fn discover_and_load_loads_tools_from_a_trusted_configured_dir() {
let dir = tmp("load-trusted");
write_manifest(
&dir,
"demo",
"name = \"demo\"\n\n[[tools]]\nname = \"echo_it\"\ncommand = \"echo\"\n",
);
let config = crate::Config::builder()
.model("m")
.plugins_enabled(true)
.plugins_dirs(vec![dir.clone()])
.trust_enabled(true)
.trust_default(TrustDecision::Always)
.build();
match discover_and_load(&config) {
PluginLoadOutcome::Loaded(loaded) => {
assert_eq!(loaded.loaded_plugin_names, vec!["demo".to_string()]);
assert_eq!(loaded.tools.len(), 1);
assert_eq!(loaded.tools[0].name(), "plugin__demo__echo_it");
}
other => panic!("expected Loaded, got {other:?}"),
}
std::fs::remove_dir_all(&dir).ok();
}
#[test]
fn discover_and_load_records_a_warning_for_an_unparseable_manifest_without_failing_others() {
let dir = tmp("load-warn");
write_manifest(&dir, "bad", "not valid toml [[[");
write_manifest(
&dir,
"good",
"name = \"good\"\n\n[[tools]]\nname = \"t\"\ncommand = \"echo\"\n",
);
let config = crate::Config::builder()
.model("m")
.plugins_enabled(true)
.plugins_dirs(vec![dir.clone()])
.trust_enabled(true)
.trust_default(TrustDecision::Always)
.build();
match discover_and_load(&config) {
PluginLoadOutcome::Loaded(loaded) => {
assert_eq!(loaded.tools.len(), 1, "the good plugin still loads");
assert_eq!(loaded.warnings.len(), 1);
assert!(loaded.warnings[0].contains("bad"));
}
other => panic!("expected Loaded, got {other:?}"),
}
std::fs::remove_dir_all(&dir).ok();
}
#[test]
fn register_into_is_a_true_noop_when_plugins_disabled() {
let config = crate::Config::builder().model("m").build();
let mut registry = crate::tools::ToolRegistry::new();
register_into(&config, &mut registry);
assert_eq!(registry.len(), 0);
}
#[test]
fn register_into_registers_nothing_when_untrusted() {
let dir = tmp("register-untrusted");
write_manifest(
&dir,
"demo",
"name = \"demo\"\n\n[[tools]]\nname = \"t\"\ncommand = \"echo\"\n",
);
let config = crate::Config::builder()
.model("m")
.plugins_enabled(true)
.plugins_dirs(vec![dir.clone()])
.trust_enabled(true)
.trust_default(TrustDecision::Never)
.build();
let mut registry = crate::tools::ToolRegistry::new();
register_into(&config, &mut registry);
assert_eq!(registry.len(), 0, "untrusted plugin must never register");
std::fs::remove_dir_all(&dir).ok();
}
#[test]
fn register_into_registers_trusted_tools() {
let dir = tmp("register-trusted");
write_manifest(
&dir,
"demo",
"name = \"demo\"\n\n[[tools]]\nname = \"t\"\ncommand = \"echo\"\n",
);
let config = crate::Config::builder()
.model("m")
.plugins_enabled(true)
.plugins_dirs(vec![dir.clone()])
.trust_enabled(true)
.trust_default(TrustDecision::Always)
.build();
let mut registry = crate::tools::ToolRegistry::new();
register_into(&config, &mut registry);
assert_eq!(registry.len(), 1);
assert!(registry.get("plugin__demo__t").is_some());
std::fs::remove_dir_all(&dir).ok();
}
fn ctx(cwd: PathBuf) -> ToolContext {
ToolContext::new(cwd)
}
#[tokio::test]
async fn plugin_tool_executes_and_returns_stdout() {
let dir = tmp("exec-basic");
let spec = PluginToolSpec {
command: "echo".to_string(),
args: vec!["hello-plugin".to_string()],
description: String::new(),
params: serde_json::json!({"type": "object"}),
};
let tool = PluginTool::new("demo", "say", &spec);
let out = tool
.execute(serde_json::json!({}), &ctx(dir.clone()))
.await
.unwrap();
assert!(out.contains("hello-plugin"), "{out}");
std::fs::remove_dir_all(&dir).ok();
}
#[tokio::test]
async fn plugin_tool_args_are_never_shell_spliced() {
let dir = tmp("exec-no-splice");
let marker = dir.join("PWNED");
let spec = PluginToolSpec {
command: "cat".to_string(),
args: vec![],
description: String::new(),
params: serde_json::json!({"type": "object"}),
};
let tool = PluginTool::new("demo", "cat_args", &spec);
let evil = format!("$(touch {})", marker.display());
let out = tool
.execute(serde_json::json!({"payload": evil}), &ctx(dir.clone()))
.await
.unwrap();
assert!(
out.contains("touch"),
"cat should echo the literal, unevaluated JSON back: {out}"
);
assert!(
!marker.exists(),
"shell metacharacters in tool-call args must never be evaluated"
);
std::fs::remove_dir_all(&dir).ok();
}
#[tokio::test]
async fn plugin_tool_bounds_output_and_marks_it_truncated() {
let dir = tmp("exec-bounded");
let spec = PluginToolSpec {
command: "sh".to_string(),
args: vec![
"-c".to_string(),
format!(
"head -c {} /dev/zero | tr '\\0' 'a'",
PLUGIN_TOOL_MAX_OUTPUT_BYTES * 2
),
],
description: String::new(),
params: serde_json::json!({"type": "object"}),
};
let tool = PluginTool::new("demo", "flood", &spec);
let out = tool
.execute(serde_json::json!({}), &ctx(dir.clone()))
.await
.unwrap();
assert!(out.contains("truncated"), "{}", &out[..out.len().min(200)]);
assert!(
out.len() < PLUGIN_TOOL_MAX_OUTPUT_BYTES * 2,
"retained output must be bounded well below what the child wrote"
);
std::fs::remove_dir_all(&dir).ok();
}
#[tokio::test]
async fn plugin_tool_timeout_is_bounded_and_reported() {
let dir = tmp("exec-timeout");
let spec = PluginToolSpec {
command: "sleep".to_string(),
args: vec!["3600".to_string()],
description: String::new(),
params: serde_json::json!({"type": "object"}),
};
let tool = PluginTool::new("demo", "hang", &spec).with_timeout(Duration::from_millis(300));
let started = std::time::Instant::now();
let result = tokio::time::timeout(
Duration::from_secs(10),
tool.execute(serde_json::json!({}), &ctx(dir.clone())),
)
.await
.expect("must not hang past the plugin tool's own timeout");
assert!(result.is_err(), "a hanging plugin tool must error out");
assert!(
result.unwrap_err().to_string().contains("timed out"),
"error should say it timed out"
);
assert!(
started.elapsed() < Duration::from_secs(5),
"took {:?}, expected to bail out near the configured timeout",
started.elapsed()
);
std::fs::remove_dir_all(&dir).ok();
}
#[cfg(unix)]
#[tokio::test]
async fn plugin_tool_reaps_grandchild_worker_processes() {
let dir = tmp("exec-grandchild");
let pidfile = dir.join("worker.pid");
let script = dir.join("spawn_worker.sh");
std::fs::write(
&script,
format!(
"#!/bin/sh\nsleep 3600 &\necho $! > {}\nwait\n",
pidfile.display()
),
)
.unwrap();
let spec = PluginToolSpec {
command: "sh".to_string(),
args: vec![script.to_string_lossy().into_owned()],
description: String::new(),
params: serde_json::json!({"type": "object"}),
};
let tool = PluginTool::new("demo", "spawns_worker", &spec)
.with_timeout(Duration::from_millis(300));
let handle = tokio::spawn({
let dir = dir.clone();
async move { tool.execute(serde_json::json!({}), &ctx(dir)).await }
});
let mut grandchild_pid: Option<i32> = None;
for _ in 0..150 {
if let Ok(s) = std::fs::read_to_string(&pidfile) {
if let Ok(pid) = s.trim().parse::<i32>() {
grandchild_pid = Some(pid);
break;
}
}
tokio::time::sleep(Duration::from_millis(20)).await;
}
let grandchild_pid = grandchild_pid.expect("worker must have recorded its pid");
assert!(
unsafe { libc::kill(grandchild_pid, 0) == 0 },
"grandchild worker must be alive before the plugin tool call completes"
);
let _ = tokio::time::timeout(Duration::from_secs(10), handle).await;
let mut still_alive = true;
for _ in 0..150 {
let alive = unsafe { libc::kill(grandchild_pid, 0) == 0 };
if !alive {
still_alive = false;
break;
}
tokio::time::sleep(Duration::from_millis(20)).await;
}
assert!(
!still_alive,
"grandchild worker pid {grandchild_pid} must be dead — it must not orphan"
);
std::fs::remove_dir_all(&dir).ok();
}
}