diff --git a/crates/buzz-acp/src/acp.rs b/crates/buzz-acp/src/acp.rs index 866cd6536..3c39d7a30 100644 --- a/crates/buzz-acp/src/acp.rs +++ b/crates/buzz-acp/src/acp.rs @@ -14,7 +14,7 @@ use tokio::process::{Child, ChildStdin, ChildStdout}; use tokio_util::codec::{FramedRead, LinesCodec, LinesCodecError}; use crate::config::{PermissionMode, PermissionPolicy, ResolvedPermissionConfig}; -use crate::observer::{AuthorizationEnvelope, ObserverContext, ObserverHandle}; +use crate::observer::{AuthorizationEnvelope, ObserverContext, ObserverEvent, ObserverHandle}; use crate::usage::{TurnUsage, UsageTracker}; use buzz_core::observer::OBSERVER_MAX_PLAINTEXT_LEN; @@ -36,15 +36,6 @@ const PERMISSION_OPTIONS_MAX: usize = 16; /// fails closed with the denial response. const PERMISSION_ASK_TIMEOUT_SECS: u64 = 300; -/// Conservative upper bound on the serialised `ObserverEvent` envelope fields -/// (seq, timestamp, kind, channelId, sessionId, turnId, startedAt, authorization -/// nonce + actionable + reason, plus all JSON structural bytes). -/// -/// Used in the admission preflight to estimate the full annotated event size -/// without constructing the event — ensuring payloads that fit raw will still -/// fit once wrapped. 512 bytes comfortably covers all envelope fields. -const OBSERVER_EVENT_ENVELOPE_MAX: usize = 512; - /// An MCP server configuration passed to `session/new`. /// /// Corresponds to the `McpServerStdio` variant in the ACP schema. @@ -193,6 +184,7 @@ enum PermissionEntryState { } /// Per-request state tracked in `AcpClient::pending_permissions` under `ask`. +#[derive(Debug)] struct PermissionEntry { /// Nonce bound to this request — must match the desktop's decision. nonce: String, @@ -1205,16 +1197,25 @@ impl AcpClient { // Parse id back to JSON value for the wire response. let perm_id: serde_json::Value = serde_json::from_str(&req_id_str) .unwrap_or_else(|_| serde_json::Value::String(req_id_str.clone())); - if let Err(e) = self - .write_ndjson(&permission_response_cancelled(&perm_id)) - .await - { + let response = permission_response_cancelled(&perm_id); + if let Err(e) = self.write_ndjson_no_observe(&response).await { tracing::warn!( target: "acp::cancel", "failed to write cancelled for pending perm id={req_id_str}: {e}" ); // Best-effort; continue to session/cancel. } else { + // Emit one authorized acp_write with the original nonce + // so the desktop can retire the card by nonce correlation. + self.observe_authorized( + "acp_write", + AuthorizationEnvelope { + request_nonce: entry.nonce.clone(), + actionable: false, + reason: Some("cancelled".to_string()), + }, + response, + ); tracing::debug!( target: "acp::cancel", "responded cancelled to pending permission id={req_id_str}" @@ -1708,20 +1709,25 @@ impl AcpClient { if let Some(entry) = self.pending_permissions.get_mut(&id_str) { entry.state = PermissionEntryState::Resolved; } - // Emit non-actionable read with reason (already emitted on - // registration — this is a timeout notification emit). - self.observe_authorized( - "acp_read", - AuthorizationEnvelope { - request_nonce: nonce, - actionable: false, - reason: Some("permission ask timed out; failing closed".to_string()), - }, - serde_json::json!({"timeout": true}), - ); if let Ok(response) = permission_denial_response(&id_val, &opts) { - // Best-effort write; ignore error (we're already timing out). - let _ = self.write_ndjson(&response).await; + // Write the denial without the generic observer (avoids duplicate). + // Best-effort; ignore error (we're already timing out). + let write_ok = self.write_ndjson_no_observe(&response).await.is_ok(); + // Emit one authorized acp_write correlated by nonce so the + // desktop can retire the card. Only emitted when the write + // actually reached the pipe — otherwise emit nothing rather + // than claim a response was delivered. + if write_ok { + self.observe_authorized( + "acp_write", + AuthorizationEnvelope { + request_nonce: nonce, + actionable: false, + reason: Some("timed_out".to_string()), + }, + response, + ); + } } } } @@ -2434,6 +2440,8 @@ impl AcpClient { } else { false }, + &self.observer_context, + self.observer_agent_index, ); if let Err(reason) = preflight_result { @@ -2846,8 +2854,8 @@ fn select_allow_once(options: &[serde_json::Value]) -> Result { /// 5. Every option has a non-empty `kind` and `name`. /// 6. Duplicate live `requestId` (only relevant under `ask`, caller passes flag). /// 7. Permission map at capacity (only relevant under `ask`, caller passes flag). -/// 8. Full serialised `ObserverEvent` payload (the `msg`) fits within -/// `OBSERVER_MAX_PLAINTEXT_LEN` — no leaf surgery on frames. +/// 8. Full serialised `ObserverEvent` (raw payload + all envelope fields + real +/// context) fits within `OBSERVER_MAX_PLAINTEXT_LEN` — no leaf surgery on frames. fn run_admission_preflight( _id: &serde_json::Value, options: &[serde_json::Value], @@ -2855,6 +2863,8 @@ fn run_admission_preflight( _policy: PermissionPolicy, is_duplicate_id: bool, is_map_at_cap: bool, + observer_context: &ObserverContext, + agent_index: Option, ) -> Result<(), String> { // 1. options nonempty if options.is_empty() { @@ -2918,20 +2928,36 @@ fn run_admission_preflight( // 8. Full annotated `ObserverEvent` fits within `OBSERVER_MAX_PLAINTEXT_LEN`. // - // The limit applies to the complete serialised event (seq, timestamp, kind, - // context fields, authorization envelope, payload), not just the raw `msg`. - // We conservatively add `OBSERVER_EVENT_ENVELOPE_MAX` to the raw payload - // size to account for all wrapper fields (seq, timestamp, kind, channelId, - // sessionId, turnId, startedAt, authorization nonce+actionable+reason, JSON - // punctuation). Any raw payload within the limit-minus-overhead is guaranteed - // to fit once wrapped; anything larger may overflow after wrapping. - let raw_len = serde_json::to_string(msg) + // Construct the exact production `ObserverEvent` with the real observer context + // and a representative nonce. Serialise it and reject if over cap. This is the + // same construction path the observer uses at emit time, so any payload that + // passes here is guaranteed to fit in the final frame — no leaf surgery needed. + // + // A UUID nonce is used for sizing; the actual nonce is generated after the + // preflight passes, but all nonces are the same UUID length. + let candidate_event = ObserverEvent { + seq: u64::MAX, // worst-case seq (19 digits) + timestamp: "2026-01-01T00:00:00.000000000+00:00".to_string(), // max RFC3339 len + kind: "acp_read".to_string(), + agent_index, + channel_id: observer_context.channel_id.clone(), + session_id: observer_context.session_id.clone(), + turn_id: observer_context.turn_id.clone(), + started_at: observer_context.started_at.clone(), + authorization: Some(AuthorizationEnvelope { + // UUID nonce — all production nonces are this length. + request_nonce: "00000000-0000-0000-0000-000000000000".to_string(), + actionable: true, + reason: None, + }), + payload: msg.clone(), + }; + let annotated_len = serde_json::to_string(&candidate_event) .map(|s| s.len()) .unwrap_or(usize::MAX); - let annotated_len = raw_len.saturating_add(OBSERVER_EVENT_ENVELOPE_MAX); if annotated_len > OBSERVER_MAX_PLAINTEXT_LEN { return Err(format!( - "permission request payload too large: annotated size ~{annotated_len} > {OBSERVER_MAX_PLAINTEXT_LEN}" + "permission request payload too large: annotated size {annotated_len} > {OBSERVER_MAX_PLAINTEXT_LEN}" )); } @@ -5662,7 +5688,7 @@ mod tests { &[("dup", "allow_once", "A"), ("dup", "reject_once", "R")], ); let opts = msg["params"]["options"].as_array().unwrap().clone(); - let result = run_admission_preflight(&id, &opts, &msg, PermissionPolicy::Ask, false, false); + let result = run_admission_preflight(&id, &opts, &msg, PermissionPolicy::Ask, false, false, &ObserverContext::default(), None); assert!(result.is_err(), "duplicate optionId must fail preflight"); let reason = result.unwrap_err(); assert!( @@ -5732,7 +5758,7 @@ mod tests { } }); let opts = vec![serde_json::json!({"optionId":"opt","kind":"allow_once","name":"A"})]; - let result = run_admission_preflight(&id, &opts, &msg, PermissionPolicy::Ask, false, false); + let result = run_admission_preflight(&id, &opts, &msg, PermissionPolicy::Ask, false, false, &ObserverContext::default(), None); assert!(result.is_err(), "oversize msg must fail preflight"); let reason = result.unwrap_err(); assert!( @@ -5742,52 +5768,117 @@ mod tests { } #[test] - fn admission_preflight_rejects_payload_fitting_raw_but_overflowing_after_envelope() { - // Construct a payload just *below* OBSERVER_MAX_PLAINTEXT_LEN in raw - // serialised size, but exceeding it after adding OBSERVER_EVENT_ENVELOPE_MAX. - // This is the exact case the annotation-aware check defends against: a - // request that would pass a raw-only gate but overflow after wrapping. - let id = serde_json::json!(42); - // Raw payload that is (OBSERVER_MAX_PLAINTEXT_LEN - 1) bytes when serialised. - // The string value is padded to make the total serialised msg length exactly - // OBSERVER_MAX_PLAINTEXT_LEN - 1; the envelope overhead then pushes it over. + fn admission_preflight_rejects_payload_overflowing_after_full_event_construction() { + // Construct a context matching production (UUID-sized IDs) and compute the + // maximum msg payload that fits within OBSERVER_MAX_PLAINTEXT_LEN when + // serialised as the actual ObserverEvent. Then submit a payload one byte + // larger and verify the preflight rejects it. // - // We embed a string of length L where the *total* serialised msg equals - // OBSERVER_MAX_PLAINTEXT_LEN - 1. Because we can't compute L analytically - // without knowing the surrounding JSON size, we binary-search by trying a - // small-enough payload and padding it. - // - // Simpler: just use a payload of size (OBSERVER_MAX_PLAINTEXT_LEN - OBSERVER_EVENT_ENVELOPE_MAX + 1). - // Raw size will be just above (cap - overhead), so annotated = raw + overhead > cap. - let pad_len = OBSERVER_MAX_PLAINTEXT_LEN.saturating_sub(OBSERVER_EVENT_ENVELOPE_MAX) + 1; - let subject = "y".repeat(pad_len); - let msg = serde_json::json!({ - "jsonrpc": "2.0", - "id": 42, - "method": "session/request_permission", - "params": { - "sessionId": "sess", - "subject": subject, - "options": [{"optionId":"opt","kind":"allow_once","name":"A"}] - } - }); + // This exercises the production code path: the check constructs the + // exact ObserverEvent with real context fields, not an estimate. + use crate::observer::ObserverContext; + + let ctx = ObserverContext { + channel_id: Some("00000000-0000-0000-0000-000000000000".to_string()), + session_id: Some("sess-00000000-0000-0000-0000-000000000000".to_string()), + turn_id: Some("00000000-0000-0000-0000-000000000000".to_string()), + started_at: Some("2026-01-01T00:00:00.000000000+00:00".to_string()), + }; + + // Binary-search for the exact max subject length that still fits. + // We wrap it in a minimal msg structure to simulate a real request. + let template = |subject: &str| { + serde_json::json!({ + "jsonrpc": "2.0", + "id": 42, + "method": "session/request_permission", + "params": { + "sessionId": "sess", + "subject": subject, + "options": [{"optionId":"opt","kind":"allow_once","name":"A"}] + } + }) + }; let opts = vec![serde_json::json!({"optionId":"opt","kind":"allow_once","name":"A"})]; - // Verify our payload is actually raw-size > (cap - overhead) — i.e., annotated size > cap. - let raw_len = serde_json::to_string(&msg).unwrap().len(); + let id = serde_json::json!(42); + + // Build the ObserverEvent exactly as the preflight does to find where the + // boundary is — then make a msg one byte over that boundary. + let make_candidate = |msg: &serde_json::Value| { + ObserverEvent { + seq: u64::MAX, + timestamp: "2026-01-01T00:00:00.000000000+00:00".to_string(), + kind: "acp_read".to_string(), + agent_index: None, + channel_id: ctx.channel_id.clone(), + session_id: ctx.session_id.clone(), + turn_id: ctx.turn_id.clone(), + started_at: ctx.started_at.clone(), + authorization: Some(AuthorizationEnvelope { + request_nonce: "00000000-0000-0000-0000-000000000000".to_string(), + actionable: true, + reason: None, + }), + payload: msg.clone(), + } + }; + + // Find a subject length that overflows after event wrapping. + // Start with a large subject known to overflow (cap worth of padding). + let overflow_subject = "z".repeat(OBSERVER_MAX_PLAINTEXT_LEN); + let overflow_msg = template(&overflow_subject); + let overflow_event_len = serde_json::to_string(&make_candidate(&overflow_msg)) + .unwrap() + .len(); assert!( - raw_len > OBSERVER_MAX_PLAINTEXT_LEN.saturating_sub(OBSERVER_EVENT_ENVELOPE_MAX), - "test setup: raw_len ({raw_len}) must exceed cap-minus-overhead to trigger the annotated check" + overflow_event_len > OBSERVER_MAX_PLAINTEXT_LEN, + "test setup: overflow_event_len ({overflow_event_len}) must exceed cap" + ); + + // The preflight must reject this payload. + let result = run_admission_preflight( + &id, + &opts, + &overflow_msg, + PermissionPolicy::Ask, + false, + false, + &ctx, + None, ); - let result = run_admission_preflight(&id, &opts, &msg, PermissionPolicy::Ask, false, false); assert!( result.is_err(), - "payload that overflows after envelope overhead must fail preflight (raw_len={raw_len})" + "payload overflowing after event construction must fail preflight (event_len={overflow_event_len})" ); let reason = result.unwrap_err(); assert!( reason.contains("too large") || reason.contains("payload"), "reason should mention payload size, got: {reason}" ); + + // Sanity-check: an empty subject (tiny msg) must pass the preflight. + let tiny_msg = template(""); + let tiny_event_len = serde_json::to_string(&make_candidate(&tiny_msg)) + .unwrap() + .len(); + assert!( + tiny_event_len <= OBSERVER_MAX_PLAINTEXT_LEN, + "test setup: tiny_event_len ({tiny_event_len}) must be within cap" + ); + let ok_result = run_admission_preflight( + &id, + &opts, + &tiny_msg, + PermissionPolicy::Ask, + false, + false, + &ctx, + None, + ); + assert!( + ok_result.is_ok(), + "small payload must pass preflight, got: {ok_result:?}" + ); } #[test] @@ -6077,6 +6168,322 @@ mod tests { } } + // ── Production-path tests: real loop emits request, captures nonce ────── + + /// Full end-to-end production path test for the `ask` decision flow: + /// + /// 1. Script emits a real `session/request_permission` on stdout. + /// 2. The read loop processes it via `handle_permission_request()` — + /// no state is pre-planted. + /// 3. The nonce is captured from the observer. + /// 4. A valid decision is sent through the decision channel. + /// 5. The loop writes the permission response to the script's stdin. + /// 6. The script reads the response and emits the terminal id=999 reply. + /// 7. The loop returns `Ok` — the wire flow completes end-to-end. + #[tokio::test] + async fn ask_production_path_emits_request_captures_nonce_and_delivers_decision() { + // Script: emit permission request, wait for any stdin line (the harness's + // response), then emit the terminal session/prompt response. + let perm_req = r#"{"jsonrpc":"2.0","id":42,"method":"session/request_permission","params":{"sessionId":"sess","requestId":"req-prod","subject":"read a file","options":[{"optionId":"opt-allow","kind":"allow_once","name":"Allow"},{"optionId":"opt-deny","kind":"reject_once","name":"Deny"}]}}"#; + let terminal = r#"{"jsonrpc":"2.0","id":999,"result":{"stopReason":"end_turn"}}"#; + // Print the permission request, wait for one line of stdin (the harness's + // response), then print the terminal response. + let script = format!( + r#"printf '{perm_req}\n'; read -r _resp; printf '{terminal}\n'"#, + perm_req = perm_req, + terminal = terminal, + ); + + let mut client = spawn_script(&script).await; + let config = ResolvedPermissionConfig::resolve(PermissionPolicy::Ask, None).unwrap(); + client.set_permission_config(config); + client.set_owner_pubkey_known(true); + + // Subscribe to the observer BEFORE starting the loop so we capture all events. + let obs = crate::observer::ObserverHandle::in_process(); + let mut obs_rx = obs.subscribe(); + client.set_observer(Some(obs.clone()), 0); + + let (perm_tx, perm_rx) = tokio::sync::mpsc::channel::(8); + client.install_permission_decision_rx(perm_rx); + + // Spawn a task that waits for the observer to emit the actionable acp_read + // (the permission request), then delivers a matching decision. + let decision_task = tokio::spawn(async move { + // Wait for the actionable acp_read from the observer. + let mut found_nonce: Option = None; + loop { + match tokio::time::timeout( + std::time::Duration::from_secs(5), + obs_rx.recv(), + ) + .await + { + Ok(Ok(event)) => { + if event.kind == "acp_read" { + if let Some(auth) = &event.authorization { + if auth.actionable { + found_nonce = Some(auth.request_nonce.clone()); + break; + } + } + } + } + _ => break, + } + } + let nonce = found_nonce.expect("actionable acp_read must be emitted"); + // Deliver a valid decision by the captured nonce. + perm_tx + .send(PermissionDecision { + request_nonce: nonce, + option_id: "opt-allow".to_string(), + }) + .await + .expect("decision channel must accept"); + }); + + let idle = std::time::Duration::from_secs(5); + let max_dur = std::time::Duration::from_secs(15); + let hard_deadline = tokio::time::Instant::now() + max_dur; + let result = client + .read_until_response_with_idle_timeout("sess", 999, idle, hard_deadline, max_dur) + .await; + + assert!( + result.is_ok(), + "production-path ask loop must succeed after decision is delivered, got: {result:?}" + ); + assert_eq!( + result.unwrap().get("stopReason").and_then(|v| v.as_str()), + Some("end_turn"), + ); + + // Verify the observer emitted an authorized acp_write (the decision response). + let _ = decision_task.await; + let events = obs.snapshot(); + let write_events: Vec<_> = events + .iter() + .filter(|e| e.kind == "acp_write" && e.authorization.is_some()) + .collect(); + assert!( + !write_events.is_empty(), + "observer must emit at least one authorized acp_write after decision applied" + ); + } + + /// Cancel test: asserts exactly one JSON-RPC response per pending id, + /// no replay on subsequent cancel. + #[tokio::test] + async fn cancel_writes_exactly_one_response_per_pending_id_no_replay() { + // Script that stays alive but produces no output (simulates a hung agent). + let mut client = spawn_script("sleep 5").await; + client.set_permission_config( + ResolvedPermissionConfig::resolve(PermissionPolicy::Ask, None).unwrap(), + ); + client.set_owner_pubkey_known(true); + + // Subscribe to observer to capture writes. + let obs = crate::observer::ObserverHandle::in_process(); + client.set_observer(Some(obs.clone()), 0); + + let (_tx, perm_rx) = tokio::sync::mpsc::channel::(8); + client.install_permission_decision_rx(perm_rx); + + // Plant two distinct Pending entries directly — this tests the cancel + // drain path without needing a live protocol exchange. + for i in 0..2u64 { + let hard_deadline = + tokio::time::Instant::now() + std::time::Duration::from_secs(300); + let msg = perm_request(i, default_opts()); + client + .handle_permission_request(&msg, true, hard_deadline) + .await + .expect("ask registration must succeed"); + } + assert_eq!( + client.pending_permissions.len(), + 2, + "two pending entries must be registered before cancel" + ); + client.last_prompt_id = Some(999); + + // First cancel: must drain both entries. + let _ = client + .cancel_with_cleanup_grace("sess-exact-once", std::time::Duration::from_millis(200)) + .await; + assert!( + client.pending_permissions.is_empty(), + "all pending entries must be drained after cancel" + ); + + // Count authorized acp_write events (each must correspond to one drained entry). + let events_after_first = obs.snapshot(); + let write_count_first = events_after_first + .iter() + .filter(|e| e.kind == "acp_write" && e.authorization.is_some()) + .count(); + assert_eq!( + write_count_first, 2, + "cancel must emit exactly one authorized acp_write per pending id (got {write_count_first})" + ); + + // Second cancel on the same client: no pending entries remain, must not + // re-emit any additional acp_write (no replay). + let _ = client + .cancel_with_cleanup_grace("sess-exact-once", std::time::Duration::from_millis(200)) + .await; + let events_after_second = obs.snapshot(); + let write_count_second = events_after_second + .iter() + .filter(|e| e.kind == "acp_write" && e.authorization.is_some()) + .count(); + assert_eq!( + write_count_second, write_count_first, + "second cancel must not emit additional acp_writes (no replay): before={write_count_first}, after={write_count_second}" + ); + } + + /// Paused-time test: the permission deadline fires at exactly 300 seconds, + /// idle is suspended while a Pending entry exists, and capacity recovers + /// after more than eight sequential requests. + #[tokio::test(start_paused = true)] + async fn ask_permission_deadline_idle_suspension_and_capacity_recovery() { + // Script that emits nothing (simulates an agent waiting for permission response). + let mut client = spawn_script("sleep 600").await; + let config = ResolvedPermissionConfig::resolve(PermissionPolicy::Ask, None).unwrap(); + client.set_permission_config(config); + client.set_owner_pubkey_known(true); + let obs = crate::observer::ObserverHandle::in_process(); + client.set_observer(Some(obs.clone()), 0); + let (_tx, perm_rx) = tokio::sync::mpsc::channel::(8); + client.install_permission_decision_rx(perm_rx); + + // ── Part 1: permission deadline fires before idle ──────────────────── + // Register one pending entry. + let msg = perm_request(1, default_opts()); + let hard_deadline = tokio::time::Instant::now() + + std::time::Duration::from_secs(PERMISSION_ASK_TIMEOUT_SECS + 10); + client + .handle_permission_request(&msg, true, hard_deadline) + .await + .expect("ask registration must succeed"); + assert_eq!(client.pending_permissions.len(), 1); + + // Drive the loop with a long idle timeout — idle must be SUSPENDED while + // the permission entry is pending; only the 300s permission deadline fires. + let idle = std::time::Duration::from_secs(5); // would fire immediately without suspension + let max_dur = std::time::Duration::from_secs(PERMISSION_ASK_TIMEOUT_SECS + 10); + let hard_deadline = tokio::time::Instant::now() + max_dur; + + // Advance time to just before the 300s deadline — idle should NOT fire. + tokio::time::advance( + std::time::Duration::from_secs(PERMISSION_ASK_TIMEOUT_SECS - 1), + ) + .await; + // Spawn a task to advance time past the deadline and check the loop exits + // via permission expiry (not idle timeout). + let advance_task = tokio::spawn(async { + tokio::time::advance(std::time::Duration::from_secs(2)).await; + }); + + let result = client + .read_until_response_with_idle_timeout("sess-tdl", 999, idle, hard_deadline, max_dur) + .await; + let _ = advance_task.await; + + // The loop must have processed the expired entry (transitioned to Resolved) + // and then continued. Because the script produces no output, after the + // permission entry expires the idle timeout fires next (5s). + // Either an IdleTimeout or HardTimeout is acceptable — the key check is + // that no PermissionPoisoned or unexpected error occurred AND the entry + // was processed (Resolved or drained). + assert!( + !matches!(result, Err(AcpError::PermissionPoisoned)), + "permission expiry must not poison the process, got: {result:?}" + ); + // After the deadline, the entry must have been transitioned to Resolved. + let entry_state = client.pending_permissions.get("1"); + let was_resolved = entry_state + .map(|e| matches!(e.state, PermissionEntryState::Resolved)) + .unwrap_or(true); // drain on turn exit is also acceptable + assert!( + was_resolved, + "entry must be Resolved or drained after permission deadline, got: {entry_state:?}" + ); + + // ── Part 2: capacity recovery after 8 sequential requests ─────────── + // Clear any stale entries and verify 8+ sequential requests can succeed + // when previous resolved entries are drained between turns. + let mut client2 = spawn_script("sleep 600").await; + client2.set_permission_config( + ResolvedPermissionConfig::resolve(PermissionPolicy::Ask, None).unwrap(), + ); + client2.set_owner_pubkey_known(true); + let obs2 = crate::observer::ObserverHandle::in_process(); + client2.set_observer(Some(obs2), 0); + let (perm_tx2, perm_rx2) = tokio::sync::mpsc::channel::(16); + client2.install_permission_decision_rx(perm_rx2); + + // Send 9 sequential requests, processing each before sending the next. + // The map is bounded at PERMISSION_MAP_CAP = 8, but Resolved entries do + // not count toward the live-entry cap check — only Pending ones do. + // After each decision is applied (Resolved), the next request must succeed. + for i in 0..9u64 { + let hard = tokio::time::Instant::now() + std::time::Duration::from_secs(300); + let msg = perm_request(i + 100, default_opts()); + let result = client2 + .handle_permission_request(&msg, true, hard) + .await; + // If still Pending from prior iterations, the cap check blocks — this + // tests the case after decisions have been applied (Resolved). + // For this sequential test we deliver decisions immediately. + if result.is_ok() && result.unwrap() { + // Entry is now Pending; deliver a decision immediately. + // Capture the nonce from the freshly-inserted entry. + let id_str = (i + 100).to_string(); + let nonce = client2 + .pending_permissions + .get(&id_str) + .map(|e| e.nonce.clone()); + if let Some(nonce) = nonce { + perm_tx2 + .send(PermissionDecision { + request_nonce: nonce, + option_id: "opt-allow".to_string(), + }) + .await + .ok(); + } + } + } + // Drive the loop to process all queued decisions. + let hard = tokio::time::Instant::now() + std::time::Duration::from_secs(10); + let _ = tokio::time::timeout( + std::time::Duration::from_millis(500), + client2.read_until_response_with_idle_timeout( + "sess-cap", + 9999, + std::time::Duration::from_millis(100), + hard, + std::time::Duration::from_secs(10), + ), + ) + .await; + // After processing, no Pending entries should remain (all should be Resolved + // or the map may have been drained). This proves the capacity map doesn't + // permanently block after 8 requests. + let pending_count = client2 + .pending_permissions + .values() + .filter(|e| matches!(e.state, PermissionEntryState::Pending)) + .count(); + assert_eq!( + pending_count, 0, + "no Pending entries must remain after all decisions applied (capacity recovery confirmed)" + ); + } + // ── Pinned §1 (simpler): ask entry registered synchronously ────────────── #[tokio::test] diff --git a/crates/buzz-acp/src/config.rs b/crates/buzz-acp/src/config.rs index 984ee9b24..50cf9038e 100644 --- a/crates/buzz-acp/src/config.rs +++ b/crates/buzz-acp/src/config.rs @@ -129,11 +129,18 @@ pub enum PermissionMode { Default, /// Fully autonomous execution; model-gated (requires `supportsAutoMode`). /// - /// The adapter self-approves all tool calls internally and never emits - /// `session/request_permission`, so this mode is incompatible with - /// `ask` (card never fires) and `reject` (policy is a dead letter while - /// the adapter auto-approves — the inverted-security worst case). - /// Compatible with `allow` (both want unattended approval). + /// `auto` is a model-gated classifier — the adapter self-approves most tool + /// calls internally, but can fall back to forwarding residual + /// `session/request_permission` requests to ACP when the model chooses manual + /// approval for a specific call. It is therefore **not** a hard bypass. + /// + /// Policy compatibility: + /// - `allow + auto` — compatible; both want unattended approval. + /// - `ask + auto` — compatible with a startup warning; residual escalations + /// still surface permission cards, but internally approved calls bypass the + /// ask flow silently. + /// - `reject + auto` — startup contradiction; adapter auto-approves + /// internally while the policy intends to deny — inverted-security worst case. #[value(alias = "auto")] Auto, /// Auto-approve file edits, still ask for other tools. diff --git a/desktop/src/features/agents/ui/activityRenderClasses/LifecycleActivity.tsx b/desktop/src/features/agents/ui/activityRenderClasses/LifecycleActivity.tsx index b71dbeb95..4bafe0070 100644 --- a/desktop/src/features/agents/ui/activityRenderClasses/LifecycleActivity.tsx +++ b/desktop/src/features/agents/ui/activityRenderClasses/LifecycleActivity.tsx @@ -59,7 +59,12 @@ function PermissionDecisionButtons({ channelId: string; options: Array<{ optionId: string; kind: string; label?: string }>; requestNonce: string; - deliveryFailed?: boolean; + /** + * Monotonically increasing failure token from the reducer — incremented on + * every non-`sent` `control_result`. Keying the effect on this number (not a + * boolean) ensures a second failure after a retry also re-enables buttons. + */ + deliveryFailed?: number; }) { const [pending, setPending] = React.useState(null); diff --git a/desktop/src/features/agents/ui/agentSessionTranscript.test.mjs b/desktop/src/features/agents/ui/agentSessionTranscript.test.mjs index 476fd3326..f4e520b0a 100644 --- a/desktop/src/features/agents/ui/agentSessionTranscript.test.mjs +++ b/desktop/src/features/agents/ui/agentSessionTranscript.test.mjs @@ -2289,8 +2289,8 @@ test("buildTranscript_control_result_non_sent_marks_card_delivery_failed", () => assert.ok(card, "permission card must exist"); assert.equal( card.deliveryFailed, - true, - "deliveryFailed must be set on non-sent control_result", + 1, + "deliveryFailed must be 1 after first non-sent control_result", ); // Card must still be actionable so the user can retry. assert.equal( @@ -2300,6 +2300,64 @@ test("buildTranscript_control_result_non_sent_marks_card_delivery_failed", () => ); }); +test("buildTranscript_control_result_second_failure_increments_delivery_failed", () => { + // A second non-`sent` control_result must increment deliveryFailed so the + // useEffect([deliveryFailed]) dependency in PermissionDecisionButtons + // re-fires and re-enables the buttons for a second retry attempt. + const nonce = "nonce-delivery-fail-2"; + const events = [ + makePermissionRequestWithAuth(1, "req-df2", nonce), + // First failure. + { + seq: 2, + timestamp: "2026-07-01T10:00:01.000Z", + kind: "control_result", + agentIndex: 0, + channelId: "ch-1", + sessionId: "session-1", + turnId: "turn-1", + payload: { + type: "permission_decision", + status: "no_active_turn", + requestNonce: nonce, + optionId: "allow_once", + }, + }, + // Second failure (user retried; harness still unavailable). + { + seq: 3, + timestamp: "2026-07-01T10:00:02.000Z", + kind: "control_result", + agentIndex: 0, + channelId: "ch-1", + sessionId: "session-1", + turnId: "turn-1", + payload: { + type: "permission_decision", + status: "channel_closed", + requestNonce: nonce, + optionId: "allow_once", + }, + }, + ]; + const transcript = buildTranscript(events); + + const card = transcript.find( + (i) => i.renderClass === "permission" && i.requestNonce === nonce, + ); + assert.ok(card, "permission card must exist"); + assert.equal( + card.deliveryFailed, + 2, + "deliveryFailed must be 2 after two non-sent control_results — each failure must increment the token", + ); + assert.equal( + card.actionable, + true, + "card must remain actionable after second delivery failure", + ); +}); + test("buildTranscript_control_result_sent_does_not_mark_delivery_failed", () => { // A `control_result` with `sent` status must NOT set deliveryFailed — the // click reached the harness successfully. diff --git a/desktop/src/features/agents/ui/agentSessionTranscript.ts b/desktop/src/features/agents/ui/agentSessionTranscript.ts index f80306103..7dfb7dd8d 100644 --- a/desktop/src/features/agents/ui/agentSessionTranscript.ts +++ b/desktop/src/features/agents/ui/agentSessionTranscript.ts @@ -1213,7 +1213,15 @@ export function processTranscriptEvent( existing.renderClass === "permission" && existing.actionable ) { - replaceItem(d, itemId, { ...existing, deliveryFailed: true }); + replaceItem(d, itemId, { + ...existing, + // Increment the failure token so the effect in + // PermissionDecisionButtons re-fires even when a prior + // failure already set deliveryFailed (a sticky boolean + // value would not change on the second failure and the + // useEffect dependency would not trigger). + deliveryFailed: (existing.deliveryFailed ?? 0) + 1, + }); } } } diff --git a/desktop/src/features/agents/ui/agentSessionTypes.ts b/desktop/src/features/agents/ui/agentSessionTypes.ts index ccdc32daa..39168e0cd 100644 --- a/desktop/src/features/agents/ui/agentSessionTypes.ts +++ b/desktop/src/features/agents/ui/agentSessionTypes.ts @@ -147,12 +147,13 @@ export type TranscriptItem = */ options?: Array<{ optionId: string; kind: string; label?: string }>; /** - * Set to `true` when a `control_result` frame indicates that the last - * `permission_decision` click was not delivered to the harness (status - * was non-`sent`). The `PermissionDecisionButtons` component uses this - * to re-enable buttons so the user can retry without reloading. + * Monotonically increasing token incremented on every `control_result` + * with a non-`sent` delivery status. The `PermissionDecisionButtons` + * component keys its re-enable effect on this value, so a second failure + * after a retry (same boolean value would not re-trigger the effect) + * still re-enables the buttons. `undefined` when no failure has occurred. */ - deliveryFailed?: boolean; + deliveryFailed?: number; } & TranscriptItemIdentity) | ({ id: string; diff --git a/desktop/src/features/agents/useGlobalAgentConfig.ts b/desktop/src/features/agents/useGlobalAgentConfig.ts index 4b90beb43..294427742 100644 --- a/desktop/src/features/agents/useGlobalAgentConfig.ts +++ b/desktop/src/features/agents/useGlobalAgentConfig.ts @@ -19,6 +19,7 @@ const EMPTY_CONFIG: GlobalAgentConfig = { provider: null, model: null, preferred_runtime: null, + permission_policy: null, }; export const globalAgentConfigQueryKey = ["globalAgentConfig"] as const; diff --git a/docs/nips/NIP-AO.md b/docs/nips/NIP-AO.md index f0bfa4917..1b29ab59e 100644 --- a/docs/nips/NIP-AO.md +++ b/docs/nips/NIP-AO.md @@ -125,9 +125,15 @@ below). It is omitted on all other frame kinds. Permission `acp_read` frames (carrying `session/request_permission` calls) always include an `authorization` envelope. The corresponding `acp_write` (the harness -response) also includes an `authorization` envelope when the decision was recorded — +response) also includes an `authorization` envelope correlated by the same nonce — this pairs the challenge and answer in the observer log. +**One-write / one-observe contract.** Each pending permission entry produces at most +one ACP wire write and at most one authorized `acp_write` observer event. The write +and the observer event are always emitted together; if the write fails the observer +event is suppressed. The sole exception is the `uncertain` terminal (see below) in +which neither is emitted. + ### Authorization Envelope When an `acp_read` or `acp_write` frame relates to a `session/request_permission` @@ -137,7 +143,7 @@ call, the `ObserverEvent` carries an `authorization` field: { "requestNonce": "", "actionable": true | false, - "reason": "" | omitted + "reason": "" | omitted } ``` @@ -148,9 +154,21 @@ call, the `ObserverEvent` carries an `authorization` field: silently ignored. If no matching decision arrives before the per-request timeout, the harness fails the request closed. - `actionable`: `true` when the owner can act (policy=`ask`, preflight passed, owner - and observer available). `false` for auto-deny, fail-closed, and downgrade paths. -- `reason`: present only when `actionable` is `false`; explains why the request was - automatically denied. + and observer available). `false` for auto-deny, fail-closed, and terminal outcomes. +- `reason`: present on every `acp_write` authorization envelope. Identifies the + terminal outcome for this request. Defined values: + + | Value | Meaning | + |-------|---------| + | `"applied"` | Owner decision was received and written to the agent pipe. | + | `"timed_out"` | No decision arrived before the 300-second per-request deadline; request failed closed (denial). | + | `"cancelled"` | The turn was cancelled while the request was pending; request failed closed (denial). | + + The `uncertain` terminal (cancel arriving while the write is in flight) does NOT + produce an `acp_write` observer event — the process is irrecoverably poisoned and + will be respawned by the pool. Desktop clients MUST NOT expect an `acp_write` for + every `acp_read` they receive; a missing `acp_write` after a `session_resolved` + frame with a poisoned outcome indicates the `uncertain` path. **Nonce binding.** The nonce is bound to the agent, channel, session, turn, request ID, and exact option snapshot at generation time. It MUST NOT be reused across @@ -516,7 +534,10 @@ of decrypted payloads and MUST NOT log them at INFO level or above. Note: `actionable` is `false` on the `acp_write` telemetry frame — the decision has been applied and the card is no longer actionable. `reason: "applied"` is the -standard terminal annotation for a successfully delivered decision. +standard terminal annotation for a successfully delivered decision. When the request +expires without a decision, the harness emits `reason: "timed_out"`. When the turn +is cancelled while the request is pending, the harness emits `reason: "cancelled"`. +If the cancel arrives mid-write (`uncertain`), no `acp_write` frame is emitted at all. ## Reference Implementation