mirror of
https://github.com/block/buzz.git
synced 2026-08-18 06:50:31 +02:00
Add ACP inactivity self-termination
Co-authored-by: npub1mprnacetjua2xx3p5eddmhxyk6wv929ymm5py8kd2xfxurxahspqqlgyta <d8473ee32b973aa31a21a65adddcc4b69cc2a8a4dee8121ecd51926e0cddbc02@buzz.block.builderlab.xyz> Signed-off-by: npub1mprnacetjua2xx3p5eddmhxyk6wv929ymm5py8kd2xfxurxahspqqlgyta <d8473ee32b973aa31a21a65adddcc4b69cc2a8a4dee8121ecd51926e0cddbc02@buzz.block.builderlab.xyz>
This commit is contained in:
parent
28ae6cd217
commit
8d68063a31
@@ -474,6 +474,11 @@ pub struct CliArgs {
|
||||
#[arg(long, env = "BUZZ_ACP_RELAY_OBSERVER", default_value_t = false)]
|
||||
pub relay_observer: bool,
|
||||
|
||||
/// Exit after this many seconds with no dispatched events and no turn in flight.
|
||||
/// 0 disables inactivity self-termination.
|
||||
#[arg(long, env = "BUZZ_ACP_EXIT_AFTER_INACTIVITY", default_value_t = 0)]
|
||||
pub exit_after_inactivity: u64,
|
||||
|
||||
/// Connect and subscribe before starting the ACP/LLM subprocess pool.
|
||||
#[arg(long, env = "BUZZ_ACP_LAZY_POOL", default_value_t = false)]
|
||||
pub lazy_pool: bool,
|
||||
@@ -550,6 +555,8 @@ pub struct Config {
|
||||
pub has_generated_codex_config: bool,
|
||||
/// Whether to publish encrypted observer frames through the relay.
|
||||
pub relay_observer: bool,
|
||||
/// Seconds without dispatched events before an idle harness exits. 0 = disabled.
|
||||
pub exit_after_inactivity_secs: u64,
|
||||
/// Whether ACP/LLM subprocess initialization is deferred until accepted work arrives.
|
||||
pub lazy_pool: bool,
|
||||
/// Agent owner pubkey (hex). Used for `--respond-to=owner-only` gate.
|
||||
@@ -1098,6 +1105,7 @@ impl Config {
|
||||
persona_env_vars,
|
||||
has_generated_codex_config,
|
||||
relay_observer: args.relay_observer,
|
||||
exit_after_inactivity_secs: args.exit_after_inactivity,
|
||||
lazy_pool: args.lazy_pool,
|
||||
agent_owner: args.agent_owner.map(|s| s.trim().to_ascii_lowercase()),
|
||||
no_base_prompt: args.no_base_prompt,
|
||||
@@ -1468,6 +1476,7 @@ mod tests {
|
||||
persona_env_vars: vec![],
|
||||
has_generated_codex_config: false,
|
||||
relay_observer: false,
|
||||
exit_after_inactivity_secs: 0,
|
||||
lazy_pool: false,
|
||||
agent_owner: None,
|
||||
no_base_prompt: false,
|
||||
@@ -2167,6 +2176,22 @@ channels = "ALL"
|
||||
assert!(err.to_string().contains("turn liveness interval must be 0"));
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn inactivity_exit_defaults_disabled_and_accepts_cli_value() {
|
||||
let key = "0".repeat(64);
|
||||
let default = CliArgs::parse_from(["buzz-acp", "--private-key", &key]);
|
||||
assert_eq!(default.exit_after_inactivity, 0);
|
||||
|
||||
let configured = CliArgs::parse_from([
|
||||
"buzz-acp",
|
||||
"--private-key",
|
||||
&key,
|
||||
"--exit-after-inactivity",
|
||||
"120",
|
||||
]);
|
||||
assert_eq!(configured.exit_after_inactivity, 120);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn lazy_pool_defaults_off() {
|
||||
let key = "0".repeat(64);
|
||||
|
||||
+111
-8
@@ -1228,6 +1228,59 @@ impl Drop for RespawnGuard {
|
||||
// sync entry point — `std::env::set_var` is only safe before tokio spawns
|
||||
// worker threads (Rust 2024 edition safety requirement).
|
||||
|
||||
fn inactivity_expired(
|
||||
last_activity: tokio::time::Instant,
|
||||
now: tokio::time::Instant,
|
||||
bound: Duration,
|
||||
turn_in_flight: bool,
|
||||
) -> bool {
|
||||
!bound.is_zero() && !turn_in_flight && now.duration_since(last_activity) >= bound
|
||||
}
|
||||
|
||||
#[cfg(test)]
|
||||
mod inactivity_tests {
|
||||
use super::*;
|
||||
|
||||
#[test]
|
||||
fn zero_disables_expiry_and_in_flight_turns_defer_it() {
|
||||
let started = tokio::time::Instant::now();
|
||||
let after_bound = started + Duration::from_secs(61);
|
||||
|
||||
assert!(!inactivity_expired(
|
||||
started,
|
||||
after_bound,
|
||||
Duration::ZERO,
|
||||
false
|
||||
));
|
||||
assert!(!inactivity_expired(
|
||||
started,
|
||||
after_bound,
|
||||
Duration::from_secs(60),
|
||||
true
|
||||
));
|
||||
assert!(inactivity_expired(
|
||||
started,
|
||||
after_bound,
|
||||
Duration::from_secs(60),
|
||||
false
|
||||
));
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn dispatched_activity_restarts_the_inactivity_bound() {
|
||||
let started = tokio::time::Instant::now();
|
||||
let dispatched = started + Duration::from_secs(50);
|
||||
let checked = started + Duration::from_secs(61);
|
||||
|
||||
assert!(!inactivity_expired(
|
||||
dispatched,
|
||||
checked,
|
||||
Duration::from_secs(60),
|
||||
false
|
||||
));
|
||||
}
|
||||
}
|
||||
|
||||
pub fn run() -> Result<()> {
|
||||
config::propagate_legacy_env_vars();
|
||||
tokio_main()
|
||||
@@ -1601,6 +1654,21 @@ async fn tokio_main() -> Result<()> {
|
||||
let mut typing_channels: HashMap<Uuid, ThreadTags> = HashMap::new();
|
||||
let mut presence_task: Option<tokio::task::JoinHandle<()>> = None;
|
||||
|
||||
// Independent of pool readiness: a never-mentioned lazy agent must still
|
||||
// self-terminate. The watch interval is capped so small configured bounds
|
||||
// remain reasonably precise without waking long-lived agents frequently.
|
||||
let inactivity_bound = Duration::from_secs(config.exit_after_inactivity_secs);
|
||||
let mut last_activity = tokio::time::Instant::now();
|
||||
let mut inactivity_reaper = if inactivity_bound.is_zero() {
|
||||
None
|
||||
} else {
|
||||
let interval = inactivity_bound.min(Duration::from_secs(30));
|
||||
Some(tokio::time::interval_at(
|
||||
tokio::time::Instant::now() + interval,
|
||||
interval,
|
||||
))
|
||||
};
|
||||
|
||||
// Runs at the TOP of every loop iteration via Instant check — cannot be
|
||||
// starved by the biased select. Slot refill spawns background tasks so
|
||||
// spawn_and_init never blocks the main loop.
|
||||
@@ -1774,7 +1842,9 @@ async fn tokio_main() -> Result<()> {
|
||||
// called on relay events or pool results, neither of which
|
||||
// arrive when the channel is silent.
|
||||
if queue.has_flushable_work() {
|
||||
for (channel_id, thread_tags) in dispatch_pending(&mut pool, &mut queue, &ctx) {
|
||||
for (channel_id, thread_tags) in
|
||||
dispatch_pending(&mut pool, &mut queue, &ctx, &mut last_activity)
|
||||
{
|
||||
typing_channels.insert(channel_id, thread_tags);
|
||||
}
|
||||
}
|
||||
@@ -1810,7 +1880,9 @@ async fn tokio_main() -> Result<()> {
|
||||
// this, batches requeued during crash recovery sit idle until the
|
||||
// next relay event arrives — which can be minutes on quiet channels.
|
||||
if respawn_collected {
|
||||
for (channel_id, thread_tags) in dispatch_pending(&mut pool, &mut queue, &ctx) {
|
||||
for (channel_id, thread_tags) in
|
||||
dispatch_pending(&mut pool, &mut queue, &ctx, &mut last_activity)
|
||||
{
|
||||
typing_channels.insert(channel_id, thread_tags);
|
||||
}
|
||||
}
|
||||
@@ -2258,7 +2330,7 @@ async fn tokio_main() -> Result<()> {
|
||||
}
|
||||
if pool_ready {
|
||||
for (channel_id, thread_tags) in
|
||||
dispatch_pending(&mut pool, &mut queue, &ctx)
|
||||
dispatch_pending(&mut pool, &mut queue, &ctx, &mut last_activity)
|
||||
{
|
||||
typing_channels.insert(channel_id, thread_tags);
|
||||
}
|
||||
@@ -2275,6 +2347,27 @@ async fn tokio_main() -> Result<()> {
|
||||
}
|
||||
None
|
||||
}
|
||||
_ = async {
|
||||
match inactivity_reaper.as_mut() {
|
||||
Some(timer) => timer.tick().await,
|
||||
None => std::future::pending().await,
|
||||
}
|
||||
} => {
|
||||
let _ = result_rx;
|
||||
if inactivity_expired(
|
||||
last_activity,
|
||||
tokio::time::Instant::now(),
|
||||
inactivity_bound,
|
||||
queue.has_in_flight() || heartbeat_in_flight,
|
||||
) {
|
||||
tracing::info!(
|
||||
inactivity_seconds = config.exit_after_inactivity_secs,
|
||||
"inactivity bound reached — exiting gracefully"
|
||||
);
|
||||
let _ = shutdown_tx.send(());
|
||||
}
|
||||
None
|
||||
}
|
||||
_ = async {
|
||||
match heartbeat.as_mut() {
|
||||
Some(hb) => hb.tick().await,
|
||||
@@ -2287,7 +2380,7 @@ async fn tokio_main() -> Result<()> {
|
||||
} else if queue.has_flushable_work() {
|
||||
tracing::debug!("heartbeat_skipped_events");
|
||||
for (channel_id, thread_tags) in
|
||||
dispatch_pending(&mut pool, &mut queue, &ctx)
|
||||
dispatch_pending(&mut pool, &mut queue, &ctx, &mut last_activity)
|
||||
{
|
||||
typing_channels.insert(channel_id, thread_tags);
|
||||
}
|
||||
@@ -2385,7 +2478,9 @@ async fn tokio_main() -> Result<()> {
|
||||
{
|
||||
break;
|
||||
}
|
||||
for (channel_id, thread_tags) in dispatch_pending(&mut pool, &mut queue, &ctx) {
|
||||
for (channel_id, thread_tags) in
|
||||
dispatch_pending(&mut pool, &mut queue, &ctx, &mut last_activity)
|
||||
{
|
||||
typing_channels.insert(channel_id, thread_tags);
|
||||
}
|
||||
}
|
||||
@@ -2408,7 +2503,9 @@ async fn tokio_main() -> Result<()> {
|
||||
tracing::error!("all agents dead — exiting");
|
||||
break;
|
||||
}
|
||||
for (channel_id, thread_tags) in dispatch_pending(&mut pool, &mut queue, &ctx) {
|
||||
for (channel_id, thread_tags) in
|
||||
dispatch_pending(&mut pool, &mut queue, &ctx, &mut last_activity)
|
||||
{
|
||||
typing_channels.insert(channel_id, thread_tags);
|
||||
}
|
||||
}
|
||||
@@ -2550,7 +2647,9 @@ async fn tokio_main() -> Result<()> {
|
||||
// tear down the in-flight task; on its completion the
|
||||
// queue drains. We still try here in case the in-flight
|
||||
// task has already returned.
|
||||
for (channel_id, thread_tags) in dispatch_pending(&mut pool, &mut queue, &ctx) {
|
||||
for (channel_id, thread_tags) in
|
||||
dispatch_pending(&mut pool, &mut queue, &ctx, &mut last_activity)
|
||||
{
|
||||
typing_channels.insert(channel_id, thread_tags);
|
||||
}
|
||||
}
|
||||
@@ -2577,7 +2676,7 @@ async fn tokio_main() -> Result<()> {
|
||||
None,
|
||||
);
|
||||
for (channel_id, thread_tags) in
|
||||
dispatch_pending(&mut pool, &mut queue, &ctx)
|
||||
dispatch_pending(&mut pool, &mut queue, &ctx, &mut last_activity)
|
||||
{
|
||||
typing_channels.insert(channel_id, thread_tags);
|
||||
}
|
||||
@@ -2911,6 +3010,7 @@ fn dispatch_pending(
|
||||
pool: &mut AgentPool,
|
||||
queue: &mut EventQueue,
|
||||
ctx: &Arc<PromptContext>,
|
||||
last_activity: &mut tokio::time::Instant,
|
||||
) -> Vec<(Uuid, ThreadTags)> {
|
||||
let mut dispatched_channels = Vec::new();
|
||||
loop {
|
||||
@@ -2990,6 +3090,7 @@ fn dispatch_pending(
|
||||
},
|
||||
);
|
||||
dispatched_channels.push((channel_id, typing_scope));
|
||||
*last_activity = tokio::time::Instant::now();
|
||||
}
|
||||
tracing::debug!(
|
||||
dispatched = dispatched_channels.len(),
|
||||
@@ -5031,6 +5132,7 @@ mod build_mcp_servers_tests {
|
||||
persona_env_vars: vec![],
|
||||
has_generated_codex_config: false,
|
||||
relay_observer: false,
|
||||
exit_after_inactivity_secs: 0,
|
||||
lazy_pool: false,
|
||||
agent_owner: None,
|
||||
no_base_prompt: false,
|
||||
@@ -5252,6 +5354,7 @@ mod error_outcome_emission_tests {
|
||||
persona_env_vars: vec![],
|
||||
has_generated_codex_config: false,
|
||||
relay_observer: false,
|
||||
exit_after_inactivity_secs: 0,
|
||||
lazy_pool: false,
|
||||
agent_owner: None,
|
||||
no_base_prompt: false,
|
||||
|
||||
@@ -646,6 +646,11 @@ impl EventQueue {
|
||||
self.in_flight_channels.contains(&channel_id)
|
||||
}
|
||||
|
||||
/// Whether any channel currently has a turn in flight.
|
||||
pub fn has_in_flight(&self) -> bool {
|
||||
!self.in_flight_channels.is_empty()
|
||||
}
|
||||
|
||||
// ── Goose-native steer withhold (side table) ──────────────────────────
|
||||
//
|
||||
// While a goose-native `_goose/unstable/session/steer` write is in flight
|
||||
|
||||
@@ -77,6 +77,10 @@ pub(crate) const RESERVED_ENV_KEYS: &[&str] = &[
|
||||
"BUZZ_ACP_RESPOND_TO",
|
||||
"BUZZ_ACP_RESPOND_TO_ALLOWLIST",
|
||||
"BUZZ_ACP_AGENT_OWNER",
|
||||
// Remote lifetime/presence policy: user env must not disable the
|
||||
// desktop/provider-owned bounds while the saved record still promises them.
|
||||
"BUZZ_ACP_EXIT_AFTER_INACTIVITY",
|
||||
"BUZZ_ACP_NO_PRESENCE",
|
||||
// Readiness handoff: desktop is the ONLY readiness source. A saved or
|
||||
// ambient env var must not be able to forge setup mode (NotReady) on a
|
||||
// Ready agent or suppress it (empty/stale payload) on a NotReady one.
|
||||
|
||||
@@ -158,6 +158,15 @@ fn reserved_keys_include_respond_to_gate() {
|
||||
}
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn reserved_keys_include_remote_lifetime_policy() {
|
||||
for key in ["BUZZ_ACP_EXIT_AFTER_INACTIVITY", "BUZZ_ACP_NO_PRESENCE"] {
|
||||
assert!(is_reserved_env_key(key), "{key} should be reserved");
|
||||
let agent = map(&[(key, "0")]);
|
||||
assert!(merged_user_env(&BTreeMap::new(), &agent).is_empty());
|
||||
}
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn reserved_keys_include_code_execution_surface() {
|
||||
// The agent/MCP command + args are what Buzz actually exec's.
|
||||
|
||||
@@ -0,0 +1,31 @@
|
||||
import assert from "node:assert/strict";
|
||||
import { describe, it } from "node:test";
|
||||
|
||||
import { coerceConfigValues } from "./ProviderConfigFields.tsx";
|
||||
|
||||
const schema = {
|
||||
properties: {
|
||||
inactivity_seconds: { type: "integer" },
|
||||
threshold: { type: "number" },
|
||||
label: { type: "string" },
|
||||
},
|
||||
};
|
||||
|
||||
describe("coerceConfigValues", () => {
|
||||
it("omits cleared numeric fields without losing explicit zero", () => {
|
||||
assert.deepEqual(
|
||||
coerceConfigValues(
|
||||
{ inactivity_seconds: "", threshold: "0", label: "" },
|
||||
schema,
|
||||
),
|
||||
{ threshold: 0, label: "" },
|
||||
);
|
||||
});
|
||||
|
||||
it("preserves nonempty invalid numeric input for provider validation", () => {
|
||||
assert.deepEqual(
|
||||
coerceConfigValues({ inactivity_seconds: "not-a-number" }, schema),
|
||||
{ inactivity_seconds: "not-a-number" },
|
||||
);
|
||||
});
|
||||
});
|
||||
@@ -14,7 +14,8 @@ export function coerceConfigValues(
|
||||
for (const [key, value] of Object.entries(config)) {
|
||||
const prop = properties[key] as Record<string, unknown> | undefined;
|
||||
const schemaType = prop?.type;
|
||||
if ((schemaType === "integer" || schemaType === "number") && value !== "") {
|
||||
if (schemaType === "integer" || schemaType === "number") {
|
||||
if (value === "") continue;
|
||||
const num = Number(value);
|
||||
result[key] = Number.isNaN(num) ? value : num;
|
||||
} else if (schemaType === "boolean") {
|
||||
|
||||
Reference in New Issue
Block a user