feat(acp): propagate automatic reply context

Co-authored-by: npub1x4hk035p3p9q39a3fcrd2fe30lpkrhr5dwe0cqzzjphxyyh8m0gsq4vqap <356f67c681884a0897b14e06d527317fc361dc746bb2fc0042906e6212e7dbd1@buzz.block.builderlab.xyz>
Signed-off-by: npub1x4hk035p3p9q39a3fcrd2fe30lpkrhr5dwe0cqzzjphxyyh8m0gsq4vqap <356f67c681884a0897b14e06d527317fc361dc746bb2fc0042906e6212e7dbd1@buzz.block.builderlab.xyz>
(cherry picked from commit 1535e6e5b3a42d0bca12b18627c83ae1164f3275)
Signed-off-by: Michael Neale <michael.neale@gmail.com>
This commit is contained in:
npub1x4hk035p3p9q39a3fcrd2fe30lpkrhr5dwe0cqzzjphxyyh8m0gsq4vqap
2026-07-29 17:01:53 +10:00
committed by Michael Neale
parent 9b09d1ba26
commit 3fcd7acef9
9 changed files with 326 additions and 24 deletions
Generated
+1
View File
@@ -795,6 +795,7 @@ dependencies = [
"serde",
"serde_json",
"sha2 0.11.0",
"tempfile",
"thiserror 2.0.18",
"tokio",
"tokio-tungstenite 0.29.0",
+3
View File
@@ -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]
+2 -1
View File
@@ -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);
}
}
+5
View File
@@ -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<PathBuf>,
/// 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,
+104 -8
View File
@@ -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<std::path::PathBuf> = (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<AgentPool> = 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<String>,
extra_env: Vec<(String, String)>,
reply_to_files: Vec<std::path::PathBuf>,
has_generated_codex_config: bool,
model: Option<String>,
observer: Option<observer::ObserverHandle>,
}
impl PoolStartup {
fn from_config(config: &Config, observer: Option<observer::ObserverHandle>) -> Self {
fn from_config(
config: &Config,
observer: Option<observer::ObserverHandle>,
reply_to_files: Vec<std::path::PathBuf>,
) -> 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<Option<OwnedAgent>> = 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<observer::ObserverHandle>,
) -> 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,
+104
View File
@@ -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<std::path::PathBuf>,
}
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<String> = None;
let prompt_sections: Vec<String> = 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![],
}
}
+28 -14
View File
@@ -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<String> {
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<Str
// - top-level → anchor to the triggering event (it becomes the root)
// Agent↔agent turns get no forced anchor — deep nesting is intentional
// there. DMs are always 1:1 with a human, so they always anchor.
let sender_pubkey = last_event.event.pubkey.to_hex();
let reply_anchor = if is_dm {
thread_tags
.root_event_id
.is_some()
.then(|| last_event.event.id.to_hex())
} else {
resolve_reply_anchor(
&sender_pubkey,
&thread_tags,
&last_event.event.id.to_hex(),
args.profile_lookup,
)
};
let reply_anchor = reply_anchor_for_batch(batch, args.channel_info, args.profile_lookup);
sections.push(format_context_hints(
batch.channel_id,
args.channel_info,
+1
View File
@@ -61,6 +61,7 @@ const PASSTHROUGH_ENV: &[&str] = &[
"BUZZ_PRIVATE_KEY",
"BUZZ_RELAY_URL",
"BUZZ_AUTH_TAG",
"BUZZ_REPLY_TO_FILE",
// Agent display name — dev-mcp uses it as the git author name. On the
// Desktop path this arrives via the wire `mcpServers[].env` declaration
// (which wins here anyway); the allowlist entry covers ACP clients that
+78 -1
View File
@@ -480,6 +480,33 @@ pub struct SendMessageParams {
pub files: Vec<String>,
}
fn reply_to_from_sources(
explicit: Option<String>,
env_value: Option<String>,
file_path: Option<&std::ffi::OsStr>,
) -> Option<String> {
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<String>) -> Option<String> {
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!([