From 1c0a4c838bfb21e6fc7413f868ff12182a4dcb7b Mon Sep 17 00:00:00 2001 From: npub12gtutshhh76rx0jx697f32f9tffd4hhp3hx58fp4x6u4uemkm7sqf8f757 <5217c5c2f7bfb4333e46d17c98a9255a52dadee18dcd43a43536b95e6776dfa0@sprout-oss.stage.blox.sqprod.co> Date: Sun, 19 Jul 2026 10:45:59 -0400 Subject: [PATCH] feat(acp): model lazy pool lifecycle Co-authored-by: npub12gtutshhh76rx0jx697f32f9tffd4hhp3hx58fp4x6u4uemkm7sqf8f757 <5217c5c2f7bfb4333e46d17c98a9255a52dadee18dcd43a43536b95e6776dfa0@sprout-oss.stage.blox.sqprod.co> Signed-off-by: npub12gtutshhh76rx0jx697f32f9tffd4hhp3hx58fp4x6u4uemkm7sqf8f757 <5217c5c2f7bfb4333e46d17c98a9255a52dadee18dcd43a43536b95e6776dfa0@sprout-oss.stage.blox.sqprod.co> Co-authored-by: Tyler Longwell Signed-off-by: Tyler Longwell --- TESTING.md | 8 + crates/buzz-acp/README.md | 1 + crates/buzz-acp/src/config.rs | 20 + crates/buzz-acp/src/lib.rs | 414 +++++++++++++----- crates/buzz-acp/src/pool_lifecycle.rs | 312 +++++++++++++ crates/buzz-acp/tests/pool_lifecycle_state.rs | 4 + 6 files changed, 658 insertions(+), 101 deletions(-) create mode 100644 crates/buzz-acp/src/pool_lifecycle.rs create mode 100644 crates/buzz-acp/tests/pool_lifecycle_state.rs diff --git a/TESTING.md b/TESTING.md index ad5225244..ee225c738 100644 --- a/TESTING.md +++ b/TESTING.md @@ -235,6 +235,14 @@ The justfile also ships `just goose key="$AGENT_NSEC"` (foreground) and same env. See `crates/buzz-acp/README.md` for parallel agents, heartbeats, respond-to gates, and forum subscriptions. +To exercise deferred ACP startup, add `BUZZ_ACP_LAZY_POOL=true` before launching +`buzz-acp`. The harness should connect, authenticate, subscribe, and publish +online presence without starting the configured ACP child. The first accepted, +flushable mention should start exactly one child and then dispatch the queued +message. Automated coverage in `pool_lifecycle_state` pins single-wake, +retry/backoff, and stale-result behavior; it does not replace this real +relay/process smoke test. + Send the agent a task — switch your shell back to the **sender** identity from step 4 and @mention the agent: diff --git a/crates/buzz-acp/README.md b/crates/buzz-acp/README.md index b06891dd8..c95bdb8b3 100644 --- a/crates/buzz-acp/README.md +++ b/crates/buzz-acp/README.md @@ -115,6 +115,7 @@ All configuration is via environment variables (or CLI flags — every env var h | Flag | Env Var | Default | Description | |------|---------|---------|-------------| | `--agents` | `BUZZ_ACP_AGENTS` | `1` | Number of agent subprocesses (1–32). | +| `--lazy-pool` | `BUZZ_ACP_LAZY_POOL` | `false` | Connect, subscribe, and queue accepted work before starting ACP/LLM subprocesses. The first accepted event wakes one pool initialization task; failures retry with bounded exponential backoff while work remains. | | `--heartbeat-interval` | `BUZZ_ACP_HEARTBEAT_INTERVAL` | `0` | Seconds between heartbeat prompts. `0` = disabled. Must be `0` or ≥10 when enabled. | | `--heartbeat-prompt` | `BUZZ_ACP_HEARTBEAT_PROMPT` | (built-in) | Custom heartbeat prompt text. Conflicts with `--heartbeat-prompt-file`. | | `--heartbeat-prompt-file` | `BUZZ_ACP_HEARTBEAT_PROMPT_FILE` | — | Read heartbeat prompt from a file. Conflicts with `--heartbeat-prompt`. | diff --git a/crates/buzz-acp/src/config.rs b/crates/buzz-acp/src/config.rs index befb7aa6a..fd5fd4dd3 100644 --- a/crates/buzz-acp/src/config.rs +++ b/crates/buzz-acp/src/config.rs @@ -467,6 +467,10 @@ pub struct CliArgs { /// Publish encrypted ACP observer frames over the relay. #[arg(long, env = "BUZZ_ACP_RELAY_OBSERVER", default_value_t = false)] pub relay_observer: bool, + + /// 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, } /// Merged NIP-01 subscription filter for a single channel. @@ -537,6 +541,8 @@ pub struct Config { pub has_generated_codex_config: bool, /// Whether to publish encrypted observer frames through the relay. pub relay_observer: bool, + /// 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. /// Replaces the old REST-based owner lookup. pub agent_owner: Option, @@ -996,6 +1002,7 @@ impl Config { persona_env_vars, has_generated_codex_config, relay_observer: args.relay_observer, + lazy_pool: args.lazy_pool, agent_owner: args.agent_owner.map(|s| s.trim().to_ascii_lowercase()), no_base_prompt: args.no_base_prompt, base_prompt_content, @@ -1368,6 +1375,7 @@ mod tests { persona_env_vars: vec![], has_generated_codex_config: false, relay_observer: false, + lazy_pool: false, agent_owner: None, no_base_prompt: false, base_prompt_content: None, @@ -2033,6 +2041,18 @@ channels = "ALL" assert!(err.to_string().contains("turn liveness interval must be 0")); } + #[test] + fn lazy_pool_defaults_off() { + assert!(!CliArgs::parse_from(["buzz-acp"]).lazy_pool); + } + + #[test] + fn lazy_pool_cli_flag_enables_deferred_startup() { + let args = CliArgs::try_parse_from(["buzz-acp", "--lazy-pool=true"]); + assert!(args.is_err(), "bool flags do not take an explicit value"); + assert!(CliArgs::parse_from(["buzz-acp", "--lazy-pool"]).lazy_pool); + } + #[test] fn test_summary_includes_agents_and_heartbeat() { let config = test_config(SubscribeMode::Mentions); diff --git a/crates/buzz-acp/src/lib.rs b/crates/buzz-acp/src/lib.rs index 4d155e218..83df90c7a 100644 --- a/crates/buzz-acp/src/lib.rs +++ b/crates/buzz-acp/src/lib.rs @@ -6,6 +6,7 @@ mod engram_fetch; mod filter; mod observer; mod pool; +mod pool_lifecycle; mod queue; mod relay; mod setup_mode; @@ -39,6 +40,7 @@ use pool::{ AgentPool, ControlSignal, IdleSwitchResult, OwnedAgent, PromptContext, PromptOutcome, PromptResult, PromptSource, SessionState, TimeoutKind, }; +use pool_lifecycle::PoolLifecycle; use queue::{CancelReason, EventQueue, FlushBatch, QueuedEvent, ThreadTags}; use relay::{HarnessRelay, RelayEventPublisher}; use tokio::sync::{mpsc, watch}; @@ -88,6 +90,28 @@ async fn publish_presence( Ok(()) } +fn emit_runtime_lifecycle( + observer: Option<&observer::ObserverHandle>, + pubkey: &str, + relay_url: &str, + lifecycle: &str, + error: Option<&str>, +) { + if let Some(observer) = observer { + observer.emit( + "managed_agent_runtime_lifecycle", + None, + &observer::ObserverContext::default(), + serde_json::json!({ + "pubkey": pubkey, + "relayUrl": relay_url, + "lifecycle": lifecycle, + "error": error, + }), + ); + } +} + /// Resolve the agent's owner pubkey at startup. /// /// Priority: @@ -1157,96 +1181,13 @@ async fn tokio_main() -> Result<()> { ); } - // - // Finding #10: one agent failing to start must not kill the whole pool. - // We attempt each spawn under a 60-second timeout; failures are logged and - // skipped. If ALL agents fail we return an error. A partial pool is valid — - // the harness continues with reduced capacity and logs a warning. - let mut agent_slots: Vec> = Vec::with_capacity(config.agents as usize); - for i in 0..config.agents as usize { - // Spawn OUTSIDE the timeout so we always own the child for cleanup. - // This matches the run_models pattern and prevents zombie leaks on - // init timeout (the cancelled future would drop the AcpClient via - // Drop which is best-effort only). - let spawn_result = AcpClient::spawn( - &config.agent_command, - &config.agent_args, - &config.persona_env_vars, - config.has_generated_codex_config, - ) - .await; - match spawn_result { - Ok(mut acp) => { - acp.set_observer(observer.clone(), i); - match tokio::time::timeout(Duration::from_secs(60), acp.initialize()).await { - Ok(Ok(init_result)) => { - tracing::info!(agent = i, "agent initialized: {init_result}"); - let protocol_version = - init_result["protocolVersion"].as_u64().unwrap_or(1) as u32; - tracing::info!( - agent = i, - name = init_result - .get("agentInfo") - .or_else(|| init_result.get("serverInfo")) - .and_then(|info| info.get("name")) - .and_then(|v| v.as_str()) - .unwrap_or("unknown"), - "agent initialized — non-cancelling steer enabled (try-and-tolerate)" - ); - acp.observe( - "agent_initialized", - serde_json::json!({ - "agentIndex": i, - "initializeResult": init_result, - }), - ); - let agent_name = normalized_agent_name(&init_result); - agent_slots.push(Some(OwnedAgent { - index: i, - acp, - state: SessionState::default(), - model_capabilities: None, - desired_model: config.model.clone(), - model_overridden: false, - agent_name, - goose_system_prompt_supported: None, - protocol_version, - })); - } - Ok(Err(e)) => { - tracing::error!(agent = i, "agent initialize failed: {e}"); - acp.shutdown().await; - agent_slots.push(None); - } - Err(_) => { - tracing::error!(agent = i, "agent timed out during init (60s)"); - acp.shutdown().await; - agent_slots.push(None); - } - } - } - Err(e) => { - tracing::error!(agent = i, "agent failed to spawn: {e}"); - agent_slots.push(None); - } - } - } - let live_count = agent_slots.iter().filter(|s| s.is_some()).count(); - if live_count == 0 { - return Err(anyhow::anyhow!( - "all {} agents failed to start — cannot continue", - config.agents - )); - } - if live_count < config.agents as usize { - tracing::warn!( - "started {}/{} agents — continuing with reduced pool", - live_count, - config.agents - ); - } - tracing::info!("agent_pool_ready agents={}", live_count); - let mut pool = AgentPool::from_slots(agent_slots); + 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? + }; + let mut pool_ready = !config.lazy_pool; + let mut pool_lifecycle: PoolLifecycle = PoolLifecycle::listening(); // // Finding #22: capture a startup watermark BEFORE connecting to the relay. @@ -1428,6 +1369,10 @@ async fn tokio_main() -> Result<()> { )); } + let dedup_mode = config.dedup_mode; + let mut queue = + EventQueue::new(dedup_mode).with_in_flight_deadline(config.max_turn_duration_secs); + // Online means the harness can receive work, not merely that its socket is // connected. Publishing after channel subscriptions gives desktop callers // a durable readiness boundary before they send a startup mention. @@ -1438,9 +1383,15 @@ async fn tokio_main() -> Result<()> { } } - let dedup_mode = config.dedup_mode; - let mut queue = - EventQueue::new(dedup_mode).with_in_flight_deadline(config.max_turn_duration_secs); + if config.lazy_pool { + emit_runtime_lifecycle( + observer.as_ref(), + &pubkey_hex, + &config.relay_url, + "listening", + None, + ); + } let base_prompt_content = config.base_prompt_content.take(); let ctx = Arc::new(PromptContext { @@ -1528,6 +1479,8 @@ async fn tokio_main() -> Result<()> { let (respawn_tx, mut respawn_rx) = mpsc::channel::(config.agents as usize); // JoinSet for respawn tasks so shutdown can abort them. let mut respawn_tasks: tokio::task::JoinSet<()> = tokio::task::JoinSet::new(); + let (wake_tx, mut wake_rx) = mpsc::channel::<(u32, Result)>(1); + let mut wake_tasks: tokio::task::JoinSet<()> = tokio::task::JoinSet::new(); // Channel for non-cancelling steer ack watchers to forward outcomes back // to the main loop. Each `pool.send_steer(...) == Ok(())` spawns a @@ -1614,10 +1567,39 @@ async fn tokio_main() -> Result<()> { Result(Box), Panic(tokio::task::JoinError), SteerAck(SteerAckEvent), + Wake(u32, Result), } loop { - if last_maintenance.elapsed() >= maintenance_interval { + if config.lazy_pool && !pool_ready { + if let Some(attempt) = pool_lifecycle + .start_wake_if_due(queue.has_flushable_work(), tokio::time::Instant::now()) + { + emit_runtime_lifecycle( + observer.as_ref(), + &pubkey_hex, + &config.relay_url, + "waking", + None, + ); + let startup = PoolStartup::from_config(&config, observer.clone()); + let wake_tx = wake_tx.clone(); + let wake_shutdown = shutdown_rx.clone(); + wake_tasks.spawn(async move { + let result = initialize_agent_pool(&startup, Some(wake_shutdown)) + .await + .map_err(|error| error.to_string()); + if let Err(error) = wake_tx.send((attempt, result)).await { + let (_attempt, result) = error.0; + if let Ok(mut abandoned_pool) = result { + shutdown_agent_pool(&mut abandoned_pool).await; + } + } + }); + } + } + + if pool_ready && last_maintenance.elapsed() >= maintenance_interval { last_maintenance = std::time::Instant::now(); queue.compact_expired_state(); @@ -1699,7 +1681,7 @@ async fn tokio_main() -> Result<()> { biased; // Finding #24: recv() returning None means all senders dropped // (pool was torn down). Break cleanly instead of panicking. - r = result_rx.recv() => match r { + r = result_rx.recv(), if pool_ready => match r { Some(result) => Some(PoolEvent::Result(Box::new(result))), None => { tracing::info!("result channel closed — exiting main loop"); @@ -1720,6 +1702,34 @@ async fn tokio_main() -> Result<()> { Some(ack_event) = steer_ack_rx.recv() => { Some(PoolEvent::SteerAck(ack_event)) } + Some((attempt, result)) = wake_rx.recv(), if config.lazy_pool && !pool_ready => { + Some(PoolEvent::Wake(attempt, result)) + } + _ = async { + match pool_lifecycle.retry_at() { + Some(retry_at) if !pool_ready => tokio::time::sleep_until(retry_at).await, + _ => std::future::pending().await, + } + } => None, + Some(Err(error)) = wake_tasks.join_next(), if !wake_tasks.is_empty() => { + if let Some(attempt) = pool_lifecycle.waking_attempt() { + let message = format!("pool wake task failed: {error}"); + if pool_lifecycle.cancel_wake( + attempt, + message.clone(), + tokio::time::Instant::now(), + ) { + emit_runtime_lifecycle( + observer.as_ref(), + &pubkey_hex, + &config.relay_url, + "failed", + Some(&message), + ); + } + } + None + } control_event = async { match relay_observer_control_rx.as_mut() { Some(rx) => rx.recv().await, @@ -1829,7 +1839,11 @@ async fn tokio_main() -> Result<()> { // complete normally (the relay may reject actions if // the agent lost access). let drained_ids = queue.drain_channel(ch); - let invalidated = pool.invalidate_channel_sessions(ch); + let invalidated = if pool_ready { + pool.invalidate_channel_sessions(ch) + } else { + 0 + }; // Track removed channels so checked-out agents get // their sessions stripped when they return to the pool. removed_channels.insert(ch); @@ -2087,10 +2101,12 @@ async fn tokio_main() -> Result<()> { } } } - for (channel_id, thread_tags) in - dispatch_pending(&mut pool, &mut queue, &ctx) - { - typing_channels.insert(channel_id, thread_tags); + if pool_ready { + for (channel_id, thread_tags) in + dispatch_pending(&mut pool, &mut queue, &ctx) + { + typing_channels.insert(channel_id, thread_tags); + } } } None => { @@ -2111,7 +2127,9 @@ async fn tokio_main() -> Result<()> { } } => { let _ = result_rx; - if queue.has_flushable_work() { + if !pool_ready { + tracing::debug!("heartbeat_skipped_pool_not_ready"); + } else if queue.has_flushable_work() { tracing::debug!("heartbeat_skipped_events"); for (channel_id, thread_tags) in dispatch_pending(&mut pool, &mut queue, &ctx) @@ -2356,10 +2374,56 @@ async fn tokio_main() -> Result<()> { typing_channels.insert(channel_id, thread_tags); } } + Some(PoolEvent::Wake(attempt, result)) => { + let completion = result.as_ref().map(|_| ()).map_err(|error| error.clone()); + if let Err(error) = + pool_lifecycle.complete_wake(attempt, result, tokio::time::Instant::now()) + { + tracing::warn!(attempt, error, "discarding stale pool wake result"); + continue; + } + match completion { + Ok(()) => { + pool = pool_lifecycle + .take_ready() + .expect("successful wake stores a ready pool"); + pool_ready = true; + emit_runtime_lifecycle( + observer.as_ref(), + &pubkey_hex, + &config.relay_url, + "ready", + None, + ); + for (channel_id, thread_tags) in + dispatch_pending(&mut pool, &mut queue, &ctx) + { + typing_channels.insert(channel_id, thread_tags); + } + } + Err(error) => { + debug_assert_eq!(pool_lifecycle.failed_error(), Some(error.as_str())); + emit_runtime_lifecycle( + observer.as_ref(), + &pubkey_hex, + &config.relay_url, + "failed", + Some(&error), + ); + } + } + } None => {} // relay/heartbeat/shutdown branches handled inline above } } + wake_tasks.shutdown().await; + while let Ok((_attempt, result)) = wake_rx.try_recv() { + if let Ok(mut awakened_pool) = result { + shutdown_agent_pool(&mut awakened_pool).await; + } + } + tracing::info!("shutdown: waiting for in-flight prompts"); // 30 s is generous for in-flight prompts to be cancelled; using // max_turn_duration here would cause Ctrl+C to hang for up to an hour. @@ -3367,6 +3431,152 @@ fn normalized_agent_name(init_result: &serde_json::Value) -> String { .to_ascii_lowercase() } +async fn shutdown_agent_slots(slots: &mut [Option]) { + for slot in slots { + if let Some(mut agent) = slot.take() { + agent.acp.shutdown().await; + } + } +} + +async fn shutdown_agent_pool(pool: &mut AgentPool) { + pool.join_set.shutdown().await; + while let Ok(mut result) = pool.result_rx_try_recv() { + result.agent.acp.shutdown().await; + } + for slot in pool.agents_mut() { + if let Some(mut agent) = slot.take() { + agent.acp.shutdown().await; + } + } +} + +struct PoolStartup { + agents: u32, + command: String, + args: Vec, + extra_env: Vec<(String, String)>, + has_generated_codex_config: bool, + model: Option, + observer: Option, +} + +impl PoolStartup { + fn from_config(config: &Config, observer: Option) -> Self { + Self { + agents: config.agents, + command: config.agent_command.clone(), + args: config.agent_args.clone(), + extra_env: config.persona_env_vars.clone(), + has_generated_codex_config: config.has_generated_codex_config, + model: config.model.clone(), + observer, + } + } +} + +async fn initialize_agent_pool( + startup: &PoolStartup, + mut shutdown: Option>, +) -> Result { + // Finding #10: one agent failing to start must not kill the whole 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 spawn_result = AcpClient::spawn( + &startup.command, + &startup.args, + &startup.extra_env, + startup.has_generated_codex_config, + ) + .await; + match spawn_result { + Ok(mut acp) => { + acp.set_observer(startup.observer.clone(), i); + let initialize = tokio::time::timeout(Duration::from_secs(60), acp.initialize()); + let initialize_result = match shutdown.as_mut() { + Some(shutdown) => tokio::select! { + biased; + _ = shutdown.changed() => { + acp.shutdown().await; + shutdown_agent_slots(&mut agent_slots).await; + return Err(anyhow::anyhow!("pool initialization cancelled by shutdown")); + } + result = initialize => result, + }, + None => initialize.await, + }; + match initialize_result { + Ok(Ok(init_result)) => { + tracing::info!(agent = i, "agent initialized: {init_result}"); + let protocol_version = + init_result["protocolVersion"].as_u64().unwrap_or(1) as u32; + tracing::info!( + agent = i, + name = init_result + .get("agentInfo") + .or_else(|| init_result.get("serverInfo")) + .and_then(|info| info.get("name")) + .and_then(|v| v.as_str()) + .unwrap_or("unknown"), + "agent initialized — non-cancelling steer enabled (try-and-tolerate)" + ); + acp.observe( + "agent_initialized", + serde_json::json!({ + "agentIndex": i, + "initializeResult": init_result, + }), + ); + let agent_name = normalized_agent_name(&init_result); + agent_slots.push(Some(OwnedAgent { + index: i, + acp, + state: SessionState::default(), + model_capabilities: None, + desired_model: startup.model.clone(), + model_overridden: false, + agent_name, + goose_system_prompt_supported: None, + protocol_version, + })); + } + Ok(Err(e)) => { + tracing::error!(agent = i, "agent initialize failed: {e}"); + acp.shutdown().await; + agent_slots.push(None); + } + Err(_) => { + tracing::error!(agent = i, "agent timed out during init (60s)"); + acp.shutdown().await; + agent_slots.push(None); + } + } + } + Err(e) => { + tracing::error!(agent = i, "agent failed to spawn: {e}"); + agent_slots.push(None); + } + } + } + let live_count = agent_slots.iter().filter(|slot| slot.is_some()).count(); + if live_count == 0 { + return Err(anyhow::anyhow!( + "all {} agents failed to start — cannot continue", + startup.agents + )); + } + if live_count < startup.agents as usize { + tracing::warn!( + "started {}/{} agents — continuing with reduced pool", + live_count, + startup.agents + ); + } + tracing::info!("agent_pool_ready agents={}", live_count); + Ok(AgentPool::from_slots(agent_slots)) +} + // ── spawn_and_init ──────────────────────────────────────────────────────────── /// Spawn an agent subprocess and run the MCP `initialize` handshake. /// @@ -4172,6 +4382,7 @@ mod build_mcp_servers_tests { persona_env_vars: vec![], has_generated_codex_config: false, relay_observer: false, + lazy_pool: false, agent_owner: None, no_base_prompt: false, base_prompt_content: None, @@ -4337,6 +4548,7 @@ mod error_outcome_emission_tests { persona_env_vars: vec![], has_generated_codex_config: false, relay_observer: false, + lazy_pool: false, agent_owner: None, no_base_prompt: false, base_prompt_content: None, diff --git a/crates/buzz-acp/src/pool_lifecycle.rs b/crates/buzz-acp/src/pool_lifecycle.rs new file mode 100644 index 000000000..a1099a2dc --- /dev/null +++ b/crates/buzz-acp/src/pool_lifecycle.rs @@ -0,0 +1,312 @@ +//! Lazy agent-pool lifecycle state. +//! +//! Relay connection, subscription, and event buffering live outside this +//! module. This state machine owns only whether a deferred pool has not started, +//! is waking, is ready, or is waiting to retry after a failed wake. + +use std::time::Duration; +use tokio::time::Instant; + +const INITIAL_RETRY_DELAY: Duration = Duration::from_secs(5); +const MAX_RETRY_DELAY: Duration = Duration::from_secs(300); + +#[derive(Debug)] +pub(crate) enum PoolLifecycle

{ + Listening, + Waking { + attempt: u32, + }, + Ready(P), + Failed { + attempt: u32, + retry_at: Instant, + error: String, + }, +} + +impl

PoolLifecycle

{ + pub(crate) fn listening() -> Self { + Self::Listening + } + + /// Start the first wake, or a due retry, when buffered work exists. + /// + /// Returns the attempt token exactly once per transition into `Waking`; + /// callers attach it to the single pool-initialization task and return it + /// with the result. + pub(crate) fn start_wake_if_due( + &mut self, + has_pending_work: bool, + now: Instant, + ) -> Option { + if !has_pending_work { + return None; + } + + let next_attempt = match self { + Self::Listening => Some(1), + Self::Failed { + attempt, retry_at, .. + } if now >= *retry_at => Some(attempt.saturating_add(1)), + Self::Waking { .. } | Self::Ready(_) | Self::Failed { .. } => None, + }; + + if let Some(attempt) = next_attempt { + *self = Self::Waking { attempt }; + } + next_attempt + } + + pub(crate) fn take_ready(&mut self) -> Option

{ + match std::mem::replace(self, Self::Listening) { + Self::Ready(pool) => Some(pool), + other => { + *self = other; + None + } + } + } + + pub(crate) fn waking_attempt(&self) -> Option { + match self { + Self::Waking { attempt } => Some(*attempt), + _ => None, + } + } + + pub(crate) fn retry_at(&self) -> Option { + match self { + Self::Failed { retry_at, .. } => Some(*retry_at), + _ => None, + } + } + + pub(crate) fn failed_error(&self) -> Option<&str> { + match self { + Self::Failed { error, .. } => Some(error), + _ => None, + } + } + + pub(crate) fn cancel_wake(&mut self, attempt: u32, error: String, now: Instant) -> bool { + self.complete_wake(attempt, Err(error), now).is_ok() + } + + /// Complete the matching in-flight wake attempt. + /// + /// A failure remains retryable. A result returned outside `Waking`, or from + /// an older attempt, is rejected: accepting it could replace a newer pool. + pub(crate) fn complete_wake( + &mut self, + completed_attempt: u32, + result: Result, + now: Instant, + ) -> Result<(), &'static str> { + let attempt = match self { + Self::Waking { attempt } if *attempt == completed_attempt => *attempt, + Self::Waking { .. } => return Err("wake result attempt did not match Waking attempt"), + _ => return Err("wake completed while lifecycle was not Waking"), + }; + + *self = match result { + Ok(pool) => Self::Ready(pool), + Err(error) => Self::Failed { + attempt, + retry_at: now + retry_delay(attempt), + error, + }, + }; + Ok(()) + } +} + +fn retry_delay(attempt: u32) -> Duration { + let exponent = attempt.saturating_sub(1).min(63); + let multiplier = 1_u64.checked_shl(exponent).unwrap_or(u64::MAX); + Duration::from_secs( + INITIAL_RETRY_DELAY + .as_secs() + .saturating_mul(multiplier) + .min(MAX_RETRY_DELAY.as_secs()), + ) +} + +#[cfg(test)] +mod tests { + use super::*; + + #[tokio::test(start_paused = true)] + async fn first_pending_event_starts_exactly_one_wake() { + let now = Instant::now(); + let mut lifecycle = PoolLifecycle::<()>::listening(); + + assert_eq!(lifecycle.start_wake_if_due(false, now), None); + assert_eq!(lifecycle.start_wake_if_due(true, now), Some(1)); + assert_eq!(lifecycle.start_wake_if_due(true, now), None); + assert!(matches!(lifecycle, PoolLifecycle::Waking { attempt: 1 })); + } + + #[tokio::test(start_paused = true)] + async fn failure_retries_only_when_work_exists_and_deadline_is_due() { + let now = Instant::now(); + let mut lifecycle = PoolLifecycle::<()>::listening(); + assert_eq!(lifecycle.start_wake_if_due(true, now), Some(1)); + lifecycle + .complete_wake(1, Err("provider unavailable".into()), now) + .unwrap(); + + assert_eq!( + lifecycle.start_wake_if_due(true, now + Duration::from_secs(4)), + None + ); + assert_eq!( + lifecycle.start_wake_if_due(false, now + Duration::from_secs(5)), + None + ); + assert_eq!( + lifecycle.start_wake_if_due(true, now + Duration::from_secs(5)), + Some(2) + ); + assert!(matches!(lifecycle, PoolLifecycle::Waking { attempt: 2 })); + } + + #[tokio::test(start_paused = true)] + async fn retry_backoff_doubles_and_caps_at_five_minutes() { + let mut now = Instant::now(); + let mut lifecycle = PoolLifecycle::<()>::listening(); + + for attempt in 1..=9 { + assert_eq!(lifecycle.start_wake_if_due(true, now), Some(attempt)); + assert!(matches!( + lifecycle, + PoolLifecycle::Waking { attempt: actual } if actual == attempt + )); + lifecycle + .complete_wake(attempt, Err("no brain".into()), now) + .unwrap(); + + let expected = retry_delay(attempt); + let retry_at = match &lifecycle { + PoolLifecycle::Failed { retry_at, .. } => *retry_at, + _ => panic!("failure must enter Failed"), + }; + assert_eq!(retry_at, now + expected); + assert!(expected <= MAX_RETRY_DELAY); + now = retry_at; + } + + assert_eq!(retry_delay(7), MAX_RETRY_DELAY); + assert_eq!(retry_delay(u32::MAX), MAX_RETRY_DELAY); + } + + #[tokio::test(start_paused = true)] + async fn successful_retry_consumes_pool_and_stops_future_wakes() { + let now = Instant::now(); + let mut lifecycle = PoolLifecycle::listening(); + assert_eq!(lifecycle.start_wake_if_due(true, now), Some(1)); + lifecycle + .complete_wake(1, Err("first attempt failed".into()), now) + .unwrap(); + + let retry_at = match &lifecycle { + PoolLifecycle::Failed { retry_at, .. } => *retry_at, + _ => panic!("expected Failed"), + }; + assert_eq!(lifecycle.start_wake_if_due(true, retry_at), Some(2)); + lifecycle.complete_wake(2, Ok("pool"), retry_at).unwrap(); + + assert!(matches!(lifecycle, PoolLifecycle::Ready("pool"))); + assert_eq!( + lifecycle.start_wake_if_due(true, retry_at + Duration::from_secs(600)), + None + ); + } + + #[tokio::test(start_paused = true)] + async fn stale_or_duplicate_wake_result_is_rejected() { + let now = Instant::now(); + let mut lifecycle = PoolLifecycle::<()>::listening(); + assert_eq!( + lifecycle.complete_wake(1, Ok(()), now), + Err("wake completed while lifecycle was not Waking") + ); + + assert_eq!(lifecycle.start_wake_if_due(true, now), Some(1)); + lifecycle.complete_wake(1, Ok(()), now).unwrap(); + assert_eq!( + lifecycle.complete_wake(1, Ok(()), now), + Err("wake completed while lifecycle was not Waking") + ); + assert!(matches!(lifecycle, PoolLifecycle::Ready(()))); + } + + #[tokio::test(start_paused = true)] + async fn stale_attempt_result_cannot_replace_current_wake() { + let now = Instant::now(); + let mut lifecycle = PoolLifecycle::<&str>::listening(); + assert_eq!(lifecycle.start_wake_if_due(true, now), Some(1)); + lifecycle + .complete_wake(1, Err("attempt one failed".into()), now) + .unwrap(); + + let retry_at = match &lifecycle { + PoolLifecycle::Failed { retry_at, .. } => *retry_at, + _ => panic!("expected Failed"), + }; + assert_eq!(lifecycle.start_wake_if_due(true, retry_at), Some(2)); + assert_eq!( + lifecycle.complete_wake(1, Ok("stale pool"), retry_at), + Err("wake result attempt did not match Waking attempt") + ); + assert!(matches!(lifecycle, PoolLifecycle::Waking { attempt: 2 })); + lifecycle + .complete_wake(2, Ok("current pool"), retry_at) + .unwrap(); + assert!(matches!(lifecycle, PoolLifecycle::Ready("current pool"))); + } + + #[tokio::test(start_paused = true)] + async fn cancelled_wake_enters_failed_and_can_retry() { + let now = Instant::now(); + let mut lifecycle = PoolLifecycle::<()>::listening(); + assert_eq!(lifecycle.start_wake_if_due(true, now), Some(1)); + assert_eq!(lifecycle.waking_attempt(), Some(1)); + assert!(lifecycle.cancel_wake(1, "task panicked".into(), now)); + assert_eq!(lifecycle.failed_error(), Some("task panicked")); + assert_eq!( + lifecycle.start_wake_if_due(true, now + Duration::from_secs(5)), + Some(2) + ); + } + + #[test] + fn take_ready_transfers_pool_exactly_once() { + let now = Instant::now(); + let mut lifecycle = PoolLifecycle::listening(); + assert_eq!(lifecycle.start_wake_if_due(true, now), Some(1)); + lifecycle.complete_wake(1, Ok("pool"), now).unwrap(); + assert_eq!(lifecycle.take_ready(), Some("pool")); + assert_eq!(lifecycle.take_ready(), None); + } + + #[test] + fn failed_state_preserves_attempt_deadline_and_error() { + let now = Instant::now(); + let mut lifecycle = PoolLifecycle::<()>::listening(); + assert_eq!(lifecycle.start_wake_if_due(true, now), Some(1)); + lifecycle.complete_wake(1, Err("boom".into()), now).unwrap(); + + match lifecycle { + PoolLifecycle::Failed { + attempt, + retry_at, + error, + } => { + assert_eq!(attempt, 1); + assert_eq!(retry_at, now + Duration::from_secs(5)); + assert_eq!(error, "boom"); + } + _ => panic!("expected Failed"), + } + } +} diff --git a/crates/buzz-acp/tests/pool_lifecycle_state.rs b/crates/buzz-acp/tests/pool_lifecycle_state.rs new file mode 100644 index 000000000..8bbc5b8e2 --- /dev/null +++ b/crates/buzz-acp/tests/pool_lifecycle_state.rs @@ -0,0 +1,4 @@ +// Compile and run the lifecycle state-machine contract as an integration target. +#[allow(dead_code)] +#[path = "../src/pool_lifecycle.rs"] +mod pool_lifecycle;