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