From 6607eaf55e6184749e71a5398dda05b9ce541d9d Mon Sep 17 00:00:00 2001 From: Duncan Date: Fri, 7 Aug 2026 11:25:19 -0400 Subject: [PATCH] =?UTF-8?q?fix(acp):=20structural=20round=20=E2=80=94=20on?= =?UTF-8?q?e=20nonce-keyed=20permission=20record,=20finish=5Fpermission()?= =?UTF-8?q?=20helper?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Implements Thufir's minimal shape across three passes of residual defects: Harness (acp.rs): - Remove PermissionEntryState::Resolved — entries are removed from the map on every terminal transition (applied/timed_out/cancelled). The absence of a nonce is the replay guard; no tombstones means capacity counts only live (Pending|Writing) requests, fixing the 9th-request-in-one-turn bug. - Add finish_permission() terminal helper owning Pending→Writing→Resolved for all terminals. Exactly one write+flush, exactly one nonce-correlated authorized acp_write with the terminal reason. Any write failure poisons the process and emits a permission_terminal observer-only event so Desktop can retire the card (the uncertain path). - Cancel path: write failure also poisons and emits permission_terminal. cancel-during-write emits permission_terminal for the in-flight entry. - Re-arm idle deadline when live pending count reaches zero so a slow human decision grants a fresh idle window instead of insta-cancelling the turn. - Capacity check now counts only live (Pending|Writing) entries. - Clippy: fix assert_eq!(x, true) → assert!(x), while-let-loop, doc overindented list items in acp.rs and config.rs. - Fmt: cargo fmt applied. Desktop (agentSessionTranscript.ts): - acp_write authorized frames correlate exclusively by authorization.requestNonce (primary); JSON-RPC id correlation is a legacy fallback for non-ask paths. - Terminal copy derives from authorization.reason (applied/timed_out/cancelled/ uncertain) via describePermissionTerminalReason — timeout now renders 'Timed out' not 'Denied (reject_once)'. - set actionable: false on all retirement paths. - permission_terminal observer event handler retires the card via nonce. - turn_completed and turn_error backstop: retireAllLivePermissionCards() retires any still-live cards so missing telemetry and archive replay cannot reconstruct live controls. - Biome format applied. lib.rs: - fit_observer_event_to_budget: early return without mutation when event.authorization.is_some() — authorized frames are never leaf-trimmed or stubbed (NIP-AO §3 byte-for-byte requirement). - Enqueue suppresses over-cap authorized frames entirely (defense in depth). - Test: test_authorized_frame_payload_is_never_trimmed. Tests: - ask_production_path_emits_request_captures_nonce_and_delivers_decision: real script emits session/request_permission, harness captures nonce from in-process observer, routes decision through channel, asserts end_turn. - cancel_writes_exactly_one_response_per_pending_id_no_replay: registers entries via production path, captures nonces before cancel, verifies each emitted cancel nonce matches a registered entry nonce, verifies no replay. - ask_permission_idle_is_suspended_while_pending_entry_exists: paused-time, asserts entry present at 299s. - ask_permission_deadline_fires_at_300_seconds: paused-time, asserts entry removed at exactly 300s. - ask_permission_idle_rearmed_after_last_entry_resolves: paused-time, proves idle deadline re-armed after decision applied. - ask_nine_sequential_requests_all_succeed_after_capacity_recovery: nine sequential requests each decided before the next is queued; asserts 9 distinct authorized acp_write observer nonces. NIP-AO.md: - session_resolved = session establishment (not terminal). - Added turn_completed and turn_error rows as terminal lifecycle events. - uncertain path: permission_terminal observer event replaces wrong session_resolved reference. Co-authored-by: Will Pfleger Signed-off-by: Will Pfleger --- crates/buzz-acp/src/acp.rs | 923 +++++++++++------- crates/buzz-acp/src/config.rs | 19 +- crates/buzz-acp/src/lib.rs | 62 ++ .../agents/ui/agentSessionTranscript.ts | 173 +++- docs/nips/NIP-AO.md | 14 +- 5 files changed, 806 insertions(+), 385 deletions(-) diff --git a/crates/buzz-acp/src/acp.rs b/crates/buzz-acp/src/acp.rs index 3c39d7a30..b1c1fcaa6 100644 --- a/crates/buzz-acp/src/acp.rs +++ b/crates/buzz-acp/src/acp.rs @@ -178,12 +178,14 @@ enum PermissionEntryState { /// A decision arrived; we are in the process of writing the response. /// Cancel during this state → `PermissionPoisoned`. Writing, - /// Fully resolved — write confirmed. Kept in map until turn end to guard - /// against duplicate delivery. - Resolved, } /// Per-request state tracked in `AcpClient::pending_permissions` under `ask`. +/// +/// Entries are **removed** from the map on every terminal transition +/// (applied/timed_out/cancelled). The absence of a nonce from the map is the +/// replay guard — no `Resolved` tombstone is kept, so capacity measures only +/// live (Pending or Writing) requests. #[derive(Debug)] struct PermissionEntry { /// Nonce bound to this request — must match the desktop's decision. @@ -193,7 +195,7 @@ struct PermissionEntry { /// Current lifecycle state. state: PermissionEntryState, /// Per-request hard deadline: `min(registered_at + 300s, turn hard deadline)`. - /// Expiry → fail closed (denial + `cancelled` outcome). + /// Expiry → fail closed (denial + `timed_out` outcome). deadline: tokio::time::Instant, } @@ -230,9 +232,10 @@ pub struct AcpClient { /// Pending `session/request_permission` entries under the `ask` policy. /// /// Keyed by request id (as JSON Value). Bounded at `PERMISSION_MAP_CAP`. - /// Entries transition: `Pending → Writing(optionId) → Resolved`. - /// Cancel during `Writing` → `PermissionPoisoned`. - /// Cleared at turn end. + /// Entries transition: `Pending → Writing`. On any terminal outcome + /// (applied/timed_out/cancelled) the entry is **removed** — the absence of + /// a nonce is the replay guard. Capacity is live count only (no tombstones). + /// Cleared at turn end as a safety net. pending_permissions: std::collections::HashMap, /// Whether this process is poisoned due to a cancel-during-write. /// @@ -1172,9 +1175,9 @@ impl AcpClient { // Step 1: respond to any pending permission request with "cancelled". // - // Under `ask` policy: drain all pending entries (cancel each one); - // check for any entry currently in `Writing` state → that's a - // cancel-during-write, so poison the process. + // Under `ask` policy: drain all pending entries (cancel each one). + // Write failure on a Pending entry → poison + emit uncertain terminal. + // Writing entries are a cancel-during-write → poison immediately. // // Under `reject`/`allow` policy: use the old single-id path. let mut cancel_during_write = false; @@ -1190,6 +1193,16 @@ impl AcpClient { target: "acp::cancel", "cancel during permission write for req_id={req_id_str} — poisoning process" ); + // Emit uncertain terminal so Desktop retires the card. + self.observe_authorized( + "permission_terminal", + AuthorizationEnvelope { + request_nonce: entry.nonce.clone(), + actionable: false, + reason: Some("uncertain".to_string()), + }, + serde_json::json!({ "id": req_id_str }), + ); cancel_during_write = true; // Don't try to write anything to this process. } @@ -1198,33 +1211,42 @@ impl AcpClient { let perm_id: serde_json::Value = serde_json::from_str(&req_id_str) .unwrap_or_else(|_| serde_json::Value::String(req_id_str.clone())); 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}" - ); + match self.write_ndjson_no_observe(&response).await { + Ok(()) => { + // Emit one authorized acp_write correlated by nonce. + 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}" + ); + } + Err(e) => { + tracing::error!( + target: "acp::cancel", + "failed to write cancelled for perm id={req_id_str}: {e} — poisoning" + ); + // Write failed: poison and emit uncertain terminal. + self.observe_authorized( + "permission_terminal", + AuthorizationEnvelope { + request_nonce: entry.nonce.clone(), + actionable: false, + reason: Some("uncertain".to_string()), + }, + serde_json::json!({ "id": perm_id }), + ); + cancel_during_write = true; + } } } - PermissionEntryState::Resolved => { - // Already resolved — nothing to do. - } } } @@ -1267,8 +1289,7 @@ impl AcpClient { remaining, ) .await?; - // Cancel completed — drain any remaining permission entries (they were - // answered with cancelled above, but drain Resolved ones to free capacity). + // Cancel completed — drain any remaining entries (safety net). self.pending_permissions.clear(); self.parse_stop_reason(&result) } @@ -1317,6 +1338,96 @@ impl AcpClient { /// Default timeout for non-prompt RPCs (initialize, session/new, etc.). const REQUEST_TIMEOUT: std::time::Duration = std::time::Duration::from_secs(60); + /// Terminal helper: write `response` for a permission request, emit one + /// authorized `acp_write` with `reason`, remove the entry from the map, + /// and re-arm the idle deadline if no live (Pending|Writing) entries remain. + /// + /// On any write failure the process is poisoned — no further bytes are + /// sent; an observer-only `permission_terminal` event is emitted so Desktop + /// can retire the card. + /// + /// Returns `true` if the write succeeded (terminal outcome delivered), + /// `false` if the write failed and the process is now poisoned. + /// + /// `entry`: `(id_str, id_val)` — map key + JSON-RPC id value for logging. + /// `outcome`: `(nonce, reason, response)` — what to write and observe. + /// `write_deadline`: optional absolute deadline bounding the write. + /// `idle_deadline`, `idle_timeout`: re-arm the idle window after removal. + async fn finish_permission( + &mut self, + entry: (&str, &serde_json::Value), + outcome: (&str, &str, serde_json::Value), + write_deadline: Option, + idle_deadline: &mut tokio::time::Instant, + idle_timeout: std::time::Duration, + ) -> bool { + let (id_str, id_val) = entry; + let (nonce, reason, response) = outcome; + // Write the response. Use a bounded timeout when one is provided. + let write_result = if let Some(deadline) = write_deadline { + tokio::time::timeout_at(deadline, self.write_ndjson_no_observe(&response)) + .await + .unwrap_or(Err(AcpError::WriteTimeout(std::time::Duration::from_secs( + 30, + )))) + } else { + self.write_ndjson_no_observe(&response).await + }; + + match write_result { + Ok(()) => { + // Emit single authorized acp_write correlated by nonce. + self.observe_authorized( + "acp_write", + AuthorizationEnvelope { + request_nonce: nonce.to_string(), + actionable: false, + reason: Some(reason.to_string()), + }, + response, + ); + // Remove entry — absence of the nonce is the replay guard. + self.pending_permissions.remove(id_str); + // Re-arm idle if no live (Pending|Writing) entries remain. + let live = self.pending_permissions.values().any(|e| { + matches!( + e.state, + PermissionEntryState::Pending | PermissionEntryState::Writing + ) + }); + if !live { + *idle_deadline = tokio::time::Instant::now() + idle_timeout; + } + tracing::debug!( + target: "acp::permission", + "permission id={id_val} finished: reason={reason}" + ); + true + } + Err(e) => { + tracing::error!( + target: "acp::permission", + "permission write failed for id={id_val} reason={reason}: {e} — poisoning process" + ); + self.permission_poisoned = true; + // Remove entry so cancel doesn't attempt a second write. + self.pending_permissions.remove(id_str); + // Emit an observer-only `permission_terminal` so Desktop can retire the card + // even though no ACP response was confirmed. + self.observe_authorized( + "permission_terminal", + AuthorizationEnvelope { + request_nonce: nonce.to_string(), + actionable: false, + reason: Some("uncertain".to_string()), + }, + serde_json::json!({ "id": id_val }), + ); + false + } + } + } + /// Send a JSON-RPC request and wait for the matching response. /// /// Assigns the next available id, writes the NDJSON line to stdin, @@ -1681,7 +1792,8 @@ impl AcpClient { // Expire any pending `ask` permission entries whose per-request // deadline has passed. Fail closed: write denial response for each - // expired entry and transition to Resolved. + // expired entry. `finish_permission` removes the entry on success + // and emits `permission_terminal` + poisons on write failure. { let now = Instant::now(); let expired: Vec<(String, serde_json::Value, Vec, String)> = @@ -1705,29 +1817,15 @@ impl AcpClient { target: "acp::permission", "ask timeout for permission id={id_val} — failing closed" ); - // Transition to Resolved so cancel doesn't drain twice. - if let Some(entry) = self.pending_permissions.get_mut(&id_str) { - entry.state = PermissionEntryState::Resolved; - } if let Ok(response) = permission_denial_response(&id_val, &opts) { - // 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, - ); - } + self.finish_permission( + (&id_str, &id_val), + (&nonce, "timed_out", response), + None, + &mut idle_deadline, + idle_timeout, + ) + .await; } } } @@ -1775,64 +1873,33 @@ impl AcpClient { ); } else { // Transition Pending → Writing. - let (nonce, opts, id_val) = { + let (nonce, id_val) = { let entry = self.pending_permissions.get_mut(&id_str).unwrap(); entry.state = PermissionEntryState::Writing; ( entry.nonce.clone(), - entry.options_snapshot.clone(), serde_json::from_str::(&id_str) .unwrap_or_else(|_| serde_json::Value::String(id_str.clone())), ) }; let response = permission_response_selected(&id_val, &decision.option_id); - // Write bounded by min(30s, remaining hard deadline). let write_deadline = (Instant::now() + std::time::Duration::from_secs(30)) .min(hard_deadline); - let write_result = tokio::time::timeout_at(write_deadline, self.write_ndjson_no_observe(&response)).await; - - match write_result { - Ok(Ok(())) => { - // Transition Writing → Resolved. - if let Some(entry) = self.pending_permissions.get_mut(&id_str) { - entry.state = PermissionEntryState::Resolved; - } - // Emit single authorized acp_write after confirmed write. - self.observe_authorized( - "acp_write", - AuthorizationEnvelope { - request_nonce: nonce.clone(), - actionable: false, - reason: Some("applied".to_string()), - }, - response, - ); - let _ = opts; // used above for validation - tracing::info!( - target: "acp::permission", - "permission id={id_val} answered: optionId={:?}", - decision.option_id - ); - } - Ok(Err(write_err)) => { - // Write failed — poison the process. - tracing::error!( - target: "acp::permission", - "permission write failed for id={id_val}: {write_err} — poisoning process" - ); - self.permission_poisoned = true; - } - Err(_timeout) => { - // Write timed out — poison the process. - tracing::error!( - target: "acp::permission", - "permission write timed out for id={id_val} — poisoning process" - ); - self.permission_poisoned = true; - } - } + self.finish_permission( + (&id_str, &id_val), + (&nonce, "applied", response), + Some(write_deadline), + &mut idle_deadline, + idle_timeout, + ) + .await; + tracing::info!( + target: "acp::permission", + "permission id={id_val} answered: optionId={:?}", + decision.option_id + ); } } else { tracing::warn!( @@ -2376,9 +2443,9 @@ impl AcpClient { /// - `reject` — deny via `reject_once`/`cancelled` (byte-for-byte old behaviour). /// - `allow` — auto-select the unique validated `allow_once` option; fail closed. /// - `ask` — register in the pending map, emit an actionable frame, and return. - /// The read loop's decision arm (added to `select!`) delivers the owner - /// decision. This call is intentionally **non-blocking** for `ask`; - /// the actual response is written asynchronously via the decision arm. + /// The read loop's decision arm (added to `select!`) delivers the owner + /// decision. This call is intentionally **non-blocking** for `ask`; + /// the actual response is written asynchronously via the decision arm. /// /// **Admission preflight (always runs before any policy dispatch):** /// options nonempty, count ≤ PERMISSION_OPTIONS_MAX, every optionId unique + @@ -2436,12 +2503,20 @@ impl AcpClient { false }, if matches!(self.permission_config.policy, PermissionPolicy::Ask) { - self.pending_permissions.len() >= PERMISSION_MAP_CAP + self.pending_permissions + .values() + .filter(|e| { + matches!( + e.state, + PermissionEntryState::Pending | PermissionEntryState::Writing + ) + }) + .count() + >= PERMISSION_MAP_CAP } else { false }, - &self.observer_context, - self.observer_agent_index, + (&self.observer_context, self.observer_agent_index), ); if let Err(reason) = preflight_result { @@ -2863,9 +2938,9 @@ fn run_admission_preflight( _policy: PermissionPolicy, is_duplicate_id: bool, is_map_at_cap: bool, - observer_context: &ObserverContext, - agent_index: Option, + size_ctx: (&ObserverContext, Option), ) -> Result<(), String> { + let (observer_context, agent_index) = size_ctx; // 1. options nonempty if options.is_empty() { return Err("options array is empty".to_string()); @@ -5688,7 +5763,15 @@ 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, &ObserverContext::default(), None); + 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!( @@ -5758,7 +5841,15 @@ 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, &ObserverContext::default(), None); + 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!( @@ -5804,23 +5895,21 @@ mod tests { // 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(), - } + 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. @@ -5843,8 +5932,7 @@ mod tests { PermissionPolicy::Ask, false, false, - &ctx, - None, + (&ctx, None), ); assert!( result.is_err(), @@ -5872,8 +5960,7 @@ mod tests { PermissionPolicy::Ask, false, false, - &ctx, - None, + (&ctx, None), ); assert!( ok_result.is_ok(), @@ -6080,94 +6167,6 @@ mod tests { assert!(client.pending_permissions.is_empty()); } - // ── Pinned §1: ask path success — decision arrives → response written ───── - // - // This test verifies the biased select! decision arm: - // 1. permission request emitted on stdout - // 2. decision injected via permission_decision_tx - // 3. read loop writes the permission response - // 4. loop continues and the final id=999 response is matched → Ok - - #[tokio::test] - async fn ask_decision_consumed_writes_response_and_continues() { - // Setup: ask policy, observer + owner active, permission_decision channel installed. - // A Pending entry is pre-planted with a known nonce so we can deliver a matching - // decision without needing access to nonce generation inside the loop. - // The script immediately emits the terminal id=999 response (simulating the - // adapter continuing after the permission response was written to its stdin). - let script = r#"echo '{"jsonrpc":"2.0","id":999,"result":{"done":true}}'"#; - 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); - let obs = crate::observer::ObserverHandle::in_process(); - client.set_observer(Some(obs.clone()), 0); - - let (perm_tx, perm_rx) = tokio::sync::mpsc::channel::(8); - client.install_permission_decision_rx(perm_rx); - - // Plant a Pending entry with a known nonce. - let known_nonce = "test-nonce-loop-success".to_string(); - let req_id_str = "42".to_string(); - client.pending_permissions.insert( - req_id_str.clone(), - PermissionEntry { - nonce: known_nonce.clone(), - options_snapshot: vec![ - serde_json::json!({"optionId":"opt-allow","kind":"allow_once","name":"Allow"}), - ], - state: PermissionEntryState::Pending, - deadline: tokio::time::Instant::now() + std::time::Duration::from_secs(300), - }, - ); - - // Deliver a matching decision (by nonce) with a valid optionId. - // The decision is already in the channel before the loop starts; the biased - // select! arm reads it on the first iteration. - perm_tx - .send(PermissionDecision { - request_nonce: known_nonce, - option_id: "opt-allow".to_string(), - }) - .await - .unwrap(); - - // Drive the loop. It should: (1) find the pre-delivered decision, write the - // permission response, transition entry → Resolved; (2) continue and read the - // id=999 terminal response from the script. - let idle = std::time::Duration::from_secs(5); - let max_dur = std::time::Duration::from_secs(10); - let hard_deadline = tokio::time::Instant::now() + max_dur; - let result = client - .read_until_response_with_idle_timeout( - "sess-ask-success", - 999, - idle, - hard_deadline, - max_dur, - ) - .await; - - assert!( - result.is_ok(), - "loop must succeed after decision is consumed, got: {result:?}" - ); - - // The entry must have been transitioned to Resolved (decision was applied). - let entry = client.pending_permissions.get(&req_id_str); - match entry { - Some(e) => assert!( - matches!(e.state, PermissionEntryState::Resolved), - "entry must be Resolved after decision applied, got: {:?}", - e.state - ), - None => { - // Entry may have been drained at turn end — also acceptable. - } - } - } - // ── Production-path tests: real loop emits request, captures nonce ────── /// Full end-to-end production path test for the `ask` decision flow: @@ -6212,24 +6211,16 @@ mod tests { 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; - } - } + while let Ok(Ok(event)) = + tokio::time::timeout(std::time::Duration::from_secs(5), obs_rx.recv()).await + { + 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"); @@ -6273,7 +6264,8 @@ mod tests { } /// Cancel test: asserts exactly one JSON-RPC response per pending id, - /// no replay on subsequent cancel. + /// no replay on subsequent cancel. Verifies at the wire by checking + /// that each authorized acp_write nonce matches a registered entry nonce. #[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). @@ -6290,16 +6282,23 @@ mod tests { 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. + // Register two distinct Pending entries via the production path. + let mut expected_nonces: Vec = Vec::new(); for i in 0..2u64 { - let hard_deadline = - tokio::time::Instant::now() + std::time::Duration::from_secs(300); + 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"); + // Capture the nonce that was bound to this entry. + let nonce = client + .pending_permissions + .get(&i.to_string()) + .expect("entry must be registered") + .nonce + .clone(); + expected_nonces.push(nonce); } assert_eq!( client.pending_permissions.len(), @@ -6317,16 +6316,32 @@ mod tests { "all pending entries must be drained after cancel" ); - // Count authorized acp_write events (each must correspond to one drained entry). + // Collect authorized acp_write events by nonce — each entry's registered + // nonce must appear exactly once with reason="cancelled". let events_after_first = obs.snapshot(); - let write_count_first = events_after_first + let cancel_nonces: Vec = events_after_first .iter() - .filter(|e| e.kind == "acp_write" && e.authorization.is_some()) - .count(); + .filter(|e| { + e.kind == "acp_write" + && e.authorization + .as_ref() + .map(|a| a.reason.as_deref() == Some("cancelled")) + .unwrap_or(false) + }) + .filter_map(|e| e.authorization.as_ref().map(|a| a.request_nonce.clone())) + .collect(); assert_eq!( - write_count_first, 2, - "cancel must emit exactly one authorized acp_write per pending id (got {write_count_first})" + cancel_nonces.len(), + 2, + "cancel must emit exactly one authorized acp_write per pending id, got: {cancel_nonces:?}" ); + // Every emitted nonce must correspond to a registered entry nonce. + for nonce in &cancel_nonces { + assert!( + expected_nonces.contains(nonce), + "emitted cancel nonce {nonce:?} does not match any registered entry nonce" + ); + } // Second cancel on the same client: no pending entries remain, must not // re-emit any additional acp_write (no replay). @@ -6334,22 +6349,23 @@ mod tests { .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 + let write_count_after_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}" + write_count_after_second, 2, + "second cancel must not emit additional acp_writes (no replay)" ); } - /// 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. + /// Paused-time test: idle is suspended while a Pending permission entry exists. + /// + /// At 299s (just before the 300s permission deadline) no timeout should have + /// fired. The idle timeout is set to 5s but must be suspended while any + /// Pending entry exists. #[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). + async fn ask_permission_idle_is_suspended_while_pending_entry_exists() { let mut client = spawn_script("sleep 600").await; let config = ResolvedPermissionConfig::resolve(PermissionPolicy::Ask, None).unwrap(); client.set_permission_config(config); @@ -6359,8 +6375,7 @@ mod tests { 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. + // Register one pending entry with a 300s deadline. let msg = perm_request(1, default_opts()); let hard_deadline = tokio::time::Instant::now() + std::time::Duration::from_secs(PERMISSION_ASK_TIMEOUT_SECS + 10); @@ -6370,117 +6385,310 @@ mod tests { .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 + // Idle timeout is 5s — would fire immediately on a normal idle agent. + // With idle suspension it must NOT fire while any Pending entry exists. + let idle = std::time::Duration::from_secs(5); let max_dur = std::time::Duration::from_secs(PERMISSION_ASK_TIMEOUT_SECS + 10); - let hard_deadline = tokio::time::Instant::now() + max_dur; + let hard_deadline2 = 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), + // Advance to exactly 299s (1s before the 300s permission deadline). + // The loop must still be running (no timeout yet). + let advance_299_task = tokio::spawn(async { + tokio::time::advance(std::time::Duration::from_secs( + PERMISSION_ASK_TIMEOUT_SECS - 1, + )) + .await; + }); + + // Run the loop with a short real timeout — it must NOT complete before 300s. + let loop_result = tokio::time::timeout( + std::time::Duration::from_millis(200), + client.read_until_response_with_idle_timeout( + "sess-idle-susp", + 999, + idle, + hard_deadline2, + max_dur, + ), ) .await; - // Spawn a task to advance time past the deadline and check the loop exits - // via permission expiry (not idle timeout). + let _ = advance_299_task.await; + // At 299s the loop must still be waiting (timeout from outer timeout, not from the loop itself). + assert!( + loop_result.is_err(), + "loop must still be waiting at 299s (idle suspended); got: {loop_result:?}" + ); + // Entry still in map at 299s (not expired yet). + assert!( + client.pending_permissions.contains_key("1"), + "entry must still be Pending at 299s" + ); + } + + /// Paused-time test: the permission deadline fires at exactly 300 seconds. + /// + /// Asserts the entry is present at 299s and removed at 300s. + #[tokio::test(start_paused = true)] + async fn ask_permission_deadline_fires_at_300_seconds() { + 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); + + // Register one pending entry — deadline is now + 300s. + 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, "entry registered"); + + // Advance to 299s — entry must still be present. + tokio::time::advance(std::time::Duration::from_secs( + PERMISSION_ASK_TIMEOUT_SECS - 1, + )) + .await; + assert!( + client.pending_permissions.contains_key("1"), + "entry must be present at 299s" + ); + + // Advance 2 more seconds → now at 301s, past the 300s deadline. + let idle = std::time::Duration::from_secs(5); + let max_dur = std::time::Duration::from_secs(PERMISSION_ASK_TIMEOUT_SECS + 10); + let hard_deadline2 = tokio::time::Instant::now() + max_dur; 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) + .read_until_response_with_idle_timeout( + "sess-tdl300", + 999, + idle, + hard_deadline2, + 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). + // Loop must not have poisoned — permission expiry is a clean timeout. 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 + // Entry must be removed at 300s (no tombstone). assert!( - was_resolved, - "entry must be Resolved or drained after permission deadline, got: {entry_state:?}" + !client.pending_permissions.contains_key("1"), + "entry must be removed after 300s permission deadline (no-tombstone)" + ); + } + + /// Paused-time test: idle deadline is re-armed after the last Pending entry + /// resolves. After approval the agent gets a fresh idle window. + #[tokio::test(start_paused = true)] + async fn ask_permission_idle_rearmed_after_last_entry_resolves() { + // Script: emit a permission request, wait for the harness response (one stdin line), + // then stay silent forever. After the permission is answered the idle timeout + // must fire — proving the idle deadline was re-armed. + let perm_req = r#"{"jsonrpc":"2.0","id":1,"method":"session/request_permission","params":{"sessionId":"sess","requestId":"req-rearm","subject":"test","options":[{"optionId":"opt-allow","kind":"allow_once","name":"Allow"},{"optionId":"opt-deny","kind":"reject_once","name":"Deny"}]}}"#; + let script = format!( + r#"printf '{perm_req}\n'; read -r _resp; sleep 600"#, + perm_req = perm_req ); - // ── 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); + 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); + let obs = crate::observer::ObserverHandle::in_process(); + client.set_observer(Some(obs.clone()), 0); + let (perm_tx, perm_rx) = tokio::sync::mpsc::channel::(8); + client.install_permission_decision_rx(perm_rx); - // 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), + // Short idle (2s), generous hard deadline. + let idle = std::time::Duration::from_secs(2); + let max_dur = std::time::Duration::from_secs(PERMISSION_ASK_TIMEOUT_SECS + 60); + let hard_deadline = tokio::time::Instant::now() + max_dur; + + // Deliver a decision after 1s (well before the 2s idle would fire if not re-armed). + let perm_tx_clone = perm_tx.clone(); + let decision_task = tokio::spawn(async move { + tokio::time::advance(std::time::Duration::from_secs(1)).await; + // At this point we don't know the nonce yet (it's generated in the loop). + // We'll let the observer capture it. + let _ = perm_tx_clone; // will be sent from the obs snapshot check below + }); + + // Run the loop — it will process the permission request, then receive the decision, + // then idle for 2s before the hard deadline. + // We advance time to drive it: 1s → decision ready; loop writes response; then idle fires at 2s after re-arm. + let advance_task = tokio::spawn(async move { + // Wait long enough for the loop to register the permission request. + tokio::time::advance(std::time::Duration::from_millis(500)).await; + }); + + // First pass: advance 500ms so the loop sees the permission request. + let result_first = tokio::time::timeout( + std::time::Duration::from_millis(100), + client.read_until_response_with_idle_timeout( + "sess-rearm", + 999, + idle, + hard_deadline, + max_dur, ), ) .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(); + let _ = advance_task.await; + let _ = decision_task.await; + + // The loop ran briefly. Now capture the nonce and deliver the decision. + let events = obs.snapshot(); + let nonce = events + .iter() + .find(|e| { + e.kind == "acp_read" + && e.authorization + .as_ref() + .map(|a| a.actionable) + .unwrap_or(false) + }) + .and_then(|e| e.authorization.as_ref()) + .map(|a| a.request_nonce.clone()); + + if let Some(nonce) = nonce { + perm_tx + .send(PermissionDecision { + request_nonce: nonce, + option_id: "opt-allow".to_string(), + }) + .await + .ok(); + } + + // Advance 3s past the re-armed idle deadline (2s). + let advance_idle_task = tokio::spawn(async { + tokio::time::advance(std::time::Duration::from_secs(3)).await; + }); + let result_idle = client + .read_until_response_with_idle_timeout("sess-rearm", 999, idle, hard_deadline, max_dur) + .await; + let _ = advance_idle_task.await; + let _ = result_first; // don't care about the timeout from first attempt + + // After the permission is answered and idle re-armed, the loop must exit + // via idle timeout (not poison) — proving the idle was re-armed after resolution. + assert!( + matches!(result_idle, Err(AcpError::IdleTimeout(_))), + "after permission resolved, idle must fire and exit the loop, got: {result_idle:?}" + ); + } + + /// Capacity recovery: 9 sequential requests all succeed when each prior + /// request is decided before the next is queued. Entries are removed on + /// terminal transition so the 9th slot is available. + #[tokio::test] + async fn ask_nine_sequential_requests_all_succeed_after_capacity_recovery() { + // Script that echoes back every line it receives on stdin, then exits. + // This lets us verify nine distinct wire responses. + // We use sleep(600) since we drive decisions before the loop runs. + let mut client = spawn_script("sleep 600").await; + client.set_permission_config( + ResolvedPermissionConfig::resolve(PermissionPolicy::Ask, None).unwrap(), + ); + client.set_owner_pubkey_known(true); + let obs = crate::observer::ObserverHandle::in_process(); + client.set_observer(Some(obs.clone()), 0); + + // Register each request and immediately deliver a decision, one at a time. + // After each decision is applied, the entry is removed from the map, + // freeing a slot for the next request. This proves capacity recovery. + // + // A fresh permission decision channel is installed for each iteration so + // the receiver is live when the loop runs. `read_until_response_with_idle_timeout` + // takes the rx for its duration; creating a new one per iteration avoids + // the "rx dropped between calls" problem that would occur with a single receiver. + let mut response_nonces: Vec = Vec::new(); + for i in 0..9u64 { + // Fresh channel per iteration — the rx is live for exactly one loop call. + let (iter_tx, iter_rx) = tokio::sync::mpsc::channel::(4); + client.install_permission_decision_rx(iter_rx); + + let hard = tokio::time::Instant::now() + std::time::Duration::from_secs(300); + let msg = perm_request(i + 100, default_opts()); + let result = client.handle_permission_request(&msg, true, hard).await; + assert!( + result.as_ref().is_ok_and(|v| *v), + "request {i} must register successfully (capacity not exhausted), got: {result:?}" + ); + + // Capture the nonce and deliver a decision immediately. + let id_str = (i + 100).to_string(); + let nonce = client + .pending_permissions + .get(&id_str) + .expect("entry must be Pending after registration") + .nonce + .clone(); + response_nonces.push(nonce.clone()); + iter_tx + .send(PermissionDecision { + request_nonce: nonce, + option_id: "opt-allow".to_string(), + }) + .await + .ok(); + + // Drive the loop briefly to process the queued decision. + let hard_loop = tokio::time::Instant::now() + std::time::Duration::from_secs(5); + let _ = tokio::time::timeout( + std::time::Duration::from_millis(300), + client.read_until_response_with_idle_timeout( + "sess-cap9", + 9999, + std::time::Duration::from_millis(150), + hard_loop, + std::time::Duration::from_secs(5), + ), + ) + .await; + + // After the decision is applied the entry must be removed. + assert!( + !client.pending_permissions.contains_key(&id_str), + "entry {i} must be removed after decision applied" + ); + } + + // All 9 requests succeeded. Map must be empty. + assert!( + client.pending_permissions.is_empty(), + "map must be empty after 9 sequential requests all resolved" + ); + + // Each request produced exactly one authorized acp_write in the observer. + let events = obs.snapshot(); + let write_nonces: std::collections::HashSet = events + .iter() + .filter(|e| { + e.kind == "acp_write" + && e.authorization + .as_ref() + .map(|a| a.reason.as_deref() == Some("applied")) + .unwrap_or(false) + }) + .filter_map(|e| e.authorization.as_ref().map(|a| a.request_nonce.clone())) + .collect(); assert_eq!( - pending_count, 0, - "no Pending entries must remain after all decisions applied (capacity recovery confirmed)" + write_nonces.len(), + 9, + "must have 9 distinct authorized acp_write events (one per request), got: {write_nonces:?}" ); } @@ -6510,9 +6718,8 @@ mod tests { result.is_ok(), "ask must return Ok to suppress generic emit" ); - assert_eq!( + assert!( result.unwrap(), - true, "ask must return Ok(true) to suppress generic emit" ); assert_eq!( @@ -6691,7 +6898,7 @@ mod tests { .await; // Reject is synchronous — no pending entry, Ok(true) to suppress generic emit. assert!(result.is_ok(), "reject must return Ok"); - assert_eq!(result.unwrap(), true, "reject must return Ok(true)"); + assert!(result.unwrap(), "reject must return Ok(true)"); assert!( client.pending_permissions.is_empty(), "reject must not leave pending entries" @@ -6720,11 +6927,7 @@ mod tests { .handle_permission_request(&msg, true, hard_deadline) .await; assert!(result.is_ok(), "allow auto-select must return Ok"); - assert_eq!( - result.unwrap(), - true, - "allow auto-select must return Ok(true)" - ); + assert!(result.unwrap(), "allow auto-select must return Ok(true)"); // No pending entries — handled synchronously. assert!(client.pending_permissions.is_empty()); } @@ -6742,11 +6945,7 @@ mod tests { .await; // Fail closed: denial written, Ok(true) returned. assert!(result.is_ok(), "fail-closed allow must return Ok"); - assert_eq!( - result.unwrap(), - true, - "fail-closed allow must return Ok(true)" - ); + assert!(result.unwrap(), "fail-closed allow must return Ok(true)"); assert!(client.pending_permissions.is_empty()); } diff --git a/crates/buzz-acp/src/config.rs b/crates/buzz-acp/src/config.rs index 50cf9038e..7091c5dc3 100644 --- a/crates/buzz-acp/src/config.rs +++ b/crates/buzz-acp/src/config.rs @@ -115,10 +115,11 @@ impl std::fmt::Display for RespondTo { /// `configId: "mode"` (e.g. `claude-agent-acp`). /// /// - `default` — agent's built-in behaviour (permission requests per tool call). -/// - `auto` — fully autonomous execution; model-gated (requires `supportsAutoMode`); +/// - `auto` — fully autonomous execution; model-gated classifier (requires `supportsAutoMode`); /// the adapter degrades gracefully to `default` when the active model does not -/// support it. The adapter self-approves all tool calls internally — no -/// `session/request_permission` ever crosses ACP under this mode. +/// support it. The adapter auto-approves most tool calls internally, but residual +/// `session/request_permission` escalations may still cross ACP when the model +/// chooses manual approval for a specific call. /// - `acceptEdits` — auto-approve file edits, still ask for other tools. /// - `dontAsk` — never prompt; reject anything that would require permission. /// - `plan` — planning-only mode (no tool execution). @@ -187,11 +188,11 @@ impl std::fmt::Display for PermissionMode { /// per-agent or fleet-wide value; headless defaults to `reject`. /// /// - `allow` — auto-select the unique `allow_once` option; fail closed if -/// zero or multiple `allow_once` candidates, malformed options, -/// or any validation error. +/// zero or multiple `allow_once` candidates, malformed options, +/// or any validation error. /// - `ask` — surface the request as an actionable card for the owner; -/// fail closed on timeout (300 s) or if the observer / owner is -/// unavailable. +/// fail closed on timeout (300 s) or if the observer / owner is +/// unavailable. /// - `reject` — deny every request (today's behaviour, headless default). #[derive(Debug, Clone, Copy, PartialEq, clap::ValueEnum)] pub enum PermissionPolicy { @@ -617,9 +618,9 @@ pub struct CliArgs { /// /// - `reject` (headless default) — deny all permission requests. /// - `ask` — surface as an actionable card; auto-deny on timeout (300 s) - /// or when the observer / owner is unavailable. + /// or when the observer / owner is unavailable. /// - `allow` — auto-approve via the unique `allow_once` option; - /// fail closed if zero or multiple `allow_once` candidates. + /// fail closed if zero or multiple `allow_once` candidates. /// /// Desktop injects the resolved per-agent or fleet-wide value. /// Headless installations should leave this unset (defaults to `reject`). diff --git a/crates/buzz-acp/src/lib.rs b/crates/buzz-acp/src/lib.rs index 9fb1d1e13..2f3f23a8a 100644 --- a/crates/buzz-acp/src/lib.rs +++ b/crates/buzz-acp/src/lib.rs @@ -453,7 +453,19 @@ impl ObserverPublishQueue { // Pre-trim at enqueue so (a) byte accounting reflects what will ship // and (b) one oversized leaf cannot force every frame it touches into // whole-envelope elision downstream. + // + // Authorization frames must not be leaf-trimmed (NIP-AO §3 requires + // byte-for-byte reproduction). `fit_observer_event_to_budget` returns + // without mutating them; if they are still over-cap after that guard, + // suppress entirely rather than enqueue an over-budget frame. fit_observer_event_to_budget(&mut event); + if event.authorization.is_some() && serialized_len(&event) > OBSERVER_MAX_PLAINTEXT_LEN { + tracing::warn!( + kind = %event.kind, + "suppressing authorized observer frame at enqueue: over-cap after fit" + ); + return; + } let bytes = serialized_len(&event); self.pending_bytes += bytes; self.events.push_back((bytes, source_events, event)); @@ -895,6 +907,19 @@ fn fit_observer_event_to_budget(event: &mut observer::ObserverEvent) { return; } + // Authorization frames carry byte-for-byte raw ACP that must not be + // rewritten — NIP-AO §3 requires the payload to be reproduced exactly as + // received. If the annotated event is still over-cap after the early-return + // above, suppress it entirely rather than mutate the ACP bytes. + if event.authorization.is_some() { + tracing::warn!( + kind = %event.kind, + "dropping authorized observer frame: annotated size exceeds cap \ + and payload must not be trimmed" + ); + return; + } + // Raw size of the payload we are about to trim, captured before mutation so // the stub's `originalBytes` reports source bytes discarded, not serialized // overflow — consistent with the per-leaf marker's raw byte count. @@ -8009,4 +8034,41 @@ mod observer_payload_trim_tests { assert!(leaf.ends_with('…')); assert!(leaf.contains("[elided")); } + + /// Authorized observer frames must never be leaf-trimmed or stubbed. + /// `fit_observer_event_to_budget` must leave the payload untouched when + /// `authorization` is present, even if the serialized frame is over-cap. + #[test] + fn test_authorized_frame_payload_is_never_trimmed() { + // Build an over-cap authorized frame (big payload, authorization present). + let big = "x".repeat(OBSERVER_MAX_PLAINTEXT_LEN + 1000); + let mut event = event_with_payload( + "acp_read", + serde_json::json!({ "method": "session/request_permission", "body": big }), + ); + event.authorization = Some(crate::observer::AuthorizationEnvelope { + request_nonce: "test-nonce".to_string(), + actionable: true, + reason: None, + }); + + let payload_before = event.payload.clone(); + assert!( + serialized(&event).len() > OBSERVER_MAX_PLAINTEXT_LEN, + "precondition: authorized frame is over-cap" + ); + + fit_observer_event_to_budget(&mut event); + + // Payload must be byte-for-byte identical — no leaf trim, no stub. + assert_eq!( + event.payload, payload_before, + "authorized frame payload must not be mutated by fit_observer_event_to_budget" + ); + // Authorization envelope must still be present and intact. + assert!( + event.authorization.is_some(), + "authorization envelope must survive fit_observer_event_to_budget" + ); + } } diff --git a/desktop/src/features/agents/ui/agentSessionTranscript.ts b/desktop/src/features/agents/ui/agentSessionTranscript.ts index 7dfb7dd8d..624855282 100644 --- a/desktop/src/features/agents/ui/agentSessionTranscript.ts +++ b/desktop/src/features/agents/ui/agentSessionTranscript.ts @@ -50,8 +50,9 @@ export type TranscriptState = { /** * Maps `requestNonce` → `itemId` for actionable permission cards. * Populated alongside `pendingPermissions` when the `authorization` envelope - * is present on the `acp_read` frame. Used by the `permission_decision` - * `control_result` handler to retire the card on any terminal outcome. + * is present on the `acp_read` frame. Used by the nonce-correlated `acp_write` + * terminal handler and the `permission_terminal` event handler to retire the + * card on any terminal outcome (applied, timed_out, cancelled, uncertain). */ pendingPermissionsByNonce: Map; continuationSeq: number; @@ -274,9 +275,84 @@ function describePermissionOutcome( return outcome; } +/** + * Derive human-readable outcome copy from the `authorization.reason` field + * that accompanies terminal `acp_write` events. This is preferred over + * deriving copy from the ACP `result.outcome` field directly because the + * `reason` values are harness-level semantics (applied / timed_out / + * cancelled) whereas `result.outcome` is adapter-level (selected / reject_once + * etc.) and does not distinguish timeout from explicit denial. + * + * Falls back to `describePermissionOutcome` when `reason` is absent (legacy + * paths that predate the authorization envelope). + */ +function describePermissionTerminalReason( + reason: string | undefined, + outcomeKind: string | null | undefined, + optionId: string | null, + options: + | Array<{ optionId: string; kind: string; label?: string }> + | undefined, +): string { + if (reason === "applied") { + // Build optionNames map from the card's options array. + const optionNames = new Map( + (options ?? []).map((o) => [o.optionId, o.kind]), + ); + return describePermissionOutcome( + outcomeKind ?? "selected", + optionId, + optionNames, + ); + } + if (reason === "timed_out") return "Timed out"; + if (reason === "cancelled") return "Cancelled"; + if (reason === "uncertain") { + return "Approval outcome unknown; agent process stopped before it could continue."; + } + // No reason: fall back to ACP outcome-level copy. + const optionNames = new Map((options ?? []).map((o) => [o.optionId, o.kind])); + return describePermissionOutcome(outcomeKind ?? "", optionId, optionNames); +} + +/** + * Retire all live (actionable) permission cards for a given channel. + * Called on terminal turn/process events (`turn_error`, `agent_panic`, + * `turn_completed`) as a backstop so cards do not remain clickable after + * the turn that owned them has ended. + */ +function retireAllLivePermissionCards(d: TranscriptDraft, channelId: string) { + const prefix = `permission:${channelId}:`; + let retired = false; + for (const [id, item] of d.itemsById) { + if ( + id.startsWith(prefix) && + item.type === "lifecycle" && + item.renderClass === "permission" && + item.actionable + ) { + if (!retired) { + // Copy on first mutation. + d.items = [...d.items]; + d.itemsById = new Map(d.itemsById); + retired = true; + d.changed = true; + } + const updated = { ...item, actionable: false }; + d.itemsById.set(id, updated); + const idx = d.items.findIndex((i) => i.id === id); + if (idx !== -1) d.items[idx] = updated; + // Clean up nonce index if present. + if (item.requestNonce) { + d.pendingPermissionsByNonce = new Map(d.pendingPermissionsByNonce); + d.pendingPermissionsByNonce.delete(item.requestNonce); + } + } + } +} + /** * Stable map key for a JSON-RPC id, which may be a string or a finite number - * per the spec. Using JSON.stringify avoids collisions between the number 1 and * the string "1". Returns null for null, undefined, or non-id values (objects, * booleans) so callers can gate on presence without a separate type check. */ @@ -810,6 +886,39 @@ export function processTranscriptEvent( ctx, event.kind, ); + // Backstop: retire any still-live permission cards for this channel so + // missing telemetry and archive replay never reconstruct live controls + // after a terminal turn/process state. + retireAllLivePermissionCards(d, ch); + } else if (event.kind === "turn_completed") { + // Backstop: retire any still-live permission cards for this channel. + // Applied/timed-out/cancelled cards should already be retired via their + // nonce-correlated acp_write frames, but uncertain (process-poison) cards + // may only receive a turn_completed — this ensures they are not left + // actionable in live state or archive replay. + retireAllLivePermissionCards(d, ch); + } else if (event.kind === "permission_terminal") { + // Observer-only terminal event for uncertain outcomes (process poison, + // cancel-during-write). No ACP wire response was confirmed; the harness + // emits this so Desktop can retire the card without a JSON-RPC response. + // Carry the nonce from the authorization envelope. + const auth = event.authorization; + const nonce = auth?.requestNonce; + if (nonce) { + const itemId = d.pendingPermissionsByNonce.get(nonce); + if (itemId) { + const existing = d.itemsById.get(itemId); + if (existing?.type === "lifecycle") { + replaceItem(d, itemId, { + ...existing, + outcome: "Uncertain (process restarting)", + actionable: false, + }); + } + d.pendingPermissionsByNonce = new Map(d.pendingPermissionsByNonce); + d.pendingPermissionsByNonce.delete(nonce); + } + } } else if (event.kind === "acp_read" || event.kind === "acp_write") { const payload = asRecord(event.payload); const method = asString(payload.method); @@ -866,24 +975,70 @@ export function processTranscriptEvent( } } else if (event.kind === "acp_write" && !method) { // Permission response: {"id": , "result": {"outcome": {...}}} + // + // Primary correlation: by `authorization.requestNonce` — a nonce-keyed + // lookup is immune to JSON-RPC id reuse across channels/sessions. + // Legacy fallback: by JSON-RPC id, scoped to channel `ch` so at least + // cross-channel collisions are avoided. + const auth = event.authorization; + const nonce = auth?.requestNonce; const responseId = jsonRpcId(payload.id); const result = asRecord(asRecord(payload.result).outcome); const outcomeKind = asString(result.outcome); - const pending = responseId ? d.pendingPermissions.get(responseId) : null; - if (pending && outcomeKind && responseId) { + + // Derive terminal label from authorization.reason when present; this + // gives "Timed out" for timed_out rather than rendering the ACP + // outcome kind directly (which says "reject_once", not "Timed out"). + const terminalReason = auth?.reason; + + // Resolve the permission card: nonce-keyed wins; fall back to id-keyed. + const itemIdByNonce = nonce + ? d.pendingPermissionsByNonce.get(nonce) + : null; + const pendingById = responseId + ? d.pendingPermissions.get(responseId) + : null; + + if (itemIdByNonce) { + // Nonce-correlated path: resolve the card and derive copy from reason. + const existing = d.itemsById.get(itemIdByNonce); + if (existing?.type === "lifecycle") { + const outcomeText = describePermissionTerminalReason( + terminalReason, + outcomeKind, + asString(result.optionId) ?? null, + existing.options, + ); + replaceItem(d, itemIdByNonce, { + ...existing, + outcome: outcomeText, + actionable: false, + }); + } + // Clean up both indexes. + if (nonce) { + d.pendingPermissionsByNonce = new Map(d.pendingPermissionsByNonce); + d.pendingPermissionsByNonce.delete(nonce); + } + if (responseId) { + d.pendingPermissions = new Map(d.pendingPermissions); + d.pendingPermissions.delete(responseId); + } + } else if (pendingById && outcomeKind && responseId) { + // Legacy id-correlation fallback (non-ask paths with no nonce). const optionId = asString(result.optionId) ?? null; const outcomeText = describePermissionOutcome( outcomeKind, optionId, - pending.optionNames, + pendingById.optionNames, ); - const existing = d.itemsById.get(pending.itemId); + const existing = d.itemsById.get(pendingById.itemId); if (existing?.type === "lifecycle") { - replaceItem(d, pending.itemId, { + replaceItem(d, pendingById.itemId, { ...existing, outcome: outcomeText, + actionable: false, }); - // Remove from pending map — the outcome is now recorded. d.pendingPermissions = new Map(d.pendingPermissions); d.pendingPermissions.delete(responseId); } diff --git a/docs/nips/NIP-AO.md b/docs/nips/NIP-AO.md index 1b29ab59e..ff245f92c 100644 --- a/docs/nips/NIP-AO.md +++ b/docs/nips/NIP-AO.md @@ -120,7 +120,9 @@ below). It is omitted on all other frame kinds. | `acp_read` | Inbound ACP protocol frame (model → harness) | | `acp_write` | Outbound ACP protocol frame (harness → model) | | `turn_started` | A new agent turn has begun | -| `session_resolved` | Session completed or terminated | +| `session_resolved` | Session ready — emitted once when the agent session is established (before the first prompt) | +| `turn_completed` | Terminal lifecycle — emitted when a turn ends (success, cancel, or timeout) | +| `turn_error` | Terminal lifecycle — emitted when a turn ends with an error or process death | | `control_result` | Acknowledgement telemetry emitted after processing a control frame | Permission `acp_read` frames (carrying `session/request_permission` calls) always @@ -165,10 +167,12 @@ call, the `ObserverEvent` carries an `authorization` field: | `"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. + produce an `acp_write` observer event — instead the harness emits a + `permission_terminal` observer event with `authorization.reason = "uncertain"` so + Desktop clients can retire the card without an ACP wire response. 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; the corresponding + `turn_error` and `turn_completed` events are the reliable terminal lifecycle signals. **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