mirror of
https://github.com/block/buzz.git
synced 2026-08-18 06:50:31 +02:00
fix(runtime): sweep node wrapper processes hosting managed agent shims (#1296)
Signed-off-by: Will Pfleger <pfleger.will@gmail.com> Co-authored-by: npub1mn7jgtj4w2pd0g0zeuhxsa6jy6p0rewxz4kujt98my82ahfmp72sxjexk7 <dcfd242e557282d7a1e2cf2e6877522682f1e5c6156dc92ca7d90eaedd3b0f95@sprout-oss.stage.blox.sqprod.co> Co-authored-by: npub1fgdl5qqnh3k3f2xkqrvt7cujalhm623x4s7fdjdj5yrtp5fzjl9qrjpucw <4a1bfa0013bc6d14a8d600d8bf6392efefbd2a26ac3c96c9b2a106b0d12297ca@sprout-oss.stage.blox.sqprod.co>
This commit is contained in:
co-authored by
npub1mn7jgtj4w2pd0g0zeuhxsa6jy6p0rewxz4kujt98my82ahfmp72sxjexk7
npub1fgdl5qqnh3k3f2xkqrvt7cujalhm623x4s7fdjdj5yrtp5fzjl9qrjpucw
parent
8717ddf2eb
commit
f072032aaf
@@ -47,7 +47,7 @@ const overrides = new Map([
|
||||
// harness-persona-sync: persona-runtime resolution threaded into the spawn
|
||||
// path here. Load-bearing feature growth; queued to split in the resolver
|
||||
// unify refactor followup.
|
||||
["src-tauri/src/managed_agents/runtime.rs", 2001],
|
||||
["src-tauri/src/managed_agents/runtime.rs", 2031],
|
||||
["src-tauri/src/managed_agents/personas.rs", 1080],
|
||||
// Phase-2 inbound reconcile + review-fix cycle: reconcile_inbound_persona_event
|
||||
// dispatches 30175/30176/30177 inbound plus kind:5 tombstone consume
|
||||
@@ -83,7 +83,7 @@ const overrides = new Map([
|
||||
// syncs team-dir edits before all personas.json readers; run_event_sync
|
||||
// signs the persona/team retention events post-identity) layered on top of
|
||||
// main's growth. Load-bearing feature growth, queued to split with the list.
|
||||
["src-tauri/src/lib.rs", 1026],
|
||||
["src-tauri/src/lib.rs", 1034],
|
||||
// onMarkRead + isUnread prop threading (mirrors the onMarkUnread prop
|
||||
// already here) for the single-toggle mark-read/unread menu item — a small
|
||||
// overage from load-bearing per-message plumbing, not generic debt growth.
|
||||
|
||||
@@ -6,9 +6,9 @@ use tauri::{AppHandle, State};
|
||||
use crate::{
|
||||
app_state::AppState,
|
||||
managed_agents::{
|
||||
build_managed_agent_summary, default_agent_workdir, find_managed_agent_mut,
|
||||
known_acp_runtime, load_managed_agents, load_personas, managed_agent_avatar_url,
|
||||
missing_command_message, normalize_agent_args, resolve_command,
|
||||
build_managed_agent_summary, current_instance_id, default_agent_workdir,
|
||||
find_managed_agent_mut, known_acp_runtime, load_managed_agents, load_personas,
|
||||
managed_agent_avatar_url, missing_command_message, normalize_agent_args, resolve_command,
|
||||
resolve_effective_prompt_model_provider, save_managed_agents, sync_managed_agent_processes,
|
||||
try_regenerate_nest, AgentModelInfo, AgentModelsResponse, UpdateManagedAgentRequest,
|
||||
UpdateManagedAgentResponse,
|
||||
@@ -37,7 +37,7 @@ pub async fn get_agent_models(
|
||||
.managed_agent_processes
|
||||
.lock()
|
||||
.map_err(|e| e.to_string())?;
|
||||
if sync_managed_agent_processes(&mut records, &mut runtimes) {
|
||||
if sync_managed_agent_processes(&mut records, &mut runtimes, ¤t_instance_id(&app)) {
|
||||
save_managed_agents(&app, &records)?;
|
||||
}
|
||||
|
||||
@@ -166,7 +166,7 @@ pub async fn update_managed_agent(
|
||||
.managed_agent_processes
|
||||
.lock()
|
||||
.map_err(|e| e.to_string())?;
|
||||
sync_managed_agent_processes(&mut records, &mut runtimes);
|
||||
sync_managed_agent_processes(&mut records, &mut runtimes, ¤t_instance_id(&app));
|
||||
|
||||
let record = find_managed_agent_mut(&mut records, &input.pubkey)?;
|
||||
|
||||
|
||||
@@ -3,8 +3,9 @@ use tauri::{AppHandle, State};
|
||||
use crate::{
|
||||
app_state::AppState,
|
||||
managed_agents::{
|
||||
build_managed_agent_summary, find_managed_agent_mut, load_managed_agents, load_personas,
|
||||
save_managed_agents, sync_managed_agent_processes, ManagedAgentSummary,
|
||||
build_managed_agent_summary, current_instance_id, find_managed_agent_mut,
|
||||
load_managed_agents, load_personas, save_managed_agents, sync_managed_agent_processes,
|
||||
ManagedAgentSummary,
|
||||
},
|
||||
util::now_iso,
|
||||
};
|
||||
@@ -26,7 +27,7 @@ pub fn set_managed_agent_start_on_app_launch(
|
||||
.lock()
|
||||
.map_err(|error| error.to_string())?;
|
||||
|
||||
if sync_managed_agent_processes(&mut records, &mut runtimes) {
|
||||
if sync_managed_agent_processes(&mut records, &mut runtimes, ¤t_instance_id(&app)) {
|
||||
save_managed_agents(&app, &records)?;
|
||||
}
|
||||
|
||||
|
||||
@@ -4,9 +4,9 @@ use tauri::{AppHandle, State};
|
||||
use crate::{
|
||||
app_state::AppState,
|
||||
managed_agents::{
|
||||
build_managed_agent_summary, discover_provider_candidates, ensure_persona_is_active,
|
||||
find_managed_agent_mut, invoke_provider, load_managed_agents, load_personas,
|
||||
managed_agent_avatar_url, managed_agent_log_path, managed_agents_base_dir,
|
||||
build_managed_agent_summary, current_instance_id, discover_provider_candidates,
|
||||
ensure_persona_is_active, find_managed_agent_mut, invoke_provider, load_managed_agents,
|
||||
load_personas, managed_agent_avatar_url, managed_agent_log_path, managed_agents_base_dir,
|
||||
normalize_agent_args, provider_deploy, read_log_tail, resolve_provider_binary,
|
||||
save_managed_agents, spawn_key_refusal, start_managed_agent_process,
|
||||
stop_managed_agent_process, sync_managed_agent_processes, try_regenerate_nest,
|
||||
@@ -433,7 +433,7 @@ pub fn list_managed_agents(
|
||||
.lock()
|
||||
.map_err(|error| error.to_string())?;
|
||||
|
||||
if sync_managed_agent_processes(&mut records, &mut runtimes) {
|
||||
if sync_managed_agent_processes(&mut records, &mut runtimes, ¤t_instance_id(&app)) {
|
||||
save_managed_agents(&app, &records)?;
|
||||
}
|
||||
|
||||
@@ -497,7 +497,7 @@ pub async fn create_managed_agent(
|
||||
.lock()
|
||||
.map_err(|error| error.to_string())?;
|
||||
|
||||
if sync_managed_agent_processes(&mut records, &mut runtimes) {
|
||||
if sync_managed_agent_processes(&mut records, &mut runtimes, ¤t_instance_id(&app)) {
|
||||
save_managed_agents(&app, &records)?;
|
||||
}
|
||||
if let Some(persona_id) = requested_persona_id.as_deref() {
|
||||
@@ -563,7 +563,7 @@ pub async fn create_managed_agent(
|
||||
.lock()
|
||||
.map_err(|error| error.to_string())?;
|
||||
|
||||
if sync_managed_agent_processes(&mut records, &mut runtimes) {
|
||||
if sync_managed_agent_processes(&mut records, &mut runtimes, ¤t_instance_id(&app)) {
|
||||
save_managed_agents(&app, &records)?;
|
||||
}
|
||||
|
||||
@@ -949,7 +949,7 @@ pub async fn start_managed_agent(
|
||||
.lock()
|
||||
.map_err(|error| error.to_string())?;
|
||||
|
||||
if sync_managed_agent_processes(&mut records, &mut runtimes) {
|
||||
if sync_managed_agent_processes(&mut records, &mut runtimes, ¤t_instance_id(&app)) {
|
||||
save_managed_agents(&app, &records)?;
|
||||
}
|
||||
|
||||
@@ -1205,7 +1205,7 @@ pub fn stop_managed_agent(
|
||||
.lock()
|
||||
.map_err(|error| error.to_string())?;
|
||||
|
||||
if sync_managed_agent_processes(&mut records, &mut runtimes) {
|
||||
if sync_managed_agent_processes(&mut records, &mut runtimes, ¤t_instance_id(&app)) {
|
||||
save_managed_agents(&app, &records)?;
|
||||
}
|
||||
|
||||
@@ -1247,7 +1247,7 @@ pub fn delete_managed_agent(
|
||||
.lock()
|
||||
.map_err(|error| error.to_string())?;
|
||||
|
||||
if sync_managed_agent_processes(&mut records, &mut runtimes) {
|
||||
if sync_managed_agent_processes(&mut records, &mut runtimes, ¤t_instance_id(&app)) {
|
||||
save_managed_agents(&app, &records)?;
|
||||
}
|
||||
|
||||
|
||||
@@ -133,8 +133,16 @@ fn shutdown_managed_agents(app: &tauri::AppHandle) -> Result<(), String> {
|
||||
.managed_agent_processes
|
||||
.lock()
|
||||
.map_err(|error| error.to_string())?;
|
||||
let mut changed = sync_managed_agent_processes(&mut records, &mut runtimes);
|
||||
changed |= kill_stale_tracked_processes(&mut records, &runtimes);
|
||||
let mut changed = sync_managed_agent_processes(
|
||||
&mut records,
|
||||
&mut runtimes,
|
||||
&managed_agents::current_instance_id(app),
|
||||
);
|
||||
changed |= kill_stale_tracked_processes(
|
||||
&mut records,
|
||||
&runtimes,
|
||||
&managed_agents::current_instance_id(app),
|
||||
);
|
||||
|
||||
// Stop all tracked agents. Send SIGTERM to all process
|
||||
// groups first, then wait for exits in parallel to avoid serial 1s waits.
|
||||
|
||||
@@ -107,8 +107,13 @@ pub async fn restore_managed_agents_on_launch(
|
||||
.managed_agent_processes
|
||||
.lock()
|
||||
.map_err(|error| error.to_string())?;
|
||||
let mut changed = sync_managed_agent_processes(&mut records, &mut runtimes);
|
||||
changed |= kill_stale_tracked_processes(&mut records, &runtimes);
|
||||
let mut changed = sync_managed_agent_processes(
|
||||
&mut records,
|
||||
&mut runtimes,
|
||||
&super::current_instance_id(app),
|
||||
);
|
||||
changed |=
|
||||
kill_stale_tracked_processes(&mut records, &runtimes, &super::current_instance_id(app));
|
||||
|
||||
let tracked_pids: Vec<u32> = records
|
||||
.iter()
|
||||
|
||||
@@ -40,6 +40,12 @@ pub(crate) const KNOWN_AGENT_BINARIES: &[&str] = &[
|
||||
"buzz_dev_mcp",
|
||||
];
|
||||
|
||||
/// Script interpreters that may host managed agent wrappers (e.g. npm shims).
|
||||
/// A process whose name matches here is NOT immediately claimed — it must also
|
||||
/// carry `BUZZ_MANAGED_AGENT` in its environment (checked by the caller via
|
||||
/// `process_has_buzz_marker()`). This avoids sweeping unrelated node processes.
|
||||
pub(crate) const KNOWN_SCRIPT_INTERPRETERS: &[&str] = &["node"];
|
||||
|
||||
/// Check if a process name matches any of our known agent binaries.
|
||||
/// Uses exact match or prefix-with-separator to avoid false positives
|
||||
/// (e.g. `"goose"` must not match `"mongoose"`).
|
||||
@@ -54,6 +60,15 @@ fn name_matches_known_binary(name: &str) -> bool {
|
||||
})
|
||||
}
|
||||
|
||||
/// Check if a process name is a known script interpreter that may be hosting
|
||||
/// a managed agent wrapper (e.g. `node` running an npm shim for `codex-acp`).
|
||||
/// Callers must additionally verify `BUZZ_MANAGED_AGENT` ownership.
|
||||
fn name_matches_interpreter(name: &str) -> bool {
|
||||
KNOWN_SCRIPT_INTERPRETERS
|
||||
.iter()
|
||||
.any(|&interp| name == interp)
|
||||
}
|
||||
|
||||
#[cfg(unix)]
|
||||
pub(crate) fn process_is_running(pid: u32) -> bool {
|
||||
// Use libc::kill with signal 0 instead of forking a subprocess.
|
||||
@@ -89,7 +104,9 @@ pub(crate) fn process_belongs_to_us(pid: u32) -> bool {
|
||||
return false;
|
||||
}
|
||||
let name = String::from_utf8_lossy(&buf[..len as usize]);
|
||||
name_matches_known_binary(&name)
|
||||
// Fall through for script interpreters (e.g. `node` hosting an npm shim):
|
||||
// the caller's `process_has_buzz_marker()` check decides true ownership.
|
||||
name_matches_known_binary(&name) || name_matches_interpreter(&name)
|
||||
}
|
||||
|
||||
#[cfg(all(unix, not(target_os = "macos")))]
|
||||
@@ -101,13 +118,18 @@ pub(crate) fn process_belongs_to_us(pid: u32) -> bool {
|
||||
if name_matches_known_binary(name.trim()) {
|
||||
return true;
|
||||
}
|
||||
// Interpreter check: `node` is 4 bytes, never truncated.
|
||||
if name_matches_interpreter(name.trim()) {
|
||||
return true;
|
||||
}
|
||||
}
|
||||
|
||||
// Fallback: read /proc/<pid>/exe which is a symlink to the full binary path.
|
||||
// This is not subject to the 15-byte truncation limit.
|
||||
if let Ok(exe_path) = std::fs::read_link(format!("/proc/{pid}/exe")) {
|
||||
if let Some(basename) = exe_path.file_name().and_then(|n| n.to_str()) {
|
||||
return name_matches_known_binary(basename);
|
||||
// Fall through for script interpreters — caller checks the marker.
|
||||
return name_matches_known_binary(basename) || name_matches_interpreter(basename);
|
||||
}
|
||||
}
|
||||
|
||||
@@ -1190,6 +1212,7 @@ pub(crate) fn reap_dead_instance_agents(_our_instance_id: &str, _skip_pids: &[u3
|
||||
pub fn kill_stale_tracked_processes(
|
||||
records: &mut [ManagedAgentRecord],
|
||||
runtimes: &HashMap<String, ManagedAgentProcess>,
|
||||
instance_id: &str,
|
||||
) -> bool {
|
||||
use crate::managed_agents::BackendKind;
|
||||
|
||||
@@ -1202,7 +1225,7 @@ pub fn kill_stale_tracked_processes(
|
||||
continue;
|
||||
};
|
||||
if !runtimes.contains_key(&record.pubkey) {
|
||||
if process_belongs_to_us(pid) {
|
||||
if process_belongs_to_us(pid) && process_has_buzz_marker(pid, instance_id) {
|
||||
let _ = terminate_process(pid);
|
||||
}
|
||||
record.runtime_pid = None;
|
||||
@@ -1217,6 +1240,7 @@ pub fn kill_stale_tracked_processes(
|
||||
pub fn sync_managed_agent_processes(
|
||||
records: &mut [ManagedAgentRecord],
|
||||
runtimes: &mut HashMap<String, ManagedAgentProcess>,
|
||||
instance_id: &str,
|
||||
) -> bool {
|
||||
let mut changed = false;
|
||||
let mut exited = Vec::new();
|
||||
@@ -1270,7 +1294,10 @@ pub fn sync_managed_agent_processes(
|
||||
continue;
|
||||
};
|
||||
|
||||
if process_is_running(pid) && process_belongs_to_us(pid) {
|
||||
if process_is_running(pid)
|
||||
&& process_belongs_to_us(pid)
|
||||
&& process_has_buzz_marker(pid, instance_id)
|
||||
{
|
||||
continue;
|
||||
}
|
||||
|
||||
@@ -1880,7 +1907,10 @@ pub fn start_managed_agent_process(
|
||||
}
|
||||
|
||||
if let Some(pid) = record.runtime_pid {
|
||||
if process_is_running(pid) && process_belongs_to_us(pid) {
|
||||
if process_is_running(pid)
|
||||
&& process_belongs_to_us(pid)
|
||||
&& process_has_buzz_marker(pid, ¤t_instance_id(app))
|
||||
{
|
||||
record.updated_at = now_iso();
|
||||
record.last_error = None;
|
||||
return Ok(());
|
||||
|
||||
@@ -544,3 +544,35 @@ fn runtime_metadata_env_vars_injects_model_even_with_acp_model_switching() {
|
||||
]
|
||||
);
|
||||
}
|
||||
|
||||
// ── name_matches_known_binary / name_matches_interpreter tests ───────────
|
||||
|
||||
#[test]
|
||||
fn name_matches_known_binary_rejects_node() {
|
||||
// `node` must NOT be in KNOWN_AGENT_BINARIES — adding it there would
|
||||
// sweep all node processes on the machine regardless of ownership.
|
||||
assert!(!super::name_matches_known_binary("node"));
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn name_matches_interpreter_accepts_node() {
|
||||
// `node` IS a known script interpreter and must be recognized.
|
||||
assert!(super::name_matches_interpreter("node"));
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn name_matches_interpreter_rejects_unknown() {
|
||||
// Interpreters not in KNOWN_SCRIPT_INTERPRETERS must not match.
|
||||
assert!(!super::name_matches_interpreter("python3"));
|
||||
assert!(!super::name_matches_interpreter("deno"));
|
||||
assert!(!super::name_matches_interpreter("bun"));
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn name_matches_interpreter_rejects_node_prefix() {
|
||||
// A name that starts with "node" but is longer must not match —
|
||||
// exact equality is required to avoid false positives.
|
||||
assert!(!super::name_matches_interpreter("node_modules"));
|
||||
assert!(!super::name_matches_interpreter("nodejs"));
|
||||
assert!(!super::name_matches_interpreter("node-gyp"));
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user