diff --git a/crates/buzz-agent/src/lib.rs b/crates/buzz-agent/src/lib.rs index f94d60a2a..a879bd621 100644 --- a/crates/buzz-agent/src/lib.rs +++ b/crates/buzz-agent/src/lib.rs @@ -206,21 +206,40 @@ async fn async_main() { models_cache: tokio::sync::OnceCell::new(), }); let (wire_tx, wire_rx) = mpsc::channel::(64); - let writer = tokio::spawn(wire::writer_task(wire_rx)); - if let Err(e) = read_loop( - BufReader::new(tokio::io::stdin()), - app.clone(), - wire_tx, - max_line, - ) - .await - { - tracing::error!("io: reader: {e}"); + let mut writer = tokio::spawn(wire::writer_task(wire_rx)); + // Whichever ends first drives shutdown. The reader ending is the normal + // path (stdin EOF/error). The writer ending while the reader still runs + // means stdout is closed/broken: no reply can ever be written, so we must + // stop reading and cancel every session rather than leave the process + // reading input while outstanding permission asks wait out their full + // deadline for a response that can never arrive. + tokio::select! { + r = read_loop( + BufReader::new(tokio::io::stdin()), + app.clone(), + wire_tx, + max_line, + ) => { + if let Err(e) = r { + tracing::error!("io: reader: {e}"); + } + cancel_all_sessions(&app).await; + let _ = writer.await; + } + _ = &mut writer => { + tracing::error!("io: writer exited (stdout closed); shutting down connection"); + cancel_all_sessions(&app).await; + } } +} + +/// Signal every live session to cancel. Run on connection teardown so in-flight +/// prompts — including any waiting on a `session/request_permission` response — +/// resolve promptly instead of waiting out their deadline. +async fn cancel_all_sessions(app: &Arc) { for session in app.sessions.lock().await.values() { let _ = session.cancel_tx.send(true); } - let _ = writer.await; } async fn read_loop( diff --git a/crates/buzz-agent/src/permission.rs b/crates/buzz-agent/src/permission.rs index 0baef5137..09d7421cf 100644 --- a/crates/buzz-agent/src/permission.rs +++ b/crates/buzz-agent/src/permission.rs @@ -24,6 +24,11 @@ //! a later lease `Drop` is a harmless no-op. //! - **Unknown/late ids ignored.** A response whose id is not a live entry is //! logged and dropped. +//! - **Undeliverable asks are terminal.** If the output wire is closed when the +//! request is enqueued, [`PermissionBroker::request_permission`] fails closed +//! immediately (dropping the lease removes the entry and releases the permit) +//! rather than leaving a resident waiter to time out — a closed wire can never +//! carry the reply. //! - **Single absolute deadline.** Admission wait and response wait share one //! absolute deadline computed at gate entry, so a saturated call cannot live //! for two full timeout windows. @@ -51,6 +56,12 @@ pub const PERMISSION_DENIED_MSG: &str = "permission denied: the tool call was no pub const PERMISSION_TIMEOUT_MSG: &str = "permission request timed out: the tool call was not authorized"; +/// Model-visible tool error when the permission request cannot be delivered +/// because the output wire is closed. Terminal and immediate — no waiter is +/// left resident, since a closed wire can never carry a reply. +pub const PERMISSION_WIRE_CLOSED_MSG: &str = + "permission request undeliverable: the tool call was not authorized"; + /// Outcome of asking the client to authorize one tool call. #[derive(Debug, PartialEq, Eq)] pub enum PermissionDecision { @@ -77,8 +88,20 @@ pub struct PermissionBroker { next_id: AtomicU64, /// Absolute deadline budget shared by admission + response wait. timeout: Duration, + /// Test-only: invoked by the waiter the instant it observes a delivered + /// response, with `true` iff the entry was already claimed (removed) from + /// `pending` before the wake. Makes claim-before-wake ordering + /// mutation-sensitive — a wake-before-claim mutant reports `false`, which a + /// purely behavioral test cannot detect (the waiter reads the oneshot once + /// either way). + #[cfg(test)] + wake_observer: Mutex>, } +/// Test-only wake-boundary observer; see [`PermissionBroker::wake_observer`]. +#[cfg(test)] +type WakeObserver = Arc; + impl PermissionBroker { /// `max_pending` is validated `>= 1` by config; `timeout` is injectable so /// broker unit tests exercise the timeout/abort paths without a 330s wait. @@ -88,6 +111,31 @@ impl PermissionBroker { pending: Arc::new(Mutex::new(HashMap::new())), next_id: AtomicU64::new(0), timeout, + #[cfg(test)] + wake_observer: Mutex::new(None), + } + } + + /// Test-only: register a callback the waiter fires the instant it observes a + /// delivered response, with `true` iff the correlation entry was already + /// claimed (removed) before the wake. Used to prove claim-before-wake + /// ordering in a way a wake-before-claim mutant cannot satisfy. + #[cfg(test)] + pub fn set_wake_observer(&self, observer: WakeObserver) { + *self.wake_observer.lock().unwrap() = Some(observer); + } + + /// Test-only: fire the wake observer (if any) with the claimed-before-wake + /// status of `id`. Called synchronously by the waiter the moment it receives + /// its response, so the observed `pending` state is exactly the state at the + /// wake — deterministic in production (removal happens-before the send) and + /// violated by a wake-before-claim mutant. + #[cfg(test)] + fn observe_wake(&self, id: u64) { + let claimed = !self.pending.lock().unwrap().contains_key(&id); + let observer = self.wake_observer.lock().unwrap().clone(); + if let Some(observer) = observer { + observer(claimed); } } @@ -177,21 +225,39 @@ impl PermissionBroker { &call.name, &call.arguments, ); - wire::send( + if wire::send_checked( wire, wire::request_permission(lease.id_value.clone(), params), ) - .await; + .await + .is_err() + { + // The output wire is closed: this ask will never be written and no + // reply can ever arrive. Fail closed now — dropping `lease` here + // removes the entry and releases the permit synchronously — instead + // of leaving the entry resident until the deadline expires. + return PermissionDecision::Denied(PERMISSION_WIRE_CLOSED_MSG); + } // ── Response wait ────────────────────────────────────────────────── if *cancel.borrow() { return PermissionDecision::Cancelled; } + #[cfg(test)] + let id = lease.id; tokio::select! { biased; _ = cancel.changed() => PermissionDecision::Cancelled, r = &mut lease.rx => match r { - Ok(result) => evaluate(&result), + Ok(result) => { + // The waiter observes delivery here. At this instant the + // entry must already be claimed (removed) — delivery removes + // before it sends. The observer is test-only and a no-op in + // production. + #[cfg(test)] + self.observe_wake(id); + evaluate(&result) + } // Sender dropped without sending — should not happen (delivery // always sends before drop); fail closed. Err(_) => PermissionDecision::Denied(PERMISSION_DENIED_MSG), @@ -252,10 +318,14 @@ fn evaluate(result: &Value) -> PermissionDecision { } /// Recover the correlation key from an outbound request id echoed by the -/// client. Only ids we minted (`perm-`) are ours; anything else is a -/// foreign/stale id and is ignored. +/// client. Only ids we minted (`perm-`, canonical decimal) are ours; a +/// noncanonical alias (`perm-01`, `perm-+0`, `perm-00`) or any other string is +/// a foreign/stale id and is ignored. Requiring an exact round-trip means only +/// the string the broker actually minted correlates — no alias is ever live. fn parse_id(id: &Value) -> Option { - id.as_str()?.strip_prefix("perm-")?.parse().ok() + let s = id.as_str()?; + let n: u64 = s.strip_prefix("perm-")?.parse().ok()?; + (format!("perm-{n}") == s).then_some(n) } #[cfg(test)] @@ -337,6 +407,31 @@ mod tests { assert_eq!(parse_id(&Value::Null), None); } + /// Noncanonical strings that `u64::parse` would otherwise accept as aliases + /// of a minted id must NOT correlate. Only the exact string the broker + /// minted (`format!("perm-{n}")`) is live; leading zeros, a sign, or + /// whitespace make the id foreign and it is ignored. Without the exact + /// round-trip check these would resolve live asks under ids the broker + /// never issued. + #[test] + fn test_parse_id_rejects_noncanonical_aliases() { + for alias in [ + "perm-00", // extra leading zero + "perm-01", // leading zero + "perm-+0", // explicit sign + "perm-0x1", // hex + "perm- 1", // leading space + "perm-1 ", // trailing space + "perm-1_000", // digit separator + ] { + assert_eq!( + parse_id(&json!(alias)), + None, + "alias must be foreign: {alias}" + ); + } + } + // ── Delivery: exact allow / deny ────────────────────────────────────────── #[tokio::test(flavor = "multi_thread", worker_threads = 2)] @@ -439,7 +534,88 @@ mod tests { assert_eq!(broker.available_permits(), 4, "timeout releases the slot"); } - // ── Cancellation while waiting ──────────────────────────────────────────── + // ── Undeliverable ask (closed wire) is terminal ─────────────────────────── + + /// When the output wire is closed, the ask can never be written and no + /// reply can ever arrive. `request_permission` must fail closed + /// *immediately* — denying with the wire-closed reason and leaving zero + /// pending entries and zero held permits — rather than registering an entry + /// that waits out the full deadline. Uses a LONG timeout so a wrong + /// implementation that waits the deadline would visibly hang the test far + /// past its own assertions. + #[tokio::test(flavor = "multi_thread", worker_threads = 2)] + async fn test_closed_wire_denies_immediately_without_leaking_state() { + let broker = Arc::new(PermissionBroker::new(4, LONG)); + // Drop the receiver so every send fails: the writer is gone. + let (tx, rx) = mpsc::channel(8); + drop(rx); + let (_cancel_tx, mut cancel_rx) = watch::channel(false); + let call = tool_call(); + + // Bound the whole call: correct behavior returns at once; a regression + // that waits the deadline blows this timeout instead of hanging LONG. + let decision = tokio::time::timeout( + Duration::from_secs(2), + broker.request_permission(&tx, 2, "ses_a", &call, &mut cancel_rx), + ) + .await + .expect("closed wire must deny immediately, not wait the deadline"); + + assert_eq!( + decision, + PermissionDecision::Denied(PERMISSION_WIRE_CLOSED_MSG) + ); + assert_eq!( + broker.pending_count(), + 0, + "undeliverable ask leaves no resident entry" + ); + assert_eq!( + broker.available_permits(), + 4, + "undeliverable ask releases its admission slot" + ); + } + + // ── Claim-before-wake ordering (mutation-sensitive) ─────────────────────── + + /// The waiter must observe the correlation entry already *claimed* (removed + /// from `pending`) at the instant it wakes with the delivered response — + /// `deliver` removes before it sends. The wake observer fires synchronously + /// inside the waiter's response arm, so it captures the exact `pending` + /// state at the wake. A wake-before-claim mutant (send first, remove after) + /// makes the observed state `false` and fails this assertion; the behavioral + /// delivery tests cannot detect that mutant because the waiter reads the + /// oneshot exactly once either way. + #[tokio::test(flavor = "multi_thread", worker_threads = 2)] + async fn test_waiter_observes_entry_claimed_before_wake() { + let broker = Arc::new(PermissionBroker::new(4, LONG)); + let claimed_at_wake = Arc::new(Mutex::new(None::)); + let sink = Arc::clone(&claimed_at_wake); + broker.set_wake_observer(Arc::new(move |claimed| { + *sink.lock().unwrap() = Some(claimed); + })); + + let (tx, mut rx) = mpsc::channel(8); + let (_cancel_tx, mut cancel_rx) = watch::channel(false); + let b = Arc::clone(&broker); + let call = tool_call(); + let task = tokio::spawn(async move { + b.request_permission(&tx, 2, "ses_a", &call, &mut cancel_rx) + .await + }); + + let id = next_request_id(&mut rx).await; + broker.deliver(&id, selected(ALLOW_OPTION_ID)); + assert_eq!(task.await.unwrap(), PermissionDecision::Allowed); + assert_eq!( + *claimed_at_wake.lock().unwrap(), + Some(true), + "entry must be claimed (removed) before the waiter is woken" + ); + } + + // ── Cancellation while waiting ─────────────────────────────────────────── #[tokio::test(flavor = "multi_thread", worker_threads = 2)] async fn test_cancel_while_waiting_returns_cancelled_and_removes_state() { diff --git a/crates/buzz-agent/src/wire.rs b/crates/buzz-agent/src/wire.rs index 2015e2237..2d1f0d6c5 100644 --- a/crates/buzz-agent/src/wire.rs +++ b/crates/buzz-agent/src/wire.rs @@ -353,7 +353,18 @@ pub fn session_update_with_goose_meta(sid: &str, update: Value, goose_meta: Valu } pub async fn send(wire: &WireSender, msg: Value) { - let _ = wire.send(WireMsg::Notify(msg)).await; + let _ = send_checked(wire, msg).await; +} + +/// Enqueue a frame, reporting whether the writer accepted it. Unlike mpsc's +/// non-blocking `try_send`, this awaits channel capacity; it fails only when +/// the writer task has dropped its receiver, which happens exactly when the +/// writer has exited because stdout is closed/broken. A frame that fails here +/// will never be written, so callers that correlate a response — the +/// permission broker — must fail closed immediately rather than wait out a +/// deadline for a reply that can never arrive. +pub async fn send_checked(wire: &WireSender, msg: Value) -> Result<(), ()> { + wire.send(WireMsg::Notify(msg)).await.map_err(|_| ()) } pub async fn read_bounded_line( @@ -724,4 +735,25 @@ mod tests { other => panic!("expected Response, got {other:?}"), } } + + // ── send_checked: observable wire closure ──────────────────────────────── + + /// `send_checked` reports `Ok` while the writer's receiver is alive and + /// `Err` once it is gone (writer task exited on closed/broken stdout). This + /// is the contract the permission broker relies on to fail an undeliverable + /// ask closed immediately instead of waiting out its deadline for a reply + /// that can never be written. + #[tokio::test] + async fn send_checked_reports_closure_when_writer_gone() { + let (tx, rx) = mpsc::channel::(4); + assert!( + send_checked(&tx, json!({ "ok": 1 })).await.is_ok(), + "send succeeds while the writer receiver is alive" + ); + drop(rx); // writer exited → receiver dropped + assert!( + send_checked(&tx, json!({ "ok": 2 })).await.is_err(), + "send reports failure once the writer is gone" + ); + } }