Merge Lane A: ACP lazy pool lifecycle (eccf18cd7)

Opt-in --lazy-pool / BUZZ_ACP_LAZY_POOL deferred agent-pool startup with
Listening -> Waking -> Ready/Failed lifecycle, single-wake, bounded retry,
and managed_agent_runtime_lifecycle observer frames.

Co-authored-by: Tyler Longwell <tlongwell@block.xyz>
Signed-off-by: Tyler Longwell <tlongwell@block.xyz>
This commit is contained in:
npub1qyvc0c5kl4gqv2fd97fsk46tu378sqgy35vc83rvgfwne90sel7s0ed67d
2026-07-19 12:10:30 -04:00
co-authored by Tyler Longwell
6 changed files with 658 additions and 101 deletions
+8
View File
@@ -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:
+1
View File
@@ -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`. |
+20
View File
@@ -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<String>,
@@ -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);
+313 -101
View File
@@ -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<Option<OwnedAgent>> = 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<AgentPool> = 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::<RespawnResult>(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<AgentPool, String>)>(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<PromptResult>),
Panic(tokio::task::JoinError),
SteerAck(SteerAckEvent),
Wake(u32, Result<AgentPool, String>),
}
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<OwnedAgent>]) {
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<String>,
extra_env: Vec<(String, String)>,
has_generated_codex_config: bool,
model: Option<String>,
observer: Option<observer::ObserverHandle>,
}
impl PoolStartup {
fn from_config(config: &Config, observer: Option<observer::ObserverHandle>) -> 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<watch::Receiver<()>>,
) -> Result<AgentPool> {
// 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<Option<OwnedAgent>> = 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,
+312
View File
@@ -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<P> {
Listening,
Waking {
attempt: u32,
},
Ready(P),
Failed {
attempt: u32,
retry_at: Instant,
error: String,
},
}
impl<P> PoolLifecycle<P> {
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<u32> {
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<P> {
match std::mem::replace(self, Self::Listening) {
Self::Ready(pool) => Some(pool),
other => {
*self = other;
None
}
}
}
pub(crate) fn waking_attempt(&self) -> Option<u32> {
match self {
Self::Waking { attempt } => Some(*attempt),
_ => None,
}
}
pub(crate) fn retry_at(&self) -> Option<Instant> {
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<P, String>,
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"),
}
}
}
@@ -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;