mirror of
https://github.com/block/buzz.git
synced 2026-08-18 06:50:31 +02:00
fix(buzz-agent): make undeliverable permission asks terminal, prove claim-before-wake ordering, reject noncanonical ids
Address three review findings on the permission broker's ask surface. An undeliverable ask is now terminal. wire::send_checked surfaces the mpsc send failure that occurs exactly when the writer task has exited on a closed/broken stdout; request_permission fails closed immediately on that error (dropping the lease removes the entry and releases the permit synchronously) instead of leaving a resident waiter to time out. async_main now selects on both the reader and the writer JoinHandle: writer death cancels every session and closes the reader lifecycle rather than reading input while asks wait out their deadline for a reply that can never be written. Claim-before-wake ordering is now mutation-sensitive. Production enforces it structurally — the oneshot sender is consumed by remove, so a send-before-remove mutant cannot compile without swapping the channel type. A test-only wake observer fires synchronously in the waiter's response arm and asserts the entry is already absent from pending at the wake; a faithful wake-before-claim mutant makes it observe false and the test goes red while the behavioral delivery tests stay green. parse_id now requires an exact canonical round-trip, so noncanonical aliases (perm-01, perm-+0, perm-00) that u64::parse would accept are rejected as foreign ids rather than resolving live asks. Co-authored-by: Will Pfleger <pfleger.will@gmail.com> Signed-off-by: Will Pfleger <pfleger.will@gmail.com>
This commit is contained in:
@@ -206,21 +206,40 @@ async fn async_main() {
|
||||
models_cache: tokio::sync::OnceCell::new(),
|
||||
});
|
||||
let (wire_tx, wire_rx) = mpsc::channel::<WireMsg>(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<App>) {
|
||||
for session in app.sessions.lock().await.values() {
|
||||
let _ = session.cancel_tx.send(true);
|
||||
}
|
||||
let _ = writer.await;
|
||||
}
|
||||
|
||||
async fn read_loop<R: tokio::io::AsyncBufRead + Unpin>(
|
||||
|
||||
@@ -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<Option<WakeObserver>>,
|
||||
}
|
||||
|
||||
/// Test-only wake-boundary observer; see [`PermissionBroker::wake_observer`].
|
||||
#[cfg(test)]
|
||||
type WakeObserver = Arc<dyn Fn(bool) + Send + Sync>;
|
||||
|
||||
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-<n>`) are ours; anything else is a
|
||||
/// foreign/stale id and is ignored.
|
||||
/// client. Only ids we minted (`perm-<n>`, 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<u64> {
|
||||
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::<bool>));
|
||||
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() {
|
||||
|
||||
@@ -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<R: AsyncBufRead + Unpin>(
|
||||
@@ -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::<WireMsg>(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"
|
||||
);
|
||||
}
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user