diff --git a/Cargo.lock b/Cargo.lock index 40f33e13b..11f99d900 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -795,6 +795,7 @@ dependencies = [ "serde", "serde_json", "sha2 0.11.0", + "tempfile", "thiserror 2.0.18", "tokio", "tokio-tungstenite 0.29.0", diff --git a/crates/buzz-acp/Cargo.toml b/crates/buzz-acp/Cargo.toml index d04784980..37a0c69ba 100644 --- a/crates/buzz-acp/Cargo.toml +++ b/crates/buzz-acp/Cargo.toml @@ -71,6 +71,9 @@ toml = "1.0" # Filter expressions evalexpr = { workspace = true } +# Scratch files for atomic per-turn reply context updates +tempfile = "3" + # Process-group kill (safe wrapper around killpg) — Unix-only; kill_process_group # has a #[cfg(not(unix))] fallback in acp.rs. [target.'cfg(unix)'.dependencies] diff --git a/crates/buzz-acp/src/acp.rs b/crates/buzz-acp/src/acp.rs index 9eb668cbc..b44a7a51b 100644 --- a/crates/buzz-acp/src/acp.rs +++ b/crates/buzz-acp/src/acp.rs @@ -509,7 +509,8 @@ impl AcpClient { // Handled by build_codex_config_env; skip here to avoid double-setting. continue; } - if std::env::var_os(key).is_none() { + // Reply context is slot-scoped harness state, not an operator override. + if key == "BUZZ_REPLY_TO_FILE" || std::env::var_os(key).is_none() { cmd.env(key, value); } } diff --git a/crates/buzz-acp/src/config.rs b/crates/buzz-acp/src/config.rs index 19304bf18..6b7955e3a 100644 --- a/crates/buzz-acp/src/config.rs +++ b/crates/buzz-acp/src/config.rs @@ -542,6 +542,9 @@ pub struct Config { /// Per-persona env vars to inject at agent spawn time (e.g., GOOSE_PROVIDER, GOOSE_MODEL, BUZZ_AGENT_MODEL). /// Populated from persona pack resolution. Empty when no pack is configured. pub persona_env_vars: Vec<(String, String)>, + /// Stable per-slot files carrying the current automatic reply anchor. + /// Populated by the harness after setup-mode handling and before pool startup. + pub reply_to_files: Vec, /// Whether `codex_network_env()` successfully injected a `CODEX_CONFIG` entry into /// `persona_env_vars`. When true, `AcpClient::spawn` merges all `CODEX_CONFIG` entries /// and forces `sandbox_workspace_write.network_access = true` via `build_codex_config_env`. @@ -1096,6 +1099,7 @@ impl Config { respond_to_allowlist, allowed_respond_to, persona_env_vars, + reply_to_files: vec![], has_generated_codex_config, relay_observer: args.relay_observer, lazy_pool: args.lazy_pool, @@ -1466,6 +1470,7 @@ mod tests { respond_to_allowlist: HashSet::new(), allowed_respond_to: Vec::new(), persona_env_vars: vec![], + reply_to_files: vec![], has_generated_codex_config: false, relay_observer: false, lazy_pool: false, diff --git a/crates/buzz-acp/src/lib.rs b/crates/buzz-acp/src/lib.rs index d63f720c6..606592f7d 100644 --- a/crates/buzz-acp/src/lib.rs +++ b/crates/buzz-acp/src/lib.rs @@ -1294,6 +1294,15 @@ async fn tokio_main() -> Result<()> { tracing::info!("buzz-acp starting: {}", config.summary()); + let reply_to_dir = tempfile::tempdir()?; + let reply_to_files: Vec = (0..config.agents) + .map(|index| reply_to_dir.path().join(format!("agent-{index}"))) + .collect(); + for path in &reply_to_files { + std::fs::write(path, "")?; + } + config.reply_to_files = reply_to_files; + let observer = config .relay_observer .then(observer::ObserverHandle::in_process); @@ -1315,7 +1324,11 @@ async fn tokio_main() -> Result<()> { let mut pool = if config.lazy_pool { AgentPool::from_slots((0..config.agents).map(|_| None).collect()) } else { - initialize_agent_pool(&PoolStartup::from_config(&config, observer.clone()), None).await? + initialize_agent_pool( + &PoolStartup::from_config(&config, observer.clone(), config.reply_to_files.clone()), + None, + ) + .await? }; let mut pool_ready = !config.lazy_pool; let mut pool_lifecycle: PoolLifecycle = PoolLifecycle::listening(); @@ -1559,6 +1572,7 @@ async fn tokio_main() -> Result<()> { memory_enabled: config.memory_enabled, harness_name: crate::config::normalize_agent_command_identity(&config.agent_command), relay_url: config.relay_url.clone(), + reply_to_files: config.reply_to_files.clone(), }); if !config.memory_enabled { @@ -1723,7 +1737,11 @@ async fn tokio_main() -> Result<()> { "waking", None, ); - let startup = PoolStartup::from_config(&config, observer.clone()); + let startup = PoolStartup::from_config( + &config, + observer.clone(), + config.reply_to_files.clone(), + ); let wake_tx = wake_tx.clone(); let wake_shutdown = shutdown_rx.clone(); wake_tasks.spawn(async move { @@ -1761,9 +1779,19 @@ async fn tokio_main() -> Result<()> { let env = config.persona_env_vars.clone(); let has_codex = config.has_generated_codex_config; let observer = observer.clone(); + let reply_to_file = config.reply_to_files.get(idx).cloned(); let guard = RespawnGuard::new(idx, respawn_tx.clone()); respawn_tasks.spawn(async move { - let result = spawn_and_init(&cmd, &args, &env, has_codex, idx, observer).await; + let result = spawn_and_init( + &cmd, + &args, + &env, + reply_to_file.as_deref(), + has_codex, + idx, + observer, + ) + .await; guard.send(result); }); } @@ -3507,13 +3535,23 @@ fn recover_panicked_agent( let cmd = config.agent_command.clone(); let args = config.agent_args.clone(); let env = config.persona_env_vars.clone(); + let reply_to_file = config.reply_to_files.get(i).cloned(); let has_codex = config.has_generated_codex_config; let guard = RespawnGuard::new(i, respawn_tx.clone()); respawn_tasks.spawn(async move { if !delay.is_zero() { tokio::time::sleep(delay).await; } - let result = spawn_and_init(&cmd, &args, &env, has_codex, i, observer).await; + let result = spawn_and_init( + &cmd, + &args, + &env, + reply_to_file.as_deref(), + has_codex, + i, + observer, + ) + .await; guard.send(result); }); } @@ -3608,6 +3646,25 @@ fn dispatch_heartbeat( #[cfg(test)] mod agent_draft_prompt_tests { + use super::spawn_env_for_slot; + + #[test] + fn spawn_env_for_slot_uses_slot_file_and_preserves_persona_env() { + let persona_env = vec![("MODEL".to_string(), "test-model".to_string())]; + let path = std::path::Path::new("/runtime/reply/agent-2"); + + assert_eq!( + spawn_env_for_slot(&persona_env, Some(path)), + vec![ + ("MODEL".to_string(), "test-model".to_string()), + ( + "BUZZ_REPLY_TO_FILE".to_string(), + path.to_string_lossy().into_owned(), + ), + ] + ); + } + #[test] fn shared_base_prompt_teaches_portable_agent_drafts() { let prompt = include_str!("base_prompt.md"); @@ -3685,6 +3742,7 @@ fn spawn_respawn_task( let cmd = config.agent_command.clone(); let args = config.agent_args.clone(); let env = config.persona_env_vars.clone(); + let reply_to_file = config.reply_to_files.get(index).cloned(); let has_codex = config.has_generated_codex_config; let guard = RespawnGuard::new(index, respawn_tx.clone()); respawn_tasks.spawn(async move { @@ -3697,7 +3755,16 @@ fn spawn_respawn_task( tokio::time::sleep(delay).await; } - let result = spawn_and_init(&cmd, &args, &env, has_codex, index, observer).await; + let result = spawn_and_init( + &cmd, + &args, + &env, + reply_to_file.as_deref(), + has_codex, + index, + observer, + ) + .await; guard.send(result); }); @@ -3740,18 +3807,24 @@ struct PoolStartup { command: String, args: Vec, extra_env: Vec<(String, String)>, + reply_to_files: Vec, has_generated_codex_config: bool, model: Option, observer: Option, } impl PoolStartup { - fn from_config(config: &Config, observer: Option) -> Self { + fn from_config( + config: &Config, + observer: Option, + reply_to_files: Vec, + ) -> Self { Self { agents: config.agents, command: config.agent_command.clone(), args: config.agent_args.clone(), extra_env: config.persona_env_vars.clone(), + reply_to_files, has_generated_codex_config: config.has_generated_codex_config, model: config.model.clone(), observer, @@ -3767,10 +3840,15 @@ async fn initialize_agent_pool( // Attempt each spawn under a 60-second timeout; a partial pool is valid. let mut agent_slots: Vec> = Vec::with_capacity(startup.agents as usize); for i in 0..startup.agents as usize { + let reply_to_file = startup + .reply_to_files + .get(i) + .ok_or_else(|| anyhow::anyhow!("missing reply context file for agent slot {i}"))?; + let extra_env = spawn_env_for_slot(&startup.extra_env, Some(reply_to_file)); let spawn_result = AcpClient::spawn( &startup.command, &startup.args, - &startup.extra_env, + &extra_env, startup.has_generated_codex_config, ) .await; @@ -3867,15 +3945,31 @@ async fn initialize_agent_pool( /// /// Takes owned args so it can run in a background `tokio::spawn` task without /// borrowing `Config`. All respawn/refill paths use this. +fn spawn_env_for_slot( + extra_env: &[(String, String)], + reply_to_file: Option<&std::path::Path>, +) -> Vec<(String, String)> { + let mut spawn_env = extra_env.to_vec(); + if let Some(reply_to_file) = reply_to_file { + spawn_env.push(( + "BUZZ_REPLY_TO_FILE".into(), + reply_to_file.to_string_lossy().into_owned(), + )); + } + spawn_env +} + async fn spawn_and_init( command: &str, args: &[String], extra_env: &[(String, String)], + reply_to_file: Option<&std::path::Path>, has_generated_codex_config: bool, agent_index: usize, observer: Option, ) -> Result<(AcpClient, u32, String)> { - let mut acp = AcpClient::spawn(command, args, extra_env, has_generated_codex_config) + let spawn_env = spawn_env_for_slot(extra_env, reply_to_file); + let mut acp = AcpClient::spawn(command, args, &spawn_env, has_generated_codex_config) .await .map_err(|e| anyhow::anyhow!("failed to spawn agent: {e}"))?; acp.set_observer(observer, agent_index); @@ -5013,6 +5107,7 @@ mod build_mcp_servers_tests { respond_to_allowlist: std::collections::HashSet::new(), allowed_respond_to: vec![], persona_env_vars: vec![], + reply_to_files: vec![], has_generated_codex_config: false, relay_observer: false, lazy_pool: false, @@ -5234,6 +5329,7 @@ mod error_outcome_emission_tests { respond_to_allowlist: HashSet::new(), allowed_respond_to: vec![], persona_env_vars: vec![], + reply_to_files: vec![], has_generated_codex_config: false, relay_observer: false, lazy_pool: false, diff --git a/crates/buzz-acp/src/pool.rs b/crates/buzz-acp/src/pool.rs index b1fd68d04..bfa20469e 100644 --- a/crates/buzz-acp/src/pool.rs +++ b/crates/buzz-acp/src/pool.rs @@ -550,6 +550,8 @@ pub struct PromptContext { /// the desktop keys per (agent, relay) pair, e.g. `session_config_captured`, /// mirroring the `managed_agent_runtime_lifecycle` frames. pub relay_url: String, + /// Slot-scoped files used to pass automatic reply anchors to tool processes. + pub reply_to_files: Vec, } impl AgentPool { @@ -859,6 +861,30 @@ async fn resolve_new_session_channel_context( (is_dm, title_channel) } +fn write_reply_anchor(path: &std::path::Path, reply_anchor: Option<&str>) -> std::io::Result<()> { + use std::io::Write; + + let parent = path.parent().ok_or_else(|| { + std::io::Error::new( + std::io::ErrorKind::InvalidInput, + "reply context path has no parent", + ) + })?; + let permissions = std::fs::metadata(path) + .ok() + .map(|metadata| metadata.permissions()); + let mut temp = tempfile::NamedTempFile::new_in(parent)?; + if let Some(anchor) = reply_anchor { + temp.write_all(anchor.as_bytes())?; + } + temp.as_file().sync_all()?; + if let Some(permissions) = permissions { + temp.as_file().set_permissions(permissions)?; + } + temp.persist(path).map_err(|error| error.error)?; + Ok(()) +} + /// Create a new ACP session via `session_new_full()`, populate model capabilities /// on the agent (first session only), and apply `desired_model` if set. /// @@ -1803,6 +1829,27 @@ pub async fn run_prompt_task( // follows as a second block. let mut slash_command: Option = None; let prompt_sections: Vec = if let Some(text) = prompt_text { + if let Some(path) = ctx.reply_to_files.get(agent.index) { + if let Err(error) = write_reply_anchor(path, None) { + tracing::error!( + target: "pool::prompt", + agent = agent.index, + path = %path.display(), + "failed to clear reply context: {error}" + ); + send_prompt_result( + &result_tx, + &turn_id, + agent, + source, + PromptOutcome::Error(AcpError::Protocol(format!( + "failed to clear reply context: {error}" + ))), + requeue_batch_if_queue(&ctx, batch), + ); + return; + } + } // Heartbeats create their session before this point, so a Goose method-not-found // probe has already selected the correct framing for this process. let text = prepend_base_for_legacy( @@ -1829,6 +1876,30 @@ pub async fn run_prompt_task( let profile_lookup = fetch_prompt_profile_lookup(b, conversation_context.as_ref(), &ctx.rest_client).await; + let reply_anchor = + crate::queue::reply_anchor_for_batch(b, channel_info.as_ref(), profile_lookup.as_ref()); + if let Some(path) = ctx.reply_to_files.get(agent.index) { + if let Err(error) = write_reply_anchor(path, reply_anchor.as_deref()) { + tracing::error!( + target: "pool::prompt", + agent = agent.index, + path = %path.display(), + "failed to update reply context: {error}" + ); + send_prompt_result( + &result_tx, + &turn_id, + agent, + source, + PromptOutcome::Error(AcpError::Protocol(format!( + "failed to update reply context: {error}" + ))), + requeue_batch_if_queue(&ctx, batch), + ); + return; + } + } + let known_names: Vec<&str> = profile_lookup .iter() .flat_map(|lookup| lookup.values()) @@ -3734,6 +3805,38 @@ mod tests { // a legacy agent WITH a base_prompt must get [Base] prepended to the user // message. This is the exact regression that shipped in the round-2 bug. + #[test] + fn reply_anchor_rewrite_replaces_and_clears_atomically() { + let dir = tempfile::tempdir().unwrap(); + let path = dir.path().join("agent-0"); + std::fs::write(&path, "stale-anchor").unwrap(); + + write_reply_anchor(&path, Some("current-anchor")).unwrap(); + assert_eq!(std::fs::read_to_string(&path).unwrap(), "current-anchor"); + + write_reply_anchor(&path, None).unwrap(); + assert_eq!(std::fs::read_to_string(&path).unwrap(), ""); + assert_eq!(std::fs::read_dir(dir.path()).unwrap().count(), 1); + } + + #[cfg(unix)] + #[test] + fn reply_anchor_rewrite_preserves_file_permissions() { + use std::os::unix::fs::PermissionsExt; + + let dir = tempfile::tempdir().unwrap(); + let path = dir.path().join("agent-0"); + std::fs::write(&path, "").unwrap(); + std::fs::set_permissions(&path, std::fs::Permissions::from_mode(0o600)).unwrap(); + + write_reply_anchor(&path, Some("current-anchor")).unwrap(); + + assert_eq!( + std::fs::metadata(&path).unwrap().permissions().mode() & 0o777, + 0o600 + ); + } + #[test] fn test_initial_message_legacy_agent_gets_base_prepended() { // protocol_version 1 + Some(base_prompt): [Base] rides along in the @@ -5388,6 +5491,7 @@ mod tests { memory_enabled: false, harness_name: "goose".to_string(), relay_url: "ws://127.0.0.1:3000".to_string(), + reply_to_files: vec![], } } diff --git a/crates/buzz-acp/src/queue.rs b/crates/buzz-acp/src/queue.rs index 029bf86db..2ac336f20 100644 --- a/crates/buzz-acp/src/queue.rs +++ b/crates/buzz-acp/src/queue.rs @@ -1383,6 +1383,33 @@ pub(crate) fn base_section(base_prompt: &str) -> String { format!("[Base]\n{}", base_prompt.trim_end()) } +/// Resolve the automatic reply target for the last event in a batch. +pub(crate) fn reply_anchor_for_batch( + batch: &FlushBatch, + channel_info: Option<&PromptChannelInfo>, + profile_lookup: Option<&PromptProfileLookup>, +) -> Option { + let last_event = batch.events.last()?; + let thread_tags = parse_thread_tags(&last_event.event); + let is_dm = channel_info + .map(|info| info.channel_type == "dm") + .unwrap_or(false); + + if is_dm { + thread_tags + .root_event_id + .is_some() + .then(|| last_event.event.id.to_hex()) + } else { + resolve_reply_anchor( + &last_event.event.pubkey.to_hex(), + &thread_tags, + &last_event.event.id.to_hex(), + profile_lookup, + ) + } +} + /// Format a [`FlushBatch`] into the per-section prompt blocks for the agent. /// /// Produces a stable prompt with these sections (in order): @@ -1464,20 +1491,7 @@ pub fn format_prompt(batch: &FlushBatch, args: &FormatPromptArgs<'_>) -> Vec, } +fn reply_to_from_sources( + explicit: Option, + env_value: Option, + file_path: Option<&std::ffi::OsStr>, +) -> Option { + explicit + .or_else(|| { + env_value.and_then(|value| { + let value = value.trim(); + (!value.is_empty()).then(|| value.to_owned()) + }) + }) + .or_else(|| { + let path = file_path?; + let value = std::fs::read_to_string(path).ok()?; + (!value.trim().is_empty()).then(|| value.trim().to_owned()) + }) +} + +fn automatic_reply_to(explicit: Option) -> Option { + reply_to_from_sources( + explicit, + std::env::var("BUZZ_REPLY_TO").ok(), + std::env::var_os("BUZZ_REPLY_TO_FILE").as_deref(), + ) +} + pub async fn cmd_send_message( client: &BuzzClient, mut p: SendMessageParams, @@ -490,6 +517,7 @@ pub async fn cmd_send_message( // bugs for agent and human users alike. p.content = read_or_stdin(&p.content)?; validate_content_size(&p.content)?; + p.reply_to = automatic_reply_to(p.reply_to); if let Some(ref r) = p.reply_to { validate_hex64(r)?; } @@ -876,7 +904,9 @@ pub async fn dispatch( #[cfg(test)] mod tests { - use super::{find_root_from_tags, match_profiles_by_name, parse_member_pubkeys}; + use super::{ + find_root_from_tags, match_profiles_by_name, parse_member_pubkeys, reply_to_from_sources, + }; use buzz_sdk::mentions::{ extract_at_mentions_with_known, extract_at_names, match_names_to_profiles, MentionProfile, }; @@ -892,6 +922,53 @@ mod tests { const PK_VALID_B: &str = "c6237ef84fa537c78dcee78efd2d4e59f728859c7f194da42ac51ededfa0be05"; const PK_VALID_C: &str = "f4a42a97e594b77bdbd8ee35191c8b28a94a4cb871d96f32921558275421fb68"; + #[test] + fn explicit_reply_to_precedes_environment_and_file() { + let file = tempfile::NamedTempFile::new().unwrap(); + std::fs::write(file.path(), ID_B).unwrap(); + assert_eq!( + reply_to_from_sources( + Some(ID_A.into()), + Some(ID_B.into()), + Some(file.path().as_os_str()), + ) + .as_deref(), + Some(ID_A) + ); + } + + #[test] + fn reply_to_environment_precedes_file() { + let file = tempfile::NamedTempFile::new().unwrap(); + std::fs::write(file.path(), ID_A).unwrap(); + assert_eq!( + reply_to_from_sources( + None, + Some(format!(" {ID_B}\n")), + Some(file.path().as_os_str()), + ) + .as_deref(), + Some(ID_B) + ); + } + + #[test] + fn reply_to_file_is_fallback_and_empty_values_are_ignored() { + let file = tempfile::NamedTempFile::new().unwrap(); + std::fs::write(file.path(), format!(" {ID_A}\n")).unwrap(); + assert_eq!( + reply_to_from_sources(None, Some(" ".into()), Some(file.path().as_os_str())) + .as_deref(), + Some(ID_A) + ); + std::fs::write(file.path(), "").unwrap(); + assert_eq!( + reply_to_from_sources(None, None, Some(file.path().as_os_str())), + None + ); + assert_eq!(reply_to_from_sources(None, None, None), None); + } + #[test] fn root_marker_wins_over_reply_marker() { let tags = json!([