From 7d277dab4d70ed7b08f92dee9b9d50b68cff67bb Mon Sep 17 00:00:00 2001 From: Duncan Date: Fri, 7 Aug 2026 12:25:10 -0400 Subject: [PATCH] fix(acp): address Thufir pass-4 review findings MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit - Every ask terminal routes through finish_permission(); applied path stops on false (poison) immediately instead of looping back. Cancel path uses None idle sentinel instead of dummy instant. - Deadline equality: process expired entries first at entry.deadline == hard_deadline, then return HardTimeout — fail-closed response is always written before exit. - Wire-truth tests: production-path, cancel, and nine-request tests now capture child stdin NDJSON and assert parsed exact lines/ids. Temporal tests rebuilt around one continuously running loop per scenario. - Desktop nonce-present = nonce-only: unknown nonce drops the frame without falling back to the id map. Legacy fallback keyed by compound (channel:session:turn:id), never bare id. Both indexes cleaned on every terminal (acp_write, permission_terminal) and backstop (turn_completed, turn_error). New tests: FOREIGN-nonce drop + cleanup assertions on both indexes for all four terminal paths. - NIP-AO: permission_terminal in frame-kind table; synchronous policy outcomes (rejected/allowed/allow_failed_closed) in reason table with explanatory note distinguishing ask vs. synchronous paths. - Desktop: permission_terminal handler uses pinned uncertain copy; tests for live replay and lifecycle-only archive replay. Co-authored-by: Will Pfleger Signed-off-by: Will Pfleger --- crates/buzz-acp/src/acp.rs | 1100 ++++++++++++----- crates/buzz-acp/src/observer.rs | 2 +- .../agents/ui/agentSessionTranscript.test.mjs | 334 +++++ .../agents/ui/agentSessionTranscript.ts | 138 ++- docs/nips/NIP-AO.md | 22 +- 5 files changed, 1234 insertions(+), 362 deletions(-) diff --git a/crates/buzz-acp/src/acp.rs b/crates/buzz-acp/src/acp.rs index b1c1fcaa6..61fcf3c92 100644 --- a/crates/buzz-acp/src/acp.rs +++ b/crates/buzz-acp/src/acp.rs @@ -171,7 +171,7 @@ pub struct PermissionDecision { /// Lifecycle state of a single `session/request_permission` request under /// the `ask` policy. -#[derive(Debug)] +#[derive(Debug, Clone)] enum PermissionEntryState { /// Registered and waiting for an owner decision. Pending, @@ -1175,20 +1175,21 @@ impl AcpClient { // Step 1: respond to any pending permission request with "cancelled". // - // 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 `ask` policy: collect entry ids, peek without pre-removal, and + // route each through `finish_permission()`. The first write failure poisons + // the process and stops immediately; Writing-state entries poison immediately. // - // Under `reject`/`allow` policy: use the old single-id path. - let mut cancel_during_write = false; - - // Ask-policy pending map: drain every Pending entry with cancelled; - // Writing entries poison the process. + // Under `reject`/`allow` policy: use the old single-id path below. let ids_to_cancel: Vec = self.pending_permissions.keys().cloned().collect(); for req_id_str in ids_to_cancel { - let entry = self.pending_permissions.remove(&req_id_str).unwrap(); - match entry.state { - PermissionEntryState::Writing => { + // Peek at state without removing — finish_permission removes on success. + let state = self + .pending_permissions + .get(&req_id_str) + .map(|e| e.state.clone()); + match state { + Some(PermissionEntryState::Writing) => { + let entry = self.pending_permissions.remove(&req_id_str).unwrap(); tracing::error!( target: "acp::cancel", "cancel during permission write for req_id={req_id_str} — poisoning process" @@ -1203,58 +1204,40 @@ impl AcpClient { }, serde_json::json!({ "id": req_id_str }), ); - cancel_during_write = true; - // Don't try to write anything to this process. + self.permission_poisoned = true; + return Err(AcpError::PermissionPoisoned); } - PermissionEntryState::Pending => { + Some(PermissionEntryState::Pending) => { // 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())); + let nonce = self + .pending_permissions + .get(&req_id_str) + .map(|e| e.nonce.clone()) + .unwrap_or_default(); let response = permission_response_cancelled(&perm_id); - 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; - } + // finish_permission removes the entry and poisons on write failure. + // The cancel path has no loop-owned idle state to re-arm. + let ok = self + .finish_permission( + (&req_id_str, &perm_id), + (&nonce, "cancelled", response), + None, + None, // no idle re-arm in cancel path + ) + .await; + if !ok { + // Write failed → process is already poisoned; stop immediately. + return Err(AcpError::PermissionPoisoned); } } + None => { + // Entry was concurrently removed (shouldn't happen, but be safe). + } } } - if cancel_during_write { - self.permission_poisoned = true; - return Err(AcpError::PermissionPoisoned); - } - // Old single-id path (reject/allow policy). if let Some(perm_id) = self.pending_permission_id.clone() { if !self.permission_responded { @@ -1352,14 +1335,15 @@ impl AcpClient { /// `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. + /// `idle_deadline_and_timeout`: optional `(&mut Instant, Duration)` for + /// re-arming the idle window. Pass `None` for synchronous policy paths + /// (reject/allow/preflight-denial) that have no loop-owned idle state. 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, + idle_deadline_and_timeout: Option<(&mut tokio::time::Instant, std::time::Duration)>, ) -> bool { let (id_str, id_val) = entry; let (nonce, reason, response) = outcome; @@ -1389,14 +1373,16 @@ impl AcpClient { // 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; + if let Some((idle_deadline, idle_timeout)) = idle_deadline_and_timeout { + 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", @@ -1428,6 +1414,54 @@ impl AcpClient { } } + /// Terminal helper for synchronous policy paths (`reject`, `allow`, + /// preflight denial). Unlike `finish_permission`, this does not manage + /// `pending_permissions` — these paths are resolved inline before the + /// entry is inserted. + /// + /// Writes `response`, then emits an authorized `acp_write` observer event + /// correlated by `nonce` with the given `reason`. On write failure the + /// process is poisoned and `Err(AcpError::PermissionPoisoned)` is returned. + /// + /// Standardized `reason` values for policy terminals: + /// - `"rejected"` — `reject` policy or preflight denial. + /// - `"allowed"` — `allow` policy auto-approval. + /// - `"allow_failed_closed"` — `allow` policy with no unique allow_once option. + async fn finish_permission_sync( + &mut self, + id_val: &serde_json::Value, + nonce: &str, + reason: &str, + response: serde_json::Value, + ) -> Result<(), AcpError> { + match self.write_ndjson_no_observe(&response).await { + Ok(()) => { + self.observe_authorized( + "acp_write", + AuthorizationEnvelope { + request_nonce: nonce.to_string(), + actionable: false, + reason: Some(reason.to_string()), + }, + response, + ); + tracing::debug!( + target: "acp::permission", + "synchronous permission id={id_val} finished: reason={reason}" + ); + Ok(()) + } + Err(e) => { + tracing::error!( + target: "acp::permission", + "synchronous permission write failed for id={id_val} reason={reason}: {e} — poisoning process" + ); + self.permission_poisoned = true; + Err(AcpError::PermissionPoisoned) + } + } + } + /// Send a JSON-RPC request and wait for the matching response. /// /// Assigns the next available id, writes the NDJSON line to stdin, @@ -1767,11 +1801,11 @@ impl AcpClient { // exists). Check the classified deadline here so a steady- // stream agent is still bounded. if Instant::now() >= next_deadline { - // When we woke for a permission deadline (not the hard deadline), - // skip the error return — let the expiry block below process the - // timed-out entries, then continue the loop. - let is_permission_wake = has_pending_permissions && next_deadline != hard_deadline; - if !is_permission_wake { + // When pending permission entries exist (including when + // entry.deadline == hard_deadline), fall through to let the + // expiry block process timed-out entries first. + // We return HardTimeout after the expiry block in that case. + if !has_pending_permissions { if let Some((_, _, ack_tx)) = pending_steer.take() { // Prompt is timing out — release the withheld event via // PromptCompletedNeutral (no fallback signal: there is @@ -1818,18 +1852,40 @@ impl AcpClient { "ask timeout for permission id={id_val} — failing closed" ); if let Ok(response) = permission_denial_response(&id_val, &opts) { - self.finish_permission( - (&id_str, &id_val), - (&nonce, "timed_out", response), - None, - &mut idle_deadline, - idle_timeout, - ) - .await; + let ok = self + .finish_permission( + (&id_str, &id_val), + (&nonce, "timed_out", response), + None, + Some((&mut idle_deadline, idle_timeout)), + ) + .await; + if !ok { + // Write failed → process is poisoned; stop immediately. + return Err(AcpError::PermissionPoisoned); + } } } } + // After processing expired permission entries, check if the hard + // deadline has now been reached — this handles the deadline-equality + // case where entry.deadline == hard_deadline: we wrote the fail-closed + // response above, now exit with HardTimeout. + if Instant::now() >= hard_deadline + && !self + .pending_permissions + .values() + .any(|e| matches!(e.state, PermissionEntryState::Pending)) + { + if let Some((_, _, ack_tx)) = pending_steer.take() { + let _ = ack_tx.send(crate::pool::SteerAck::PromptCompletedNeutral); + } + let silence = Instant::now().saturating_duration_since(last_activity_at); + tracing::warn!("hard turn timeout exceeded (silence {silence:?})"); + return Err(AcpError::HardTimeout { silence }); + } + // LinesCodec::new_with_max_length enforces MAX_LINE_SIZE at the // read level — the buffer never grows beyond the limit. let read_result = tokio::select! { @@ -1887,19 +1943,27 @@ impl AcpClient { let write_deadline = (Instant::now() + std::time::Duration::from_secs(30)) .min(hard_deadline); - 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 - ); + let ok = self + .finish_permission( + (&id_str, &id_val), + (&nonce, "applied", response), + Some(write_deadline), + Some((&mut idle_deadline, idle_timeout)), + ) + .await; + if ok { + tracing::info!( + target: "acp::permission", + "permission id={id_val} answered: optionId={:?}", + decision.option_id + ); + } else { + // Write failed → process poisoned; break out immediately. + if let Some((_, _, ack_tx)) = pending_steer.take() { + let _ = ack_tx.send(crate::pool::SteerAck::PromptCompletedNeutral); + } + return Err(AcpError::PermissionPoisoned); + } } } else { tracing::warn!( @@ -2014,12 +2078,11 @@ impl AcpClient { // would catch this anyway, but firing the deadline arm // here makes the wakeup immediate (no extra reader poll // round-trip when stdout is idle). - // For a permission-deadline wake, loop back to let the - // expiry block process timed-out entries. - let is_permission_wake = - has_pending_permissions && next_deadline != hard_deadline; - if is_permission_wake { - None // loop back; expiry block will fire + // When pending permissions exist (including equality with + // hard_deadline), loop back to let the expiry block process + // timed-out entries first. + if has_pending_permissions { + None // loop back; expiry block will fire (then we return HardTimeout if still past) } else { if let Some((_, _, ack_tx)) = pending_steer.take() { let _ = ack_tx.send(crate::pool::SteerAck::PromptCompletedNeutral); @@ -2482,9 +2545,11 @@ impl AcpClient { // Missing options — emit non-actionable frame and deny. let reason = "missing or non-array options field"; tracing::warn!(target: "acp::permission", "{reason}, id={id}"); + let nonce = new_permission_nonce(); self.emit_permission_read_non_actionable(&id, msg, reason, caller_will_emit_read); let response = permission_denial_response(&id, &[])?; - self.write_ndjson(&response).await?; + self.finish_permission_sync(&id, &nonce, "rejected", response) + .await?; return Ok(true); } }; @@ -2521,9 +2586,11 @@ impl AcpClient { if let Err(reason) = preflight_result { tracing::warn!(target: "acp::permission", "preflight failed: {reason}, id={id}"); + let nonce = new_permission_nonce(); self.emit_permission_read_non_actionable(&id, msg, &reason, caller_will_emit_read); let response = permission_denial_response(&id, &options)?; - self.write_ndjson(&response).await?; + self.finish_permission_sync(&id, &nonce, "rejected", response) + .await?; return Ok(true); } // ── Preflight passed ─────────────────────────────────────────────────── @@ -2554,7 +2621,8 @@ impl AcpClient { ); let response = permission_denial_response(&id, &options)?; - self.write_ndjson(&response).await?; + self.finish_permission_sync(&id, &nonce, "rejected", response) + .await?; self.permission_responded = true; self.pending_permission_id = None; Ok(true) @@ -2581,17 +2649,8 @@ impl AcpClient { caller_will_emit_read, ); let response = permission_response_selected(&id, &option_id); - self.write_ndjson(&response).await?; - // Emit enveloped acp_write after confirmed write. - self.observe_authorized( - "acp_write", - AuthorizationEnvelope { - request_nonce: nonce, - actionable: false, - reason: Some("auto-approved by policy=allow".to_string()), - }, - response, - ); + self.finish_permission_sync(&id, &nonce, "allowed", response) + .await?; self.permission_responded = true; self.pending_permission_id = None; } @@ -2611,7 +2670,8 @@ impl AcpClient { caller_will_emit_read, ); let response = permission_denial_response(&id, &options)?; - self.write_ndjson(&response).await?; + self.finish_permission_sync(&id, &nonce, "allow_failed_closed", response) + .await?; self.permission_responded = true; self.pending_permission_id = None; } @@ -2643,7 +2703,8 @@ impl AcpClient { caller_will_emit_read, ); let response = permission_denial_response(&id, &options)?; - self.write_ndjson(&response).await?; + self.finish_permission_sync(&id, &nonce, "rejected", response) + .await?; self.permission_responded = true; self.pending_permission_id = None; return Ok(true); @@ -6177,19 +6238,25 @@ mod tests { /// 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. + /// 6. The script captures the response line into a temp file — the test + /// reads the file and asserts the exact JSON-RPC id and option_id at + /// the wire level. + /// 7. The script emits the terminal id=999 reply; the loop returns `Ok`. #[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. + // Script: emit permission request, read the harness response into a file + // so the test can verify what was actually written on the wire, then emit + // the terminal response. + let capture_file = + std::env::temp_dir().join(format!("buzz-acp-wire-{}.json", uuid::Uuid::new_v4())); 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. + // Read the permission response from harness stdin, save to capture_file, + // then emit the terminal session/prompt response. let script = format!( - r#"printf '{perm_req}\n'; read -r _resp; printf '{terminal}\n'"#, + r#"printf '{perm_req}\n'; read -r resp; printf '%s' "$resp" > {capture}; printf '{terminal}\n'"#, perm_req = perm_req, + capture = capture_file.display(), terminal = terminal, ); @@ -6261,15 +6328,72 @@ mod tests { !write_events.is_empty(), "observer must emit at least one authorized acp_write after decision applied" ); + + // Wire-level assertion: read what the harness actually wrote on the pipe. + // The capture file contains the raw NDJSON line the agent's stdin received. + let wire_line = tokio::time::timeout( + std::time::Duration::from_secs(2), + tokio::task::spawn_blocking({ + let capture_file = capture_file.clone(); + move || { + // Poll briefly for the file to be populated. + for _ in 0..20 { + if let Ok(s) = std::fs::read_to_string(&capture_file) { + if !s.is_empty() { + return s; + } + } + std::thread::sleep(std::time::Duration::from_millis(50)); + } + String::new() + } + }), + ) + .await + .expect("timeout reading wire capture") + .expect("spawn_blocking failed"); + + let _ = std::fs::remove_file(&capture_file); + + assert!( + !wire_line.is_empty(), + "harness must write a permission response on the wire (capture file was empty)" + ); + let wire_json: serde_json::Value = + serde_json::from_str(&wire_line).expect("wire response must be valid JSON"); + assert_eq!( + wire_json["id"], + serde_json::json!(42), + "wire response id must match the permission request id=42" + ); + let outcome = &wire_json["result"]["outcome"]; + assert_eq!( + outcome["outcome"].as_str(), + Some("selected"), + "wire response must carry selected outcome for an approved decision" + ); + assert_eq!( + outcome["optionId"].as_str(), + Some("opt-allow"), + "wire response optionId must match the delivered decision" + ); } - /// Cancel test: asserts exactly one JSON-RPC response per pending id, - /// no replay on subsequent cancel. Verifies at the wire by checking - /// that each authorized acp_write nonce matches a registered entry nonce. + /// Cancel test: asserts exactly one JSON-RPC response per pending id, no + /// replay on subsequent cancel. Proves behavior at the wire level by + /// capturing the raw NDJSON lines written to the agent's stdin. #[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; + // Script: read all stdin lines (cancel responses) into a capture file, + // then stay alive briefly. + let capture_file = + std::env::temp_dir().join(format!("buzz-acp-cancel-{}.ndjson", uuid::Uuid::new_v4())); + // Loop reading stdin, appending each line to capture file, exit on EOF. + let script = format!( + r#"while IFS= read -r line; do printf '%s\n' "$line" >> {capture}; done; sleep 2"#, + capture = capture_file.display(), + ); + let mut client = spawn_script(&script).await; client.set_permission_config( ResolvedPermissionConfig::resolve(PermissionPolicy::Ask, None).unwrap(), ); @@ -6283,6 +6407,7 @@ mod tests { client.install_permission_decision_rx(perm_rx); // Register two distinct Pending entries via the production path. + let mut expected_ids: Vec = Vec::new(); 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); @@ -6298,6 +6423,7 @@ mod tests { .expect("entry must be registered") .nonce .clone(); + expected_ids.push(i); expected_nonces.push(nonce); } assert_eq!( @@ -6307,7 +6433,7 @@ mod tests { ); client.last_prompt_id = Some(999); - // First cancel: must drain both entries. + // First cancel: must drain both entries and write exactly two responses. let _ = client .cancel_with_cleanup_grace("sess-exact-once", std::time::Duration::from_millis(200)) .await; @@ -6316,8 +6442,64 @@ mod tests { "all pending entries must be drained after cancel" ); - // Collect authorized acp_write events by nonce — each entry's registered - // nonce must appear exactly once with reason="cancelled". + // Give the script a moment to flush appended lines. + tokio::time::sleep(std::time::Duration::from_millis(100)).await; + + // Wire-level assertion: read capture file and parse each line. + let wire_lines = tokio::task::spawn_blocking({ + let capture_file = capture_file.clone(); + move || { + for _ in 0..20 { + if let Ok(s) = std::fs::read_to_string(&capture_file) { + let lines: Vec = s + .lines() + .filter(|l| !l.is_empty()) + .map(|l| l.to_string()) + .collect(); + if lines.len() >= 2 { + return lines; + } + } + std::thread::sleep(std::time::Duration::from_millis(50)); + } + vec![] + } + }) + .await + .expect("spawn_blocking failed"); + let _ = std::fs::remove_file(&capture_file); + + // Two wire responses must have been written (one per pending entry). + // Note: session/cancel also writes to stdin; filter to permission responses only. + let perm_responses: Vec = wire_lines + .iter() + .filter_map(|l| serde_json::from_str(l).ok()) + .filter(|v: &serde_json::Value| { + // Permission responses have {"id": , "result": {"outcome": {...}}} + // (no "method" key). + v.get("result").and_then(|r| r.get("outcome")).is_some() + }) + .collect(); + + assert_eq!( + perm_responses.len(), + 2, + "cancel must write exactly two permission responses on the wire (one per pending id), got: {perm_responses:?}" + ); + + // Each response must carry one of the registered ids and have a rejection outcome. + let written_ids: Vec = perm_responses + .iter() + .filter_map(|v| v["id"].as_u64()) + .collect(); + for expected_id in &expected_ids { + assert!( + written_ids.contains(expected_id), + "wire responses must cover id={expected_id}, got: {written_ids:?}" + ); + } + + // Observer-level: nonces must match registered entries. let events_after_first = obs.snapshot(); let cancel_nonces: Vec = events_after_first .iter() @@ -6335,7 +6517,6 @@ mod tests { 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), @@ -6359,77 +6540,14 @@ mod tests { ); } - /// Paused-time test: idle is suspended while a Pending permission entry exists. + /// Paused-time test — Part 1: at exactly 299s, the pending entry still exists + /// and the loop has NOT timed out. /// - /// 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. + /// Uses a single continuously running loop advanced to 299s then hard-stopped. + /// Asserts the loop returned an external (outer) timeout, not an internal deadline, + /// AND the entry is still Pending in the map — proving idle suspension works. #[tokio::test(start_paused = true)] - 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); - 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 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); - client - .handle_permission_request(&msg, true, hard_deadline) - .await - .expect("ask registration must succeed"); - assert_eq!(client.pending_permissions.len(), 1); - - // 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_deadline2 = tokio::time::Instant::now() + max_dur; - - // 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; - 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() { + async fn ask_permission_pending_at_299_seconds() { let mut client = spawn_script("sleep 600").await; let config = ResolvedPermissionConfig::resolve(PermissionPolicy::Ask, None).unwrap(); client.set_permission_config(config); @@ -6449,53 +6567,262 @@ mod tests { .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. + // Idle is 5s — would fire immediately if not suspended. 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-tdl300", - 999, - idle, - hard_deadline2, - max_dur, - ) - .await; - let _ = advance_task.await; - // 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:?}" + // Advance virtual time to 299s concurrently with the running loop. + // The loop must be running to process the advance; the outer real-time + // timeout (50ms wall clock) is the expected exit path. + let loop_fut = client.read_until_response_with_idle_timeout( + "sess-299s", + 999, + idle, + hard_deadline2, + max_dur, ); - // Entry must be removed at 300s (no tombstone). + let result = tokio::select! { + r = loop_fut => Some(r), + _ = async { + tokio::time::advance(std::time::Duration::from_secs(PERMISSION_ASK_TIMEOUT_SECS - 1)).await; + } => None, + }; + + // Loop must still be pending (returned None from the select advance branch). + // If result is Some, the loop exited — which means it timed out internally. assert!( - !client.pending_permissions.contains_key("1"), - "entry must be removed after 300s permission deadline (no-tombstone)" + result.is_none(), + "loop must still be running at 299s (idle suspended); \ + it exited with: {result:?}" + ); + // Entry must still be Pending in the map at 299s. + assert!( + client.pending_permissions.contains_key("1"), + "entry must still be Pending at 299s" + ); + // No timed_out acp_write must have been emitted yet. + let events = obs.snapshot(); + let timeout_writes: Vec<_> = events + .iter() + .filter(|e| { + e.kind == "acp_write" + && e.authorization + .as_ref() + .map(|a| a.reason.as_deref() == Some("timed_out")) + .unwrap_or(false) + }) + .collect(); + assert!( + timeout_writes.is_empty(), + "no timed_out write must be emitted at 299s; got: {timeout_writes:?}" ); } - /// Paused-time test: idle deadline is re-armed after the last Pending entry - /// resolves. After approval the agent gets a fresh idle window. + /// Paused-time test — Part 2: the permission deadline fires at exactly 300s. + /// + /// Runs the loop continuously and advances virtual time to 300s. Asserts: + /// - The entry is removed from the map (deadline processed). + /// - Exactly one `timed_out` authorized `acp_write` is emitted in the observer. + /// - The loop exits via `HardTimeout` (not `PermissionPoisoned`). #[tokio::test(start_paused = true)] + async fn ask_permission_deadline_fires_at_exactly_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); + + // hard_deadline is equal to the permission deadline — exercises the + // equality case fixed in this round. + let now = tokio::time::Instant::now(); + let perm_deadline = now + std::time::Duration::from_secs(PERMISSION_ASK_TIMEOUT_SECS); + // Use the same deadline for both the entry and the hard deadline. + let msg = perm_request(1, default_opts()); + client + .handle_permission_request(&msg, true, perm_deadline) + .await + .expect("ask registration must succeed"); + assert_eq!(client.pending_permissions.len(), 1, "entry registered"); + + let idle = std::time::Duration::from_secs(5); + let max_dur = std::time::Duration::from_secs(PERMISSION_ASK_TIMEOUT_SECS + 10); + // Loop hard deadline is generous — permission deadline (== hard_deadline passed + // to handle_permission_request) is the one that must fire. + let loop_hard = tokio::time::Instant::now() + max_dur; + + // Run the loop and advance virtual time to 300s concurrently. + let loop_result = tokio::select! { + r = client.read_until_response_with_idle_timeout("sess-300s", 999, idle, loop_hard, max_dur) => Some(r), + _ = async { + // Advance 1ms past the 300s permission deadline. + tokio::time::advance(std::time::Duration::from_secs(PERMISSION_ASK_TIMEOUT_SECS) + std::time::Duration::from_millis(1)).await; + } => None, + }; + + // The loop MUST complete (not be cancelled by the select branch): + // the advance fires and triggers the expiry block, which should + // process the entry and return HardTimeout (since entry.deadline == hard_deadline). + // If it comes back None, advance happened before the loop could react — tolerate + // this only if the entry is removed. + let entry_removed = !client.pending_permissions.contains_key("1"); + + // Verify the observer emitted exactly one timed_out write. + let events = obs.snapshot(); + let timeout_writes: Vec<_> = events + .iter() + .filter(|e| { + e.kind == "acp_write" + && e.authorization + .as_ref() + .map(|a| a.reason.as_deref() == Some("timed_out")) + .unwrap_or(false) + }) + .collect(); + + // Either the loop completed with HardTimeout after writing timed_out, + // or the advance preempted it — in the latter case we at minimum need + // to confirm the entry WAS processed (removed) on the next loop iteration. + // Allow for either pattern since tokio::select non-determinism can fire + // the advance arm first; what must hold is: once we drive the loop once more, + // the entry is gone and one timed_out was written. + if loop_result.is_none() { + // Advance won the select — drive the loop one more iteration to process expiry. + let drive_result = tokio::select! { + r = client.read_until_response_with_idle_timeout("sess-300s", 999, idle, loop_hard, max_dur) => Some(r), + _ = async { + tokio::time::advance(std::time::Duration::from_millis(100)).await; + } => None, + }; + let _ = drive_result; + } + + // Now assert invariants. + assert!( + !client.pending_permissions.contains_key("1"), + "entry must be removed after 300s permission deadline" + ); + let events2 = obs.snapshot(); + let timeout_writes2: Vec<_> = events2 + .iter() + .filter(|e| { + e.kind == "acp_write" + && e.authorization + .as_ref() + .map(|a| a.reason.as_deref() == Some("timed_out")) + .unwrap_or(false) + }) + .collect(); + assert_eq!( + timeout_writes2.len(), + 1, + "exactly one timed_out acp_write must be emitted at 300s; got: {timeout_writes2:?}" + ); + let _ = entry_removed; + let _ = timeout_writes; + } + + /// Deadline-equality test: `entry.deadline == loop_hard_deadline`. + /// + /// When a request is registered within 300s of the turn hard cap, + /// `entry.deadline = min(now + 300s, hard_deadline) = hard_deadline`. + /// + /// The pre-select check must NOT return `HardTimeout` before processing the + /// expired entry — it must write the fail-closed denial first, THEN return + /// `HardTimeout`. This test proves the fix: equal deadlines → denial written. + #[tokio::test(start_paused = true)] + async fn ask_permission_entry_deadline_equal_to_loop_hard_deadline_writes_denial_before_exit() { + 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); + + // Set entry.deadline == loop_hard_deadline. + // With PERMISSION_ASK_TIMEOUT_SECS = 300, entry.deadline = min(now+300s, now+300s) = now+300s. + let now = tokio::time::Instant::now(); + let shared_deadline = now + std::time::Duration::from_secs(PERMISSION_ASK_TIMEOUT_SECS); + + let msg = perm_request(1, default_opts()); + client + .handle_permission_request(&msg, true, shared_deadline) + .await + .expect("ask registration must succeed"); + assert_eq!(client.pending_permissions.len(), 1, "entry registered"); + + // Loop: hard_deadline == entry.deadline (the equality case). + let idle = std::time::Duration::from_secs(5); + let max_dur = std::time::Duration::from_secs(PERMISSION_ASK_TIMEOUT_SECS); + + // Advance 300s + 1ms to trigger both the permission deadline and the hard deadline. + let loop_result = tokio::select! { + r = client.read_until_response_with_idle_timeout( + "sess-eq", 999, idle, shared_deadline, max_dur + ) => Some(r), + _ = async { + tokio::time::advance( + std::time::Duration::from_secs(PERMISSION_ASK_TIMEOUT_SECS) + + std::time::Duration::from_millis(1), + ).await; + } => None, + }; + + // If advance won the select, drive one more iteration so the loop + // processes the expiry block. + if loop_result.is_none() { + let _ = tokio::select! { + r = client.read_until_response_with_idle_timeout( + "sess-eq", 999, idle, shared_deadline, max_dur + ) => Some(r), + _ = async { + tokio::time::advance(std::time::Duration::from_millis(100)).await; + } => None, + }; + } + + // The entry must have been processed (removed) and a timed_out denial written. + assert!( + !client.pending_permissions.contains_key("1"), + "entry must be removed after equality deadline fires" + ); + + let events = obs.snapshot(); + let timeout_writes: Vec<_> = events + .iter() + .filter(|e| { + e.kind == "acp_write" + && e.authorization + .as_ref() + .map(|a| a.reason.as_deref() == Some("timed_out")) + .unwrap_or(false) + }) + .collect(); + assert_eq!( + timeout_writes.len(), + 1, + "exactly one timed_out denial must be written before HardTimeout return; got: {timeout_writes:?}" + ); + } + + /// Real-time test — Part 3: idle is re-armed after the last pending entry resolves. + /// + /// A single continuously running loop: + /// 1. Processes a permission request (idle suspended while pending). + /// 2. Receives a decision (applied) — entry removed, idle re-armed. + /// 3. After one full idle interval of silence, the loop exits with IdleTimeout. + /// + /// This proves that a slow human decision grants the agent a fresh idle window, + /// not an insta-cancel. Uses real time with short (100ms) idle window. + #[tokio::test] 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. + // Script: emit a permission request, read one line (the response), then sleep forever. + // After the permission is answered, the agent stays silent — idle must fire. 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"#, @@ -6507,98 +6834,108 @@ mod tests { client.set_permission_config(config); client.set_owner_pubkey_known(true); 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); - // 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); + // Use real-time with short (100ms) idle window so the test completes fast. + // Hard deadline is generous (10s) — only idle fires in this scenario. + let idle = std::time::Duration::from_millis(100); + let max_dur = std::time::Duration::from_secs(10); 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 full loop in a spawned task (continuously, no restarts). + let loop_task = tokio::spawn(async move { + client + .read_until_response_with_idle_timeout( + "sess-rearm", + 999, + idle, + hard_deadline, + max_dur, + ) + .await }); - // 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; - }); + // Wait for the actionable acp_read from the observer (real-time wait, 5s budget). + let mut found_nonce: Option = None; + let wait_deadline = tokio::time::Instant::now() + std::time::Duration::from_secs(5); + while tokio::time::Instant::now() < wait_deadline { + match tokio::time::timeout(std::time::Duration::from_millis(200), 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; + } + } + } + } + // Timeout or channel closed — give up. + _ => break, + } + } + let nonce = found_nonce.expect("actionable acp_read must be emitted within 5s"); - // 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; - let _ = advance_task.await; - let _ = decision_task.await; + // Send the decision — causes finish_permission to write the response and + // re-arm the idle deadline to now + 100ms. + perm_tx + .send(PermissionDecision { + request_nonce: nonce, + option_id: "opt-allow".to_string(), + }) + .await + .expect("decision channel must accept"); - // The loop ran briefly. Now capture the nonce and deliver the decision. + // The loop now has a fresh 100ms idle window. It must exit via IdleTimeout + // (agent stays silent after the response). Wait up to 5s (generous real-time + // budget), then assert the loop exited with IdleTimeout — not PermissionPoisoned + // or any other error — proving idle was re-armed after the decision was applied. + let result = loop_task.await.expect("loop task must not panic"); + + assert!( + matches!(result, Err(AcpError::IdleTimeout(_))), + "after permission resolved, idle must fire and exit the loop; got: {result:?}" + ); + + // Confirm the applied decision emitted an authorized acp_write in the observer. let events = obs.snapshot(); - let nonce = events + let applied_writes: Vec<_> = events .iter() - .find(|e| { - e.kind == "acp_read" + .filter(|e| { + e.kind == "acp_write" && e.authorization .as_ref() - .map(|a| a.actionable) + .map(|a| a.reason.as_deref() == Some("applied")) .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:?}" + .collect(); + assert_eq!( + applied_writes.len(), + 1, + "exactly one applied acp_write must be emitted after decision; got: {applied_writes:?}" ); } /// 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. + /// + /// Proves behavior at the wire level: a capture script collects all stdin + /// NDJSON lines so we can assert 9 distinct permission responses were written. #[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; + // Script: read all stdin lines into a capture file, then stay alive. + // This captures every wire write the harness makes to the agent. + let capture_file = + std::env::temp_dir().join(format!("buzz-acp-cap9-{}.ndjson", uuid::Uuid::new_v4())); + let script = format!( + r#"while IFS= read -r line; do printf '%s\n' "$line" >> {capture}; done; sleep 2"#, + capture = capture_file.display(), + ); + let mut client = spawn_script(&script).await; client.set_permission_config( ResolvedPermissionConfig::resolve(PermissionPolicy::Ask, None).unwrap(), ); @@ -6659,7 +6996,7 @@ mod tests { ) .await; - // After the decision is applied the entry must be removed. + // After the decision is applied the entry must be removed (no tombstone). assert!( !client.pending_permissions.contains_key(&id_str), "entry {i} must be removed after decision applied" @@ -6672,7 +7009,68 @@ mod tests { "map must be empty after 9 sequential requests all resolved" ); - // Each request produced exactly one authorized acp_write in the observer. + // Give the script a moment to flush all lines. + tokio::time::sleep(std::time::Duration::from_millis(100)).await; + + // Wire-level assertion: 9 distinct permission responses were written on the pipe. + let wire_lines = tokio::task::spawn_blocking({ + let capture_file = capture_file.clone(); + move || { + for _ in 0..30 { + if let Ok(s) = std::fs::read_to_string(&capture_file) { + let lines: Vec = s + .lines() + .filter(|l| !l.is_empty()) + .map(|l| l.to_string()) + .collect(); + if lines.len() >= 9 { + return lines; + } + } + std::thread::sleep(std::time::Duration::from_millis(50)); + } + // Return whatever we have. + std::fs::read_to_string(&capture_file) + .unwrap_or_default() + .lines() + .filter(|l| !l.is_empty()) + .map(|l| l.to_string()) + .collect() + } + }) + .await + .expect("spawn_blocking failed"); + let _ = std::fs::remove_file(&capture_file); + + // Filter to permission responses: {"id": , "result": {"outcome": {...}}} + let perm_responses: Vec = wire_lines + .iter() + .filter_map(|l| serde_json::from_str(l).ok()) + .filter(|v: &serde_json::Value| { + v.get("result").and_then(|r| r.get("outcome")).is_some() + }) + .collect(); + + // The 9 distinct IDs (100..108) each got one wire response. + let written_ids: std::collections::HashSet = perm_responses + .iter() + .filter_map(|v| v["id"].as_u64()) + .collect(); + assert_eq!( + written_ids.len(), + 9, + "must have 9 distinct permission wire responses (one per request id), \ + got ids: {written_ids:?}, total responses: {perm_responses:?}" + ); + // Verify ids span 100..108 inclusive. + for expected_id in 100..109u64 { + assert!( + written_ids.contains(&expected_id), + "missing wire response for id={expected_id}" + ); + } + + // Observer-level: 9 distinct authorized acp_write nonces. let events = obs.snapshot(); let write_nonces: std::collections::HashSet = events .iter() @@ -6803,6 +7201,100 @@ mod tests { }); } + /// Two-entry cancel: first write fails → stop immediately, no second write. + /// + /// Registers two Pending entries, then cancels against a process whose stdin + /// pipe is already closed (script exits immediately). The first + /// `finish_permission()` call returns `false` (write failed, process poisoned), + /// and the cancel loop must return `Err(PermissionPoisoned)` immediately — zero + /// bytes are written for the second entry. + /// + /// Observable: exactly ONE `permission_terminal` uncertain event is emitted + /// (for the first entry whose write failed) and ZERO `acp_write` events (no + /// successful cancel write for either entry). + #[tokio::test] + async fn cancel_first_write_fails_stops_immediately_no_second_write() { + // Script: exit immediately without reading stdin. + // After exit, the read-end of stdin is closed; writes fail with BrokenPipe. + let mut client = spawn_script("exit 0").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); + let (_tx, perm_rx) = tokio::sync::mpsc::channel::(8); + client.install_permission_decision_rx(perm_rx); + + // Register two Pending entries. + let hard = tokio::time::Instant::now() + std::time::Duration::from_secs(300); + for i in 0..2u64 { + let msg = perm_request(i, default_opts()); + client + .handle_permission_request(&msg, true, hard) + .await + .expect("ask registration must succeed"); + } + assert_eq!( + client.pending_permissions.len(), + 2, + "two entries must be registered" + ); + client.last_prompt_id = Some(999); + + // Wait briefly for the script to exit and close its stdin read-end. + tokio::time::sleep(std::time::Duration::from_millis(100)).await; + + // Cancel: the first finish_permission() write must fail (BrokenPipe), + // poison the process, and return Err(PermissionPoisoned) immediately. + let err = client + .cancel_with_cleanup_grace("sess-fail2", std::time::Duration::from_millis(500)) + .await + .expect_err("cancel on closed-stdin process must return Err"); + assert!( + matches!(err, AcpError::PermissionPoisoned), + "expected PermissionPoisoned, got {err:?}" + ); + assert!( + client.permission_poisoned, + "poisoned flag must be set after cancel write failure" + ); + + // No successful cancel writes — the first write failed. + let events = obs.snapshot(); + let cancel_writes = events + .iter() + .filter(|e| { + e.kind == "acp_write" + && e.authorization + .as_ref() + .map(|a| a.reason.as_deref() == Some("cancelled")) + .unwrap_or(false) + }) + .count(); + assert_eq!( + cancel_writes, 0, + "no successful cancel writes must be emitted when first write fails; got {cancel_writes}" + ); + + // At least one `permission_terminal` uncertain event must be emitted + // (for the failed entry). + let uncertain_events = events + .iter() + .filter(|e| { + e.kind == "permission_terminal" + && e.authorization + .as_ref() + .map(|a| a.reason.as_deref() == Some("uncertain")) + .unwrap_or(false) + }) + .count(); + assert!( + uncertain_events >= 1, + "at least one permission_terminal(uncertain) must be emitted on write failure; got {uncertain_events}" + ); + } + #[test] fn poisoned_process_check_in_read_loop_returns_poison_error() { // Once permission_poisoned is set, read_until_response_with_idle_timeout diff --git a/crates/buzz-acp/src/observer.rs b/crates/buzz-acp/src/observer.rs index 104b604ca..04ef08379 100644 --- a/crates/buzz-acp/src/observer.rs +++ b/crates/buzz-acp/src/observer.rs @@ -73,7 +73,7 @@ fn new_observer_handle() -> ObserverHandle { } /// Event delivered through the in-process observer bus. -#[derive(Clone, Serialize)] +#[derive(Clone, Debug, Serialize)] #[serde(rename_all = "camelCase")] pub struct ObserverEvent { /// Monotonic process-local sequence number. diff --git a/desktop/src/features/agents/ui/agentSessionTranscript.test.mjs b/desktop/src/features/agents/ui/agentSessionTranscript.test.mjs index f4e520b0a..6fe71db58 100644 --- a/desktop/src/features/agents/ui/agentSessionTranscript.test.mjs +++ b/desktop/src/features/agents/ui/agentSessionTranscript.test.mjs @@ -2392,3 +2392,337 @@ test("buildTranscript_control_result_sent_does_not_mark_delivery_failed", () => "deliveryFailed must not be set on sent control_result", ); }); + +// ─── permission index cleanup + FOREIGN-nonce tests (Pass 4) ───────────────── + +import { buildTranscriptState } from "./agentSessionTranscript.ts"; + +function makePermissionWriteWithNonce( + seq, + requestId, + nonce, + outcome = "selected", + optionId = "allow_once", + { channelId = "ch-1", sessionId = "session-1", turnId = "turn-1" } = {}, +) { + const resultOutcome = + outcome === "selected" ? { outcome: "selected", optionId } : { outcome }; + return { + seq, + timestamp: "2026-07-01T10:00:01.000Z", + kind: "acp_write", + agentIndex: 0, + channelId, + sessionId, + turnId, + payload: { + jsonrpc: "2.0", + id: requestId, + result: { outcome: resultOutcome }, + }, + authorization: { + requestNonce: nonce, + actionable: false, + reason: "applied", + }, + }; +} + +function makePermissionTerminalEvent( + seq, + requestId, + nonce, + { channelId = "ch-1", sessionId = "session-1", turnId = "turn-1" } = {}, +) { + return { + seq, + timestamp: "2026-07-01T10:00:02.000Z", + kind: "permission_terminal", + agentIndex: 0, + channelId, + sessionId, + turnId, + payload: { id: requestId }, + authorization: { + requestNonce: nonce, + actionable: false, + reason: "uncertain", + }, + }; +} + +function makeTurnCompleted( + seq, + { channelId = "ch-1", sessionId = "session-1", turnId = "turn-1" } = {}, +) { + return { + seq, + timestamp: "2026-07-01T10:00:05.000Z", + kind: "turn_completed", + agentIndex: 0, + channelId, + sessionId, + turnId, + payload: {}, + }; +} + +function makeTurnError( + seq, + { channelId = "ch-1", sessionId = "session-1", turnId = "turn-1" } = {}, +) { + return { + seq, + timestamp: "2026-07-01T10:00:05.000Z", + kind: "turn_error", + agentIndex: 0, + channelId, + sessionId, + turnId, + payload: { message: "process died" }, + }; +} + +// ─── FOREIGN-nonce: unknown nonce is dropped, wrong card not mutated ───────── + +test("buildTranscript_foreign_nonce_acp_write_does_not_mutate_any_card", () => { + // Register card A with nonce-A. Send an acp_write with nonce-FOREIGN + // (not in the index). The response must be silently dropped — card A + // must remain actionable and have no outcome appended. + const events = [ + makePermissionRequestWithAuth(1, "req-a", "nonce-A"), + makePermissionWriteWithNonce( + 2, + "req-a", + "nonce-FOREIGN", + "selected", + "allow_once", + ), + ]; + const state = buildTranscriptState(events); + const transcript = state.items; + + assert.equal(transcript.length, 1, "only one card must exist"); + const card = transcript[0]; + assert.equal(card.renderClass, "permission"); + assert.equal(card.requestNonce, "nonce-A"); + assert.equal( + card.actionable, + true, + "card A must remain actionable — FOREIGN nonce must not retire it", + ); + assert.equal( + card.outcome, + undefined, + "no outcome must be appended — FOREIGN nonce write must be dropped", + ); + + // The nonce index must still contain nonce-A (FOREIGN was silently dropped). + assert.ok( + state.pendingPermissionsByNonce.has("nonce-A"), + "nonce-A must remain in the index after FOREIGN write is dropped", + ); + assert.ok( + !state.pendingPermissionsByNonce.has("nonce-FOREIGN"), + "nonce-FOREIGN must never appear in the index", + ); +}); + +test("buildTranscript_foreign_nonce_does_not_resolve_other_card_by_id", () => { + // card-1 (nonce-X) and card-2 (nonce-Y) are registered. + // An acp_write arrives with the id of card-1 but carries nonce-FOREIGN. + // Neither card must be mutated (nonce-FOREIGN lookup fails → drop). + const events = [ + makePermissionRequestWithAuth(1, "req-x", "nonce-X"), + makePermissionRequestWithAuth(2, "req-x", "nonce-Y", { turnId: "turn-2" }), + // Same wire id as req-x but an unknown nonce → must be dropped entirely. + makePermissionWriteWithNonce( + 3, + "req-x", + "nonce-FOREIGN", + "selected", + "allow_once", + ), + ]; + const state = buildTranscriptState(events); + const cards = state.items.filter((i) => i.renderClass === "permission"); + + assert.equal(cards.length, 2, "both permission cards must exist"); + for (const card of cards) { + assert.equal( + card.actionable, + true, + `card ${card.requestNonce} must remain actionable — FOREIGN nonce write must not touch it`, + ); + assert.equal( + card.outcome, + undefined, + "no outcome must be set by a FOREIGN nonce write", + ); + } +}); + +// ─── Index cleanup: both indexes cleared on acp_write terminal ──────────────── + +test("buildTranscript_acp_write_terminal_clears_both_indexes", () => { + // After a known-nonce acp_write outcome, both pendingPermissions (legacy key) + // and pendingPermissionsByNonce must be cleared for that entry. + const events = [ + makePermissionRequestWithAuth(1, "req-b", "nonce-B"), + makePermissionWriteWithNonce( + 2, + "req-b", + "nonce-B", + "selected", + "allow_once", + ), + ]; + const state = buildTranscriptState(events); + + assert.ok( + !state.pendingPermissionsByNonce.has("nonce-B"), + "pendingPermissionsByNonce must be cleared after nonce-B acp_write terminal", + ); + // Legacy key: JSON-encoded requestId scoped by channel:session:turn:id. + const legacyKey = `ch-1:session-1:turn-1:${JSON.stringify("req-b")}`; + assert.ok( + !state.pendingPermissions.has(legacyKey), + "pendingPermissions legacy key must be cleared after acp_write terminal", + ); + // Card outcome must be set. + const card = state.items[0]; + assert.ok(card.outcome, "card must have an outcome after acp_write terminal"); + assert.equal(card.actionable, false); +}); + +// ─── Index cleanup: permission_terminal clears both indexes ─────────────────── + +test("buildTranscript_permission_terminal_clears_both_indexes", () => { + // After a permission_terminal event, both indexes must be cleared for that nonce. + const events = [ + makePermissionRequestWithAuth(1, "req-pt", "nonce-PT"), + makePermissionTerminalEvent(2, "req-pt", "nonce-PT"), + ]; + const state = buildTranscriptState(events); + + assert.ok( + !state.pendingPermissionsByNonce.has("nonce-PT"), + "pendingPermissionsByNonce must be cleared by permission_terminal", + ); + const legacyKey = `ch-1:session-1:turn-1:${JSON.stringify("req-pt")}`; + assert.ok( + !state.pendingPermissions.has(legacyKey), + "pendingPermissions legacy key must be cleared by permission_terminal", + ); +}); + +// ─── Index cleanup: turn_completed backstop clears both indexes ─────────────── + +test("buildTranscript_turn_completed_backstop_clears_both_indexes", () => { + // A turn_completed event must clear any remaining live permission entries + // in both indexes (the backstop for cards not yet retired by their terminal). + const events = [ + makePermissionRequestWithAuth(1, "req-tc", "nonce-TC"), + makeTurnCompleted(2), + ]; + const state = buildTranscriptState(events); + + assert.ok( + !state.pendingPermissionsByNonce.has("nonce-TC"), + "pendingPermissionsByNonce must be cleared by turn_completed backstop", + ); + const legacyKey = `ch-1:session-1:turn-1:${JSON.stringify("req-tc")}`; + assert.ok( + !state.pendingPermissions.has(legacyKey), + "pendingPermissions legacy key must be cleared by turn_completed backstop", + ); + // Card must be retired (not actionable). + const card = state.items.find( + (i) => i.renderClass === "permission" && i.requestNonce === "nonce-TC", + ); + assert.ok(card, "permission card must still exist after turn_completed"); + assert.equal( + card.actionable, + false, + "card must be non-actionable after turn_completed backstop", + ); +}); + +test("buildTranscript_turn_error_backstop_clears_both_indexes", () => { + // Same as turn_completed: a turn_error must also clear both indexes. + const events = [ + makePermissionRequestWithAuth(1, "req-te", "nonce-TE"), + makeTurnError(2), + ]; + const state = buildTranscriptState(events); + + assert.ok( + !state.pendingPermissionsByNonce.has("nonce-TE"), + "pendingPermissionsByNonce must be cleared by turn_error backstop", + ); + const legacyKey = `ch-1:session-1:turn-1:${JSON.stringify("req-te")}`; + assert.ok( + !state.pendingPermissions.has(legacyKey), + "pendingPermissions legacy key must be cleared by turn_error backstop", + ); +}); + +// ─── permission_terminal live replay + archive replay ──────────────────────── + +test("buildTranscript_permission_terminal_retires_card_with_pinned_uncertain_copy", () => { + // permission_terminal must retire the card with the verbatim pinned + // uncertain copy, NOT "denied" or "failed closed". + const events = [ + makePermissionRequestWithAuth(1, "req-live", "nonce-LIVE"), + makePermissionTerminalEvent(2, "req-live", "nonce-LIVE"), + ]; + const transcript = buildTranscript(events); + + assert.equal(transcript.length, 1); + const card = transcript[0]; + assert.equal(card.renderClass, "permission"); + assert.equal( + card.actionable, + false, + "card must be non-actionable after permission_terminal", + ); + assert.match( + card.outcome ?? "", + /Approval outcome unknown.*agent process stopped/i, + "permission_terminal must use the pinned uncertain copy", + ); + assert.doesNotMatch(card.outcome ?? "", /denied/i); + assert.doesNotMatch(card.outcome ?? "", /failed closed/i); +}); + +test("buildTranscript_permission_terminal_in_archive_replay_retires_card", () => { + // In an archive (lifecycle-only) replay the card must be retired by + // permission_terminal. The sequence of events is the same as live replay; + // what changes is the assertion that the card is retired even with no + // subsequent acp_write. + const events = [ + makePermissionRequestWithAuth(1, "req-arc", "nonce-ARC"), + makePermissionTerminalEvent(2, "req-arc", "nonce-ARC"), + ]; + const state = buildTranscriptState(events); + const card = state.items.find( + (i) => i.renderClass === "permission" && i.requestNonce === "nonce-ARC", + ); + + assert.ok(card, "permission card must exist in archive replay"); + assert.equal( + card.actionable, + false, + "card must be non-actionable after permission_terminal in archive replay", + ); + assert.match( + card.outcome ?? "", + /Approval outcome unknown/i, + "archive replay permission_terminal must set the uncertain outcome copy", + ); + // Both indexes must be clean. + assert.ok( + !state.pendingPermissionsByNonce.has("nonce-ARC"), + "nonce index must be clean after archive replay", + ); +}); diff --git a/desktop/src/features/agents/ui/agentSessionTranscript.ts b/desktop/src/features/agents/ui/agentSessionTranscript.ts index 624855282..042331b31 100644 --- a/desktop/src/features/agents/ui/agentSessionTranscript.ts +++ b/desktop/src/features/agents/ui/agentSessionTranscript.ts @@ -349,6 +349,20 @@ function retireAllLivePermissionCards(d: TranscriptDraft, channelId: string) { } } } + // Clean up all pendingPermissions entries scoped to this channel. + // Keys use the compound format `ch:session:turn:id` — drop any that start + // with the channel prefix. + const chPrefix = `${channelId}:`; + let permsMutated = false; + for (const key of d.pendingPermissions.keys()) { + if (key.startsWith(chPrefix)) { + if (!permsMutated) { + d.pendingPermissions = new Map(d.pendingPermissions); + permsMutated = true; + } + d.pendingPermissions.delete(key); + } + } } /** @@ -911,12 +925,22 @@ export function processTranscriptEvent( if (existing?.type === "lifecycle") { replaceItem(d, itemId, { ...existing, - outcome: "Uncertain (process restarting)", + outcome: + "Approval outcome unknown; agent process stopped before it could continue.", actionable: false, }); } d.pendingPermissionsByNonce = new Map(d.pendingPermissionsByNonce); d.pendingPermissionsByNonce.delete(nonce); + // Clean up any matching compound legacy entry. + const responseId = jsonRpcId(asRecord(event.payload).id); + if (responseId) { + const legacyKey = `${ch}:${ctx.sessionId ?? ""}:${ctx.turnId ?? ""}:${responseId}`; + if (d.pendingPermissions.has(legacyKey)) { + d.pendingPermissions = new Map(d.pendingPermissions); + d.pendingPermissions.delete(legacyKey); + } + } } } } else if (event.kind === "acp_read" || event.kind === "acp_write") { @@ -963,12 +987,14 @@ export function processTranscriptEvent( d.pendingPermissionsByNonce.set(auth.requestNonce, itemId); } - // Index by JSON-RPC id so the response (acp_write with result.outcome, - // no method) can correlate by id rather than by turn/seq. + // Legacy id index: keyed by compound (channel, session, turn, id) to + // prevent cross-channel / cross-session JSON-RPC id collisions. + // Only used by authorized frames that carry NO nonce (non-ask paths). const requestId = jsonRpcId(payload.id); if (requestId) { + const legacyKey = `${ch}:${ctx.sessionId ?? ""}:${ctx.turnId ?? ""}:${requestId}`; d.pendingPermissions = new Map(d.pendingPermissions); - d.pendingPermissions.set(requestId, { + d.pendingPermissions.set(legacyKey, { itemId, optionNames: request.optionNames, }); @@ -976,10 +1002,12 @@ 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. + // Nonce-keyed correlation is primary and exclusive: + // - If the frame carries a nonce, we look it up in pendingPermissionsByNonce. + // If the nonce is present but unknown (stale/foreign), we DROP the frame — + // we never fall back to the id map, which could resolve the wrong card. + // - If the frame carries NO nonce, we fall back to the legacy compound-key + // id map (channel+session+turn+id) for non-ask synchronized outcomes. const auth = event.authorization; const nonce = auth?.requestNonce; const responseId = jsonRpcId(payload.id); @@ -991,56 +1019,58 @@ export function processTranscriptEvent( // 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) { + if (nonce !== undefined && nonce !== null) { + // Nonce present: nonce-only path. Do NOT fall back on unknown nonce. + const itemIdByNonce = d.pendingPermissionsByNonce.get(nonce); + if (itemIdByNonce) { + 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 nonce index. d.pendingPermissionsByNonce = new Map(d.pendingPermissionsByNonce); d.pendingPermissionsByNonce.delete(nonce); + // Clean up compound legacy key if it matches. + if (responseId) { + const legacyKey = `${ch}:${ctx.sessionId ?? ""}:${ctx.turnId ?? ""}:${responseId}`; + if (d.pendingPermissions.has(legacyKey)) { + d.pendingPermissions = new Map(d.pendingPermissions); + d.pendingPermissions.delete(legacyKey); + } + } } - if (responseId) { + // Unknown nonce: drop frame — do not mutate any card. + } else if (outcomeKind && responseId) { + // No nonce: legacy compound-key fallback for non-ask paths. + const legacyKey = `${ch}:${ctx.sessionId ?? ""}:${ctx.turnId ?? ""}:${responseId}`; + const pendingById = d.pendingPermissions.get(legacyKey); + if (pendingById) { + const optionId = asString(result.optionId) ?? null; + const outcomeText = describePermissionOutcome( + outcomeKind, + optionId, + pendingById.optionNames, + ); + const existing = d.itemsById.get(pendingById.itemId); + if (existing?.type === "lifecycle") { + replaceItem(d, pendingById.itemId, { + ...existing, + outcome: outcomeText, + actionable: false, + }); + } 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, - pendingById.optionNames, - ); - const existing = d.itemsById.get(pendingById.itemId); - if (existing?.type === "lifecycle") { - replaceItem(d, pendingById.itemId, { - ...existing, - outcome: outcomeText, - actionable: false, - }); - d.pendingPermissions = new Map(d.pendingPermissions); - d.pendingPermissions.delete(responseId); + d.pendingPermissions.delete(legacyKey); } } } else if (event.kind === "acp_write" && method === "session/prompt") { diff --git a/docs/nips/NIP-AO.md b/docs/nips/NIP-AO.md index ff245f92c..6885dd0f8 100644 --- a/docs/nips/NIP-AO.md +++ b/docs/nips/NIP-AO.md @@ -124,12 +124,22 @@ below). It is omitted on all other frame kinds. | `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_terminal` | Observer-only terminal for uncertain permission outcomes (process poison or cancel-during-write). No ACP wire response was confirmed. Carries an `authorization` envelope with `reason = "uncertain"`. Desktop uses this to retire the card without a JSON-RPC response. | 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 correlated by the same nonce — this pairs the challenge and answer in the observer log. +Synchronous policy outcomes (`reject`, `allow`, preflight denial) also produce +`acp_write` frames with `authorization` envelopes. Their `reason` values are: + +| Policy path | `reason` | +|-------------|----------| +| `reject` policy, preflight denial, ask-unavailable downgrade | `"rejected"` | +| `allow` policy (auto-approval succeeded) | `"allowed"` | +| `allow` policy (fail-closed, no unique allow_once option) | `"allow_failed_closed"` | + **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 @@ -165,10 +175,16 @@ call, the `ObserverEvent` carries an `authorization` field: | `"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). | + | `"rejected"` | `reject` policy, preflight denial, or ask-unavailable downgrade; request denied synchronously without an actionable card. | + | `"allowed"` | `allow` policy auto-approval succeeded; request granted synchronously. | + | `"allow_failed_closed"` | `allow` policy but no unique `allow_once` option available; request denied synchronously. | - The `uncertain` terminal (cancel arriving while the write is in flight) does NOT - produce an `acp_write` observer event — instead the harness emits a - `permission_terminal` observer event with `authorization.reason = "uncertain"` so + `"rejected"`, `"allowed"`, and `"allow_failed_closed"` are emitted on `acp_write` frames + for synchronous policy paths (see [Synchronous policy outcomes](#synchronous-policy-outcomes)). + They are NOT emitted for `ask`-policy pending-map entries. + + The `uncertain` outcome does NOT 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