mirror of
https://github.com/block/buzz.git
synced 2026-08-18 06:50:31 +02:00
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 <tlongwell@block.xyz> Signed-off-by: Tyler Longwell <tlongwell@block.xyz>
This commit is contained in:
co-authored by
Tyler Longwell
parent
8908bd6b71
commit
1c0a4c838b
@@ -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:
|
||||
|
||||
|
||||
@@ -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`. |
|
||||
|
||||
@@ -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
@@ -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,
|
||||
|
||||
@@ -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;
|
||||
Reference in New Issue
Block a user