diff --git a/crates/buzz-acp/src/acp.rs b/crates/buzz-acp/src/acp.rs index cd049830688..2efcb80d5bc 100644 --- a/crates/buzz-acp/src/acp.rs +++ b/crates/buzz-acp/src/acp.rs @@ -9,11 +9,13 @@ //! 5. [`AcpClient::session_cancel`] / [`AcpClient::cancel_with_cleanup`] — cancel in-flight turn use futures_util::StreamExt; +use sha2::{Digest, Sha256}; use tokio::io::AsyncWriteExt; use tokio::process::{Child, ChildStdin, ChildStdout}; use tokio_util::codec::{FramedRead, LinesCodec, LinesCodecError}; use crate::observer::{ObserverContext, ObserverHandle}; +use crate::permission_ledger::{self, PermissionLedger, PermissionRecord, PermissionState}; use crate::usage::{ PromptResponseUsage, StandardAdapterKind, StandardUsageTracker, TurnUsage, UsageTracker, }; @@ -21,6 +23,35 @@ use crate::usage::{ /// Maximum allowed size of a single NDJSON line from the agent's stdout. /// Lines exceeding this limit are rejected to prevent OOM from rogue agents. const MAX_LINE_SIZE: usize = 10_000_000; // 10 MB +const PERMISSION_DECISION_TTL: std::time::Duration = std::time::Duration::from_secs(300); + +/// An owner decision delivered through Buzz's signed observer-control channel. +/// +/// The harness verifies every field against the still-pending ACP request before +/// writing a response. This is deliberately not an Auto-review override. +#[derive(Clone, Debug)] +pub struct PermissionResolution { + pub turn_id: String, + pub session_id: String, + pub request_id: serde_json::Value, + pub action_digest: String, + pub option_id: String, +} + +#[derive(Clone, Debug)] +struct PendingPermission { + ledger_key: String, + session_id: String, + request_id: serde_json::Value, + action_digest: String, + option_ids: Vec, + expires_at: tokio::time::Instant, +} + +fn permission_action_digest(msg: &serde_json::Value) -> Result { + let canonical = serde_json::to_vec(msg)?; + Ok(hex::encode(Sha256::digest(canonical))) +} /// Package and binary name used by Buzz's Pi ACP fork. pub(crate) const BUZZ_PI_ACP_NAME: &str = "buzz-pi-acp"; @@ -211,6 +242,14 @@ pub struct AcpClient { /// outside of a goose-native turn — the read loop's steer arm is /// disabled in that case. steer_rx: Option>, + /// Per-turn channel for signed owner resolutions of a pending ACP + /// permission request. It is independent from cancellation controls: an + /// approval must not tear down or replay the agent turn. + permission_rx: Option>, + pending_permission: Option, + /// One process-shared, Desktop-owned durable ledger. `None` means this is + /// not a managed Desktop runtime, so permission requests are cancelled. + permission_ledger: Option>>, /// Usage tracker for goose/buzz-agent's cumulative notification format. goose_usage: UsageTracker, /// Per-turn prompt-response usage and Claude's optional cumulative cost. @@ -557,6 +596,17 @@ impl AcpClient { .take() .ok_or_else(|| AcpError::Protocol("failed to open agent stdout".into()))?; + let permission_ledger = match ( + std::env::var_os("BUZZ_ACP_PERMISSION_LEDGER_PATH"), + std::env::var("BUZZ_MANAGED_AGENT_START_NONCE").ok(), + ) { + (Some(path), Some(start_nonce)) => Some( + PermissionLedger::shared(std::path::PathBuf::from(path), start_nonce) + .map_err(AcpError::Protocol)?, + ), + _ => None, + }; + Ok(Self { child, stdin, @@ -572,6 +622,9 @@ impl AcpClient { active_run_id: None, steering_supported: false, steer_rx: None, + permission_rx: None, + pending_permission: None, + permission_ledger, goose_usage: UsageTracker::default(), standard_usage: StandardUsageTracker::default(), standard_adapter, @@ -940,6 +993,27 @@ impl AcpClient { self.steer_rx = Some(rx); } + /// Install the owner-resolution receiver for one channel turn. + pub fn install_permission_rx(&mut self, rx: tokio::sync::mpsc::Receiver) { + debug_assert!( + self.permission_rx.is_none(), + "install_permission_rx: previous turn's receiver was not consumed" + ); + self.permission_rx = Some(rx); + } + + #[cfg(test)] + fn install_test_permission_ledger(&mut self, name: &str) { + let directory = std::env::temp_dir().join(format!( + "buzz-acp-permission-fixture-{name}-{}", + uuid::Uuid::new_v4() + )); + let path = directory.join("permission-lifecycle.json"); + let ledger = PermissionLedger::open(path, format!("fixture-{name}")) + .expect("test permission ledger opens"); + self.permission_ledger = Some(std::sync::Arc::new(std::sync::Mutex::new(ledger))); + } + /// Clear any installed steer receiver without consuming it. /// /// Called by `send_prompt_result` on every exit path of `run_prompt_task` @@ -948,6 +1022,29 @@ impl AcpClient { /// Idempotent — safe to call when `steer_rx` is already `None`. pub fn clear_steer_rx(&mut self) { self.steer_rx = None; + self.permission_rx = None; + if let Some(pending) = self.pending_permission.take() { + let persisted = self + .permission_ledger + .as_ref() + .and_then(|ledger| ledger.lock().ok()) + .map(|mut ledger| { + ledger.transition_pending(&pending.ledger_key, PermissionState::Abandoned) + }) + .transpose(); + self.observe( + "permission_abandoned", + serde_json::json!({ + "requestId": pending.request_id, + "sessionId": pending.session_id, + "actionDigest": pending.action_digest, + "outcome": "agent_session_ended", + "persisted": persisted.is_ok(), + }), + ); + } + self.pending_permission_id = None; + self.permission_responded = false; } /// Returns `true` if no steer receiver is currently installed. @@ -1044,6 +1141,18 @@ impl AcpClient { // but only if we haven't already responded (guards against double-response race). if let Some(perm_id) = self.pending_permission_id.clone() { if !self.permission_responded { + if let Some(pending) = self.pending_permission.as_ref() { + let ledger = self.permission_ledger.as_ref().ok_or_else(|| { + AcpError::Protocol("permission lifecycle ledger is unavailable".into()) + })?; + ledger + .lock() + .map_err(|_| { + AcpError::Protocol("permission lifecycle ledger lock poisoned".into()) + })? + .transition_pending(&pending.ledger_key, PermissionState::Cancelled) + .map_err(AcpError::Protocol)?; + } let response = permission_response_cancelled(&perm_id); self.write_ndjson(&response).await?; tracing::debug!( @@ -1053,6 +1162,7 @@ impl AcpClient { } self.pending_permission_id = None; self.permission_responded = false; + self.pending_permission = None; } // Step 2: send session/cancel notification (no id) @@ -1199,7 +1309,8 @@ impl AcpClient { /// /// While waiting, handles: /// - `session/update` notifications → logged via tracing - /// - `session/request_permission` requests → auto-approved with `allow_once` + /// - `session/request_permission` requests → held for an authenticated + /// owner decision that is bound to the exact ACP request /// - Any other messages → debug-logged and ignored; if they carry an `id` /// (i.e. they are requests, not notifications), a JSON-RPC -32601 error is sent. /// @@ -1210,10 +1321,34 @@ impl AcpClient { expected_id: u64, ) -> Result { loop { + // Setup and prompt requests both use this production dispatch loop. + // Keep owner resolution in the same loop so an adapter cannot turn a + // permission request during setup into the old unconditional allow. + let read_result = tokio::select! { + _ = async { + match self.pending_permission.as_ref() { + Some(pending) => tokio::time::sleep_until(pending.expires_at).await, + None => std::future::pending().await, + } + } => { + self.expire_pending_permission().await?; + continue; + } + Some(resolution) = async { + match self.permission_rx.as_mut() { + Some(rx) => rx.recv().await, + None => None, + } + } => { + self.resolve_pending_permission(resolution).await?; + continue; + } + read_result = self.reader.next() => read_result, + }; // LinesCodec::new_with_max_length enforces MAX_LINE_SIZE at the // read level — the buffer never grows beyond the limit, preventing // OOM from rogue agents writing infinite non-newline bytes. - let line = match self.reader.next().await { + let line = match read_result { None => return Err(AcpError::AgentExited), Some(Err(LinesCodecError::MaxLineLengthExceeded)) => { return Err(AcpError::Protocol( @@ -1277,6 +1412,11 @@ impl AcpClient { "session/request_permission" => { self.handle_permission_request(&msg).await?; } + "$/cancel_request" if msg.get("id").is_none() => { + let session_id = self.observer_context.session_id.clone(); + self.handle_permission_cancel_request(&msg, session_id.as_deref()) + .await?; + } other => { // If the unknown message has an id, it's a request expecting a reply. // Silence would cause the agent to hang waiting for a response. @@ -1344,6 +1484,7 @@ impl AcpClient { // Dropped at scope exit (return paths drain `pending_steer` first // so the ack_tx oneshot is never leaked silently). let mut steer_rx = self.steer_rx.take(); + let mut permission_rx = self.permission_rx.take(); // Tracks the in-flight steer write: `(request_id, transport, ack_tx)`. // While `Some`, the steer arm is gated off so we don't stack writes, @@ -1402,6 +1543,24 @@ impl AcpClient { // read level — the buffer never grows beyond the limit. let read_result = tokio::select! { biased; + _ = async { + match self.pending_permission.as_ref() { + Some(pending) => tokio::time::sleep_until(pending.expires_at).await, + None => std::future::pending().await, + } + } => { + self.expire_pending_permission().await?; + continue; + } + Some(resolution) = async { + match permission_rx.as_mut() { + Some(rx) => rx.recv().await, + None => None, + } + } => { + self.resolve_pending_permission(resolution).await?; + continue; + } read_result = self.reader.next() => Some(read_result), // Steer arm: gated off whenever a steer write is already in // flight so we don't stack two writes against the same @@ -1725,6 +1884,10 @@ impl AcpClient { "session/request_permission" => { self.handle_permission_request(&msg).await?; } + "$/cancel_request" if msg.get("id").is_none() => { + self.handle_permission_cancel_request(&msg, Some(session_id)) + .await?; + } other => { // If the unknown message has an id, it's a request expecting a reply. // Silence would cause the agent to hang waiting for a response. @@ -1939,89 +2102,304 @@ impl AcpClient { } } - /// Auto-approve a `session/request_permission` request from the agent. - /// - /// Finds the option with `kind == "allow_once"` and responds with its `optionId`. - /// If no `allow_once` option exists, falls back to `reject_once`. + /// Park a `session/request_permission` until its authenticated owner resolves it. /// - /// **Critical:** Never hardcode `optionId` — always find it dynamically by `kind`. - /// - /// The request `id` is stored as `serde_json::Value` to support both numeric - /// and string IDs per JSON-RPC 2.0. + /// A transcript frame is evidence only. The only route that writes an ACP + /// selection is a signed observer-control frame whose request ID, action + /// digest, and offered option ID match this exact pending request. async fn handle_permission_request(&mut self, msg: &serde_json::Value) -> Result<(), AcpError> { - // Extract id as a Value — JSON-RPC 2.0 allows both numeric and string IDs. let id = msg .get("id") .cloned() .ok_or_else(|| AcpError::Protocol("permission request missing id".into()))?; - - // Store pending permission id so cancel_with_cleanup can respond to it. + let session_id = msg + .pointer("/params/sessionId") + .and_then(serde_json::Value::as_str) + .ok_or_else(|| AcpError::Protocol("permission request missing sessionId".into()))?; + if self.pending_permission.is_some() { + return Err(AcpError::Protocol( + "received a second permission request while one is pending".into(), + )); + } self.pending_permission_id = Some(id.clone()); - // Mark as not yet responded — guards against double-response race. self.permission_responded = false; - let options = msg["params"]["options"] .as_array() .ok_or_else(|| AcpError::Protocol("permission request missing options".into()))?; - - tracing::debug!( - target: "acp::permission", - "session/request_permission id={id}, {} options", - options.len() + let action_digest = permission_action_digest(msg)?; + let option_ids = options + .iter() + .map(|option| { + option + .get("optionId") + .and_then(serde_json::Value::as_str) + .map(ToOwned::to_owned) + .ok_or_else(|| AcpError::Protocol("permission option missing optionId".into())) + }) + .collect::, _>>()?; + let expires_at = tokio::time::Instant::now() + PERMISSION_DECISION_TTL; + let Some(ledger_handle) = self.permission_ledger.as_ref().cloned() else { + // A managed owner card without an acknowledged host ledger could + // survive neither a dropped observer frame nor a process crash. + self.write_ndjson(&permission_response_cancelled(&id)) + .await?; + self.pending_permission_id = None; + self.permission_responded = true; + self.observe( + "permission_unavailable", + serde_json::json!({"reason": "durable_ledger_unavailable"}), + ); + return Ok(()); + }; + let expires_at_wall = chrono::Utc::now() + .checked_add_signed( + chrono::Duration::from_std(PERMISSION_DECISION_TTL).unwrap_or_default(), + ) + .unwrap_or_else(chrono::Utc::now) + .to_rfc3339(); + let persist_result = { + let mut ledger = ledger_handle.lock().map_err(|_| { + AcpError::Protocol("permission lifecycle ledger lock poisoned".into()) + })?; + let turn_id = self.observer_context.turn_id.clone().ok_or_else(|| { + AcpError::Protocol("permission request has no managed turn identity".into()) + })?; + let key = permission_ledger::record_key( + ledger.start_nonce(), + &turn_id, + session_id, + &id, + &action_digest, + ); + let record = PermissionRecord { + key: key.clone(), + channel_id: self.observer_context.channel_id.clone(), + start_nonce: ledger.start_nonce().to_owned(), + turn_id, + session_id: session_id.to_owned(), + request_id: id.clone(), + action_digest: action_digest.clone(), + request: msg.clone(), + options: options.clone(), + expires_at: expires_at_wall.clone(), + updated_at: chrono::Utc::now().to_rfc3339(), + state: PermissionState::Pending, + }; + ledger.record_pending(record).map(|()| key) + }; + let key = match persist_result { + Ok(key) => key, + Err(error) => { + self.write_ndjson(&permission_response_cancelled(&id)) + .await?; + self.pending_permission_id = None; + self.permission_responded = true; + self.observe( + "permission_unavailable", + serde_json::json!({"reason": "pending_persist_failed"}), + ); + return Err(AcpError::Protocol(error)); + } + }; + self.pending_permission = Some(PendingPermission { + ledger_key: key, + session_id: session_id.to_owned(), + request_id: id.clone(), + action_digest: action_digest.clone(), + option_ids: option_ids.clone(), + expires_at, + }); + self.observe( + "permission_pending", + serde_json::json!({ + "requestId": id, + "sessionId": session_id, + "actionDigest": action_digest, + "options": options, + "expiresAt": expires_at_wall, + }), ); + Ok(()) + } - // Find allow_once by kind — NEVER hardcode optionId. - let allow_once = options - .iter() - .find(|opt| opt.get("kind").and_then(|k| k.as_str()) == Some("allow_once")); + /// Cancel only the still-pending permission request named by the ACP SDK. + /// Persist a terminal state before replying, so a failed pipe write can + /// never leave an executable owner decision behind. + async fn handle_permission_cancel_request( + &mut self, + msg: &serde_json::Value, + active_session_id: Option<&str>, + ) -> Result<(), AcpError> { + let Some(request_id) = msg.pointer("/params/requestId") else { + self.observe( + "permission_cancellation_rejected", + serde_json::json!({"reason": "missing_request_id"}), + ); + return Ok(()); + }; + let Some(pending) = self.pending_permission.as_ref().cloned() else { + self.observe( + "permission_cancellation_rejected", + serde_json::json!({"reason": "no_pending_request"}), + ); + return Ok(()); + }; + let supplied_session_id = msg.pointer("/params/sessionId"); + if request_id != &pending.request_id + || self.pending_permission_id.as_ref() != Some(&pending.request_id) + || self.permission_responded + || active_session_id.is_some_and(|id| id != pending.session_id.as_str()) + || supplied_session_id + .is_some_and(|id| id.as_str() != Some(pending.session_id.as_str())) + { + self.observe( + "permission_cancellation_rejected", + serde_json::json!({"reason": "binding_mismatch"}), + ); + return Ok(()); + } + let ledger_handle = self.permission_ledger.as_ref().cloned().ok_or_else(|| { + AcpError::Protocol("permission lifecycle ledger is unavailable".into()) + })?; + ledger_handle + .lock() + .map_err(|_| AcpError::Protocol("permission lifecycle ledger lock poisoned".into()))? + .transition_pending(&pending.ledger_key, PermissionState::Cancelled) + .map_err(AcpError::Protocol)?; + let response = permission_response_cancelled(&pending.request_id); + if let Err(error) = self.write_ndjson(&response).await { + // The bytes may or may not have reached the adapter. Keep the + // owner card terminal and never replay a decision after recovery. + ledger_handle + .lock() + .map_err(|_| { + AcpError::Protocol("permission lifecycle ledger lock poisoned".into()) + })? + .transition_cancelled_to_delivery_unknown(&pending.ledger_key) + .map_err(AcpError::Protocol)?; + self.observe( + "permission_delivery_unknown", + serde_json::json!({ + "reason": "cancellation_response_write_failed", + "requestId": pending.request_id, + "sessionId": pending.session_id, + }), + ); + return Err(error); + } + self.pending_permission = None; + self.pending_permission_id = None; + self.permission_responded = true; + Ok(()) + } - let response = if let Some(opt) = allow_once { - let option_id = opt["optionId"] - .as_str() - .ok_or_else(|| AcpError::Protocol("allow_once option missing optionId".into()))?; - tracing::info!( - target: "acp::permission", - "auto-approving permission id={id} with allow_once optionId={option_id:?}" + async fn resolve_pending_permission( + &mut self, + resolution: PermissionResolution, + ) -> Result<(), AcpError> { + let Some(pending) = self.pending_permission.take() else { + self.observe( + "permission_resolution_rejected", + serde_json::json!({"reason": "no_pending_request"}), ); - permission_response_selected(&id, option_id) - } else { - // No allow_once — fall back to reject_once. - tracing::warn!( - target: "acp::permission", - "no allow_once option found in permission request id={id}, falling back to reject_once" + return Ok(()); + }; + if self.observer_context.turn_id.as_deref() != Some(resolution.turn_id.as_str()) { + self.observe( + "permission_resolution_rejected", + serde_json::json!({"reason": "turn_mismatch"}), ); - let reject = options + self.pending_permission = Some(pending); + return Ok(()); + } + if pending.expires_at <= tokio::time::Instant::now() { + self.pending_permission = Some(pending); + self.expire_pending_permission().await?; + return Ok(()); + } + if pending.session_id != resolution.session_id + || pending.request_id != resolution.request_id + || pending.action_digest != resolution.action_digest + || !pending + .option_ids .iter() - .find(|opt| opt.get("kind").and_then(|k| k.as_str()) == Some("reject_once")); - - if let Some(opt) = reject { - let option_id = opt["optionId"].as_str().unwrap_or("reject"); - permission_response_selected(&id, option_id) - } else { - return Err(AcpError::Protocol( - "no suitable permission option found (neither allow_once nor reject_once)" - .into(), - )); - } + .any(|id| id == &resolution.option_id) + { + self.observe( + "permission_resolution_rejected", + serde_json::json!({"reason": "binding_mismatch"}), + ); + self.pending_permission = Some(pending); + return Ok(()); + } + let Some(ledger_handle) = self.permission_ledger.as_ref().cloned() else { + self.pending_permission = Some(pending); + return Err(AcpError::Protocol( + "permission lifecycle ledger is unavailable".into(), + )); }; + { + let mut ledger = ledger_handle.lock().map_err(|_| { + AcpError::Protocol("permission lifecycle ledger lock poisoned".into()) + })?; + ledger + .transition_pending( + &pending.ledger_key, + PermissionState::DecisionConsumed(resolution.option_id.clone()), + ) + .map_err(AcpError::Protocol)?; + ledger + .transition_consumed_to_delivery_attempt(&pending.ledger_key) + .map_err(AcpError::Protocol)?; + } + let response = permission_response_selected(&pending.request_id, &resolution.option_id); + if let Err(error) = self.write_ndjson(&response).await { + self.observe( + "permission_delivery_unknown", + serde_json::json!({"reason": "response_write_failed"}), + ); + return Err(error); + } + ledger_handle + .lock() + .map_err(|_| AcpError::Protocol("permission lifecycle ledger lock poisoned".into()))? + .transition_delivery_attempt_to_selected(&pending.ledger_key) + .map_err(AcpError::Protocol)?; + self.permission_responded = true; + self.pending_permission_id = None; + self.observe( + "permission_resolved", + serde_json::json!({"outcome": "selected", "optionId": resolution.option_id}), + ); + Ok(()) + } - // Write the response first, then mark as responded. - // - // Previous ordering (flag-before-write) was intended to guard against a - // double-response if a timeout fires between write and flag-set. However, - // the deadlock risk is worse: if write_ndjson fails (e.g. WriteTimeout), - // the flag would be true but no response was actually sent. Then - // cancel_with_cleanup would see permission_responded=true, skip sending - // the cancelled outcome, and the agent would hang waiting for a reply - // that never arrives — a guaranteed deadlock. - // - // The correct fix: set the flag AFTER a successful write. The double- - // response window (between write completion and flag-set) is negligibly - // small and bounded by a single memory store; the deadlock window was - // unbounded. - self.write_ndjson(&response).await?; + async fn expire_pending_permission(&mut self) -> Result<(), AcpError> { + let Some(pending) = self.pending_permission.take() else { + return Ok(()); + }; + let Some(ledger_handle) = self.permission_ledger.as_ref().cloned() else { + self.pending_permission = Some(pending); + return Err(AcpError::Protocol( + "permission lifecycle ledger is unavailable".into(), + )); + }; + ledger_handle + .lock() + .map_err(|_| AcpError::Protocol("permission lifecycle ledger lock poisoned".into()))? + .transition_pending(&pending.ledger_key, PermissionState::Expired) + .map_err(AcpError::Protocol)?; + let response = permission_response_cancelled(&pending.request_id); + if let Err(error) = self.write_ndjson(&response).await { + self.pending_permission = Some(pending); + return Err(error); + } self.permission_responded = true; self.pending_permission_id = None; + self.observe( + "permission_expired", + serde_json::json!({"outcome": "cancelled"}), + ); Ok(()) } @@ -3403,6 +3781,410 @@ mod tests { assert_eq!(result.unwrap()["worked"], serde_json::json!(true)); } + #[tokio::test] + async fn sdk_cancel_request_cancels_only_exact_pending_permission_in_non_idle_loop() { + let permission = serde_json::json!({ + "jsonrpc": "2.0", "id": "17", "method": "session/request_permission", + "params": {"sessionId": "session-cedar", "options": [{"optionId": "grant-cedar"}]} + }); + let output = + std::env::temp_dir().join(format!("buzz-acp-cancel-non-idle-{}", uuid::Uuid::new_v4())); + let script = format!( + r#"read _request +echo '{}' +echo '{{"jsonrpc":"2.0","method":"$/cancel_request","params":{{"requestId":17}}}}' +if read -t 1 _unexpected; then exit 7; fi +echo '{{"jsonrpc":"2.0","method":"$/cancel_request","params":{{"requestId":"17","sessionId":"other-session"}}}}' +if read -t 1 _unexpected; then exit 7; fi +echo '{{"jsonrpc":"2.0","method":"$/cancel_request","params":{{"requestId":"17"}}}}' +read _cancelled +printf '%s\n' "$_cancelled" > '{}' +echo '{{"jsonrpc":"2.0","method":"$/cancel_request","params":{{"requestId":"17"}}}}' +sleep 0.15 +if read -t 1 _unexpected; then exit 7; fi +echo '{{"jsonrpc":"2.0","id":0,"result":{{"ok":true}}}}'"#, + serde_json::to_string(&permission) + .unwrap() + .replace('\'', "'\\''"), + output.display(), + ); + let mut client = spawn_script(&script).await; + client.observer_context.turn_id = Some("turn-cedar".into()); + client.observer_context.session_id = Some("session-cedar".into()); + let observer = crate::observer::ObserverHandle::in_process(); + client.set_observer(Some(observer.clone()), 0); + client.install_test_permission_ledger("cancel-non-idle"); + let (tx, rx) = tokio::sync::mpsc::channel(2); + client.install_permission_rx(rx); + { + let request = client.send_request("fixture/request", serde_json::json!({})); + tokio::pin!(request); + tokio::time::timeout(std::time::Duration::from_secs(5), async { + while !output.exists() { + tokio::select! { + result = &mut request => panic!("request completed before cancellation: {result:?}"), + _ = tokio::time::sleep(std::time::Duration::from_millis(10)) => {} + } + } + }) + .await + .unwrap(); + tx.send(PermissionResolution { + turn_id: "turn-cedar".into(), + session_id: "session-cedar".into(), + request_id: serde_json::json!("17"), + action_digest: permission_action_digest(&permission).unwrap(), + option_id: "grant-cedar".into(), + }) + .await + .unwrap(); + assert_eq!(request.await.unwrap()["ok"], serde_json::json!(true)); + } + assert_eq!( + std::fs::read_to_string(&output).unwrap(), + "{\"id\":\"17\",\"jsonrpc\":\"2.0\",\"result\":{\"outcome\":{\"outcome\":\"cancelled\"}}}\n" + ); + std::fs::remove_file(output).unwrap(); + assert!(client.pending_permission.is_none()); + assert!(client.pending_permission_id.is_none()); + let ledger = client.permission_ledger.as_ref().unwrap().lock().unwrap(); + let key = permission_ledger::record_key( + ledger.start_nonce(), + "turn-cedar", + "session-cedar", + &serde_json::json!("17"), + &permission_action_digest(&permission).unwrap(), + ); + assert_eq!(ledger.get(&key).unwrap().state, PermissionState::Cancelled); + assert!(observer.snapshot().iter().any(|event| { + event.kind == "permission_resolution_rejected" + && event.payload["reason"] == "no_pending_request" + })); + } + + #[tokio::test] + async fn sdk_cancel_request_failed_response_write_stays_terminal_and_refuses_owner() { + let mut client = spawn_script("exec sleep 10").await; + client.observer_context.turn_id = Some("turn-failed-write".into()); + client.observer_context.session_id = Some("session-failed-write".into()); + let observer = crate::observer::ObserverHandle::in_process(); + client.set_observer(Some(observer.clone()), 0); + client.install_test_permission_ledger("cancel-failed-write"); + let permission = serde_json::json!({ + "jsonrpc": "2.0", "id": "failure-id", "method": "session/request_permission", + "params": {"sessionId": "session-failed-write", "options": [{"optionId": "grant"}]} + }); + client.handle_permission_request(&permission).await.unwrap(); + client.child.kill().await.unwrap(); + client.child.wait().await.unwrap(); + let cancellation = serde_json::json!({ + "jsonrpc": "2.0", "method": "$/cancel_request", + "params": {"requestId": "failure-id"} + }); + assert!(client + .handle_permission_cancel_request(&cancellation, Some("session-failed-write")) + .await + .is_err()); + let ledger_key = client + .pending_permission + .as_ref() + .unwrap() + .ledger_key + .clone(); + assert_eq!( + client.pending_permission_id, + Some(serde_json::json!("failure-id")) + ); + let ledger = client.permission_ledger.as_ref().unwrap().clone(); + assert_eq!( + ledger.lock().unwrap().get(&ledger_key).unwrap().state, + PermissionState::DeliveryUnknown + ); + client + .resolve_pending_permission(PermissionResolution { + turn_id: "turn-failed-write".into(), + session_id: "session-failed-write".into(), + request_id: serde_json::json!("failure-id"), + action_digest: permission_action_digest(&permission).unwrap(), + option_id: "grant".into(), + }) + .await + .unwrap_err(); + assert_eq!( + ledger.lock().unwrap().get(&ledger_key).unwrap().state, + PermissionState::DeliveryUnknown + ); + assert!(observer.snapshot().iter().any(|event| { + event.kind == "permission_delivery_unknown" + && event.payload["requestId"] == "failure-id" + })); + } + + #[tokio::test] + async fn owner_resolution_binds_exact_permission_in_non_idle_dispatch_loop() { + // This fixture uses deliberately non-obvious option IDs. The adapter + // releases the production request only after it has read the exact + // selected response, so an old allow_once response cannot pass it. + let permission = serde_json::json!({ + "jsonrpc": "2.0", + "id": "perm-cedar-17", + "method": "session/request_permission", + "params": { + "sessionId": "session-cedar", + "title": "fixture request", + "options": [ + {"optionId": "grant-lantern-9", "kind": "allow_once", "name": "Grant lantern"}, + {"optionId": "deny-otter-4", "kind": "reject_once", "name": "Deny otter"} + ] + } + }); + let permission_wire = serde_json::to_string(&permission).unwrap(); + let script = format!( + r#"read _request +echo '{}' +read _decision +case "$_decision" in + *'"id":"perm-cedar-17"'*) ;; + *) exit 7 ;; +esac +case "$_decision" in + *'"outcome":"selected"'*) ;; + *) exit 7 ;; +esac +case "$_decision" in + *'"optionId":"grant-lantern-9"'*) echo '{{"jsonrpc":"2.0","id":0,"result":{{"ok":true}}}}' ;; + *) exit 7 ;; +esac"#, + permission_wire.replace('\'', "'\\''") + ); + let mut client = spawn_script(&script).await; + client.observer_context.turn_id = Some("turn-cedar".into()); + client.install_test_permission_ledger("cedar"); + let (tx, rx) = tokio::sync::mpsc::channel(4); + client.install_permission_rx(rx); + let digest = permission_action_digest(&permission).unwrap(); + let request = client.send_request("fixture/request", serde_json::json!({})); + tokio::pin!(request); + + tokio::select! { + result = &mut request => panic!("permission auto-resolved before owner decision: {result:?}"), + _ = tokio::time::sleep(std::time::Duration::from_millis(40)) => {} + } + tx.send(PermissionResolution { + turn_id: "turn-cedar".into(), + session_id: "session-cedar".into(), + request_id: serde_json::json!("perm-cedar-17"), + action_digest: "changed-action".into(), + option_id: "grant-lantern-9".into(), + }) + .await + .unwrap(); + tokio::select! { + result = &mut request => panic!("changed action resolved permission: {result:?}"), + _ = tokio::time::sleep(std::time::Duration::from_millis(40)) => {} + } + tx.send(PermissionResolution { + turn_id: "turn-cedar".into(), + session_id: "session-cedar".into(), + request_id: serde_json::json!("perm-cedar-17"), + action_digest: digest.clone(), + option_id: "invented-option".into(), + }) + .await + .unwrap(); + tokio::select! { + result = &mut request => panic!("unoffered option resolved permission: {result:?}"), + _ = tokio::time::sleep(std::time::Duration::from_millis(40)) => {} + } + tx.send(PermissionResolution { + turn_id: "turn-cedar".into(), + session_id: "session-cedar".into(), + request_id: serde_json::json!("perm-cedar-17"), + action_digest: digest, + option_id: "grant-lantern-9".into(), + }) + .await + .unwrap(); + assert_eq!(request.await.unwrap()["ok"], serde_json::json!(true)); + } + + #[tokio::test] + async fn sdk_cancel_request_cancels_only_exact_pending_permission_in_idle_loop() { + let permission = serde_json::json!({ + "jsonrpc": "2.0", "id": 701, "method": "session/request_permission", + "params": {"sessionId": "session-poppy", "options": [{"optionId": "grant-poppy"}]} + }); + let output = + std::env::temp_dir().join(format!("buzz-acp-cancel-idle-{}", uuid::Uuid::new_v4())); + let script = format!( + r#"read _prompt +echo '{}' +echo '{{"jsonrpc":"2.0","method":"$/cancel_request","params":{{"requestId":"701"}}}}' +if read -t 1 _unexpected; then exit 7; fi +echo '{{"jsonrpc":"2.0","method":"$/cancel_request","params":{{"requestId":701,"sessionId":"other-session"}}}}' +if read -t 1 _unexpected; then exit 7; fi +echo '{{"jsonrpc":"2.0","method":"$/cancel_request","params":{{"requestId":701}}}}' +read _cancelled +printf '%s\n' "$_cancelled" > '{}' +echo '{{"jsonrpc":"2.0","method":"$/cancel_request","params":{{"requestId":701}}}}' +sleep 0.15 +if read -t 1 _unexpected; then exit 7; fi +echo '{{"jsonrpc":"2.0","id":0,"result":{{"stopReason":"end_turn"}}}}'"#, + serde_json::to_string(&permission) + .unwrap() + .replace('\'', "'\\''"), + output.display(), + ); + let mut client = spawn_script(&script).await; + client.observer_context.turn_id = Some("turn-poppy".into()); + client.observer_context.session_id = Some("session-poppy".into()); + let observer = crate::observer::ObserverHandle::in_process(); + client.set_observer(Some(observer.clone()), 0); + client.install_test_permission_ledger("cancel-idle"); + let (tx, rx) = tokio::sync::mpsc::channel(2); + client.install_permission_rx(rx); + { + let prompt = client.session_prompt_with_idle_timeout( + "session-poppy", + "fixture prompt", + std::time::Duration::from_secs(5), + std::time::Duration::from_secs(10), + ); + tokio::pin!(prompt); + tokio::time::timeout(std::time::Duration::from_secs(5), async { + while !output.exists() { + tokio::select! { + result = &mut prompt => panic!("prompt completed before cancellation: {result:?}"), + _ = tokio::time::sleep(std::time::Duration::from_millis(10)) => {} + } + } + }) + .await + .unwrap(); + tx.send(PermissionResolution { + turn_id: "turn-poppy".into(), + session_id: "session-poppy".into(), + request_id: serde_json::json!(701), + action_digest: permission_action_digest(&permission).unwrap(), + option_id: "grant-poppy".into(), + }) + .await + .unwrap(); + assert_eq!(prompt.await.unwrap(), StopReason::EndTurn); + } + assert_eq!( + std::fs::read_to_string(&output).unwrap(), + "{\"id\":701,\"jsonrpc\":\"2.0\",\"result\":{\"outcome\":{\"outcome\":\"cancelled\"}}}\n" + ); + std::fs::remove_file(output).unwrap(); + assert!(client.pending_permission.is_none()); + assert!(client.pending_permission_id.is_none()); + let ledger = client.permission_ledger.as_ref().unwrap().lock().unwrap(); + let key = permission_ledger::record_key( + ledger.start_nonce(), + "turn-poppy", + "session-poppy", + &serde_json::json!(701), + &permission_action_digest(&permission).unwrap(), + ); + assert_eq!(ledger.get(&key).unwrap().state, PermissionState::Cancelled); + assert!(observer.snapshot().iter().any(|event| { + event.kind == "permission_resolution_rejected" + && event.payload["reason"] == "no_pending_request" + })); + } + + #[tokio::test] + async fn owner_resolution_binds_exact_permission_in_idle_prompt_dispatch_loop() { + let permission = serde_json::json!({ + "jsonrpc": "2.0", + "id": 701, + "method": "session/request_permission", + "params": { + "sessionId": "session-poppy", + "options": [{"optionId": "reject-poppy-2", "kind": "reject_once"}] + } + }); + let permission_wire = serde_json::to_string(&permission).unwrap(); + let script = format!( + r#"read _prompt +echo '{}' +read _decision +case "$_decision" in + *'"id":701'*) ;; + *) exit 7 ;; +esac +case "$_decision" in + *'"outcome":"selected"'*) ;; + *) exit 7 ;; +esac +case "$_decision" in + *'"optionId":"reject-poppy-2"'*) echo '{{"jsonrpc":"2.0","id":0,"result":{{"stopReason":"end_turn"}}}}' ;; + *) exit 7 ;; +esac"#, + permission_wire.replace('\'', "'\\''") + ); + let mut client = spawn_script(&script).await; + client.observer_context.turn_id = Some("turn-poppy".into()); + client.install_test_permission_ledger("poppy"); + let (tx, rx) = tokio::sync::mpsc::channel(2); + client.install_permission_rx(rx); + let digest = permission_action_digest(&permission).unwrap(); + let prompt = client.session_prompt_with_idle_timeout( + "session-poppy", + "fixture prompt", + std::time::Duration::from_secs(1), + std::time::Duration::from_secs(1), + ); + tokio::pin!(prompt); + tokio::select! { + result = &mut prompt => panic!("permission auto-resolved in idle loop: {result:?}"), + _ = tokio::time::sleep(std::time::Duration::from_millis(40)) => {} + } + tx.send(PermissionResolution { + turn_id: "turn-poppy".into(), + session_id: "session-poppy".into(), + request_id: serde_json::json!(701), + action_digest: digest, + option_id: "reject-poppy-2".into(), + }) + .await + .unwrap(); + assert_eq!(prompt.await.unwrap(), StopReason::EndTurn); + } + + #[tokio::test] + async fn cancel_with_cleanup_ends_pending_owner_wait_with_cancelled_outcome() { + let script = r#"read _prompt +echo '{"jsonrpc":"2.0","id":"perm-cancel-8","method":"session/request_permission","params":{"sessionId":"session-cancel","options":[{"optionId":"grant-should-not-appear"}]}}' +read decision +case "$decision" in *'"outcome":"cancelled"'*) echo '{"jsonrpc":"2.0","id":0,"result":{"stopReason":"cancelled"}}' ;; *) exit 7 ;; esac"#; + let mut client = spawn_script(script).await; + client.observer_context.turn_id = Some("turn-cancel".into()); + client.install_test_permission_ledger("cancel"); + let result = client + .session_prompt_with_idle_timeout( + "session-cancel", + "fixture prompt", + std::time::Duration::from_millis(30), + std::time::Duration::from_secs(1), + ) + .await; + assert!( + matches!(result, Err(AcpError::IdleTimeout(_))), + "expected owner wait to remain pending, got {result:?}" + ); + assert_eq!( + client + .cancel_with_cleanup_grace("session-cancel", std::time::Duration::from_secs(1)) + .await + .unwrap(), + StopReason::Cancelled + ); + assert!(client.pending_permission.is_none()); + assert!(client.pending_permission_id.is_none()); + } + #[tokio::test] async fn keepalive_resets_idle_past_deadline() { // Keepalive session/update lines every 50ms against a 100ms idle deadline. diff --git a/crates/buzz-acp/src/config.rs b/crates/buzz-acp/src/config.rs index b4d27903c62..714b47f1040 100644 --- a/crates/buzz-acp/src/config.rs +++ b/crates/buzz-acp/src/config.rs @@ -347,6 +347,12 @@ pub struct CliArgs { #[arg(long, env = "BUZZ_ACP_CONFIG", default_value = "./buzz-acp.toml")] pub config: PathBuf, + /// Desktop-owned, owner-only lifecycle ledger for ACP permission requests. + /// Managed Desktop runtimes pass a pair-scoped path. Without it, permission + /// requests are cancelled rather than represented by an unaudited card. + #[arg(long, env = "BUZZ_ACP_PERMISSION_LEDGER_PATH")] + pub permission_ledger_path: Option, + #[arg(long, env = "BUZZ_ACP_DEDUP", default_value = "queue", value_enum)] pub dedup: DedupMode, @@ -566,6 +572,7 @@ pub struct Config { pub channels_override: Option>, pub no_mention_filter: bool, pub config_path: PathBuf, + pub permission_ledger_path: Option, pub context_message_limit: u32, /// Maximum turns per session before proactive rotation. 0 = disabled. pub max_turns_per_session: u32, @@ -1179,6 +1186,7 @@ impl Config { channels_override: args.channels, no_mention_filter: args.no_mention_filter, config_path: args.config, + permission_ledger_path: args.permission_ledger_path, context_message_limit: args.context_message_limit, max_turns_per_session: args.max_turns_per_session, presence_enabled: !args.no_presence, @@ -1558,6 +1566,7 @@ mod tests { channels_override: None, no_mention_filter: false, config_path: PathBuf::from("./buzz-acp.toml"), + permission_ledger_path: None, context_message_limit: 12, max_turns_per_session: 0, presence_enabled: true, diff --git a/crates/buzz-acp/src/lib.rs b/crates/buzz-acp/src/lib.rs index 3ff1dc39898..6f42875c3a8 100644 --- a/crates/buzz-acp/src/lib.rs +++ b/crates/buzz-acp/src/lib.rs @@ -5,6 +5,7 @@ mod config; mod engram_fetch; mod filter; mod observer; +mod permission_ledger; mod pool; mod pool_lifecycle; mod prompt_framing; @@ -1614,6 +1615,9 @@ fn handle_relay_observer_control_event( Some("switch_model") => { handle_switch_model_control(&payload, pool, observer); } + Some("resolve_permission") => { + handle_resolve_permission_control(&payload, pool, observer); + } Some("publish_project_owner_announcements") => { handle_publish_project_owner_announcements_control( &payload, @@ -1628,6 +1632,69 @@ fn handle_relay_observer_control_event( } } +#[derive(serde::Deserialize)] +#[serde(rename_all = "camelCase")] +struct ResolvePermissionControl { + turn_id: String, + session_id: String, + request_id: serde_json::Value, + action_digest: String, + option_id: String, +} + +/// Route an authenticated owner decision to one exact in-flight turn. +/// +/// The owner signature and freshness window are checked before this function. +/// This handler deliberately does not use a channel ID: concurrent thread +/// sessions may share a channel. `AcpClient` performs the second, independent +/// request/session/digest/offered-option binding before it writes ACP output. +fn handle_resolve_permission_control( + payload: &serde_json::Value, + pool: &mut AgentPool, + observer: Option<&observer::ObserverHandle>, +) { + let Ok(control) = serde_json::from_value::(payload.clone()) else { + tracing::warn!("permission resolution control frame has an invalid payload"); + return; + }; + let resolution = acp::PermissionResolution { + turn_id: control.turn_id.clone(), + session_id: control.session_id, + request_id: control.request_id, + action_digest: control.action_digest, + option_id: control.option_id, + }; + let status = match pool + .task_map_mut() + .values_mut() + .find(|meta| meta.turn_id == control.turn_id) + .and_then(|meta| meta.permission_tx.as_ref()) + { + Some(tx) => match tx.try_send(resolution) { + Ok(()) => "sent", + Err(tokio::sync::mpsc::error::TrySendError::Full(_)) => "already_resolving", + Err(tokio::sync::mpsc::error::TrySendError::Closed(_)) => "turn_ending", + }, + None => "no_matching_turn", + }; + if let Some(observer) = observer { + observer.emit( + "control_result", + None, + &observer::ObserverContext { + channel_id: None, + session_id: None, + turn_id: Some(control.turn_id), + started_at: None, + }, + serde_json::json!({ + "type": "resolve_permission", + "status": status, + }), + ); + } +} + #[derive(serde::Deserialize)] #[serde(rename_all = "camelCase")] struct ProjectOwnerAnnouncementControl { @@ -2510,6 +2577,17 @@ async fn tokio_main() -> Result<()> { tracing::info!("buzz-acp starting: {}", config.summary()); + // The Desktop stamps this path after all user-provided environment layers. + // Open it before spawning an ACP client so a managed owner card can never + // exist without an acknowledged host lifecycle store. + if let Some(path) = config.permission_ledger_path.clone() { + let start_nonce = std::env::var("BUZZ_MANAGED_AGENT_START_NONCE").map_err(|_| { + anyhow::anyhow!("permission lifecycle ledger requires managed start nonce") + })?; + permission_ledger::PermissionLedger::shared(path, start_nonce) + .map_err(|error| anyhow::anyhow!("permission lifecycle ledger unavailable: {error}"))?; + } + let observer = config .relay_observer .then(observer::ObserverHandle::in_process); @@ -4522,6 +4600,8 @@ fn dispatch_pending( // Prompt text is now built inside run_prompt_task (needs async for // context fetching). Pass None for prompt_text; batch carries the data. let (control_tx, control_rx) = tokio::sync::oneshot::channel::(); + let (permission_tx, permission_rx) = + tokio::sync::mpsc::channel::(2); let turn_id = Uuid::new_v4().to_string(); let task_turn_id = turn_id.clone(); @@ -4540,7 +4620,7 @@ fn dispatch_pending( None, ctx_clone, result_tx, - Some(control_rx), + (Some(control_rx), Some(permission_rx)), task_turn_id, ) .await; @@ -4556,6 +4636,7 @@ fn dispatch_pending( recoverable_batch, control_tx: Some(control_tx), steer_tx, + permission_tx: Some(permission_tx), successful_steer_deliveries: HashSet::new(), }, ); @@ -5228,7 +5309,7 @@ fn dispatch_heartbeat( Some(prompt_text), ctx_clone, result_tx, - None, + (None, None), task_turn_id, ) .await; @@ -5244,6 +5325,7 @@ fn dispatch_heartbeat( recoverable_batch: None, control_tx: None, steer_tx: None, + permission_tx: None, successful_steer_deliveries: HashSet::new(), }, ); @@ -6070,6 +6152,7 @@ mod owner_control_command_tests { recoverable_batch: None, control_tx: Some(control_tx), steer_tx: None, + permission_tx: None, successful_steer_deliveries: HashSet::new(), }, ); @@ -6116,6 +6199,7 @@ mod owner_control_command_tests { recoverable_batch: None, control_tx: Some(control_tx), steer_tx: None, + permission_tx: None, successful_steer_deliveries: HashSet::new(), }, ); @@ -9154,6 +9238,7 @@ mod build_mcp_servers_tests { channels_override: None, no_mention_filter: false, config_path: std::path::PathBuf::from("./buzz-acp.toml"), + permission_ledger_path: None, context_message_limit: 12, max_turns_per_session: 0, presence_enabled: true, @@ -9380,6 +9465,7 @@ mod error_outcome_emission_tests { channels_override: None, no_mention_filter: false, config_path: std::path::PathBuf::from("./buzz-acp.toml"), + permission_ledger_path: None, context_message_limit: 12, max_turns_per_session: 0, presence_enabled: true, @@ -9485,6 +9571,7 @@ mod error_outcome_emission_tests { recoverable_batch: None, control_tx: None, steer_tx: None, + permission_tx: None, successful_steer_deliveries: HashSet::from([ crate::pool::SuccessfulSteerDelivery { event_id: steer_event_id.into(), @@ -9565,6 +9652,7 @@ mod error_outcome_emission_tests { recoverable_batch: None, control_tx: None, steer_tx: None, + permission_tx: None, successful_steer_deliveries: HashSet::from([ crate::pool::SuccessfulSteerDelivery { event_id: "stale-event".into(), @@ -9687,6 +9775,7 @@ mod error_outcome_emission_tests { recoverable_batch: None, control_tx: None, steer_tx: None, + permission_tx: None, successful_steer_deliveries: HashSet::from([ crate::pool::SuccessfulSteerDelivery { event_id: "stale-event".into(), @@ -9756,6 +9845,7 @@ mod error_outcome_emission_tests { recoverable_batch: None, control_tx: None, steer_tx: None, + permission_tx: None, successful_steer_deliveries: HashSet::new(), }, ); @@ -9836,6 +9926,7 @@ mod error_outcome_emission_tests { recoverable_batch: None, control_tx: None, steer_tx: None, + permission_tx: None, successful_steer_deliveries: HashSet::new(), }, ); @@ -9928,6 +10019,7 @@ mod error_outcome_emission_tests { recoverable_batch: Some(batch), control_tx: None, steer_tx: None, + permission_tx: None, successful_steer_deliveries: HashSet::new(), }, ); @@ -10027,6 +10119,7 @@ mod error_outcome_emission_tests { recoverable_batch: None, control_tx: None, steer_tx: None, + permission_tx: None, successful_steer_deliveries: HashSet::new(), }, ); @@ -10124,6 +10217,7 @@ mod error_outcome_emission_tests { recoverable_batch: None, control_tx: None, steer_tx: None, + permission_tx: None, successful_steer_deliveries: HashSet::new(), }, ); @@ -10232,6 +10326,7 @@ mod error_outcome_emission_tests { recoverable_batch: None, control_tx: None, steer_tx: None, + permission_tx: None, successful_steer_deliveries: HashSet::new(), }, ); @@ -10310,6 +10405,7 @@ mod error_outcome_emission_tests { recoverable_batch: None, control_tx: None, steer_tx: None, + permission_tx: None, successful_steer_deliveries: HashSet::new(), }, ); @@ -10407,6 +10503,7 @@ mod error_outcome_emission_tests { recoverable_batch: None, control_tx: None, steer_tx: None, + permission_tx: None, successful_steer_deliveries: HashSet::new(), }, ); @@ -10527,6 +10624,7 @@ mod error_outcome_emission_tests { recoverable_batch: None, control_tx: None, steer_tx: None, + permission_tx: None, successful_steer_deliveries: HashSet::new(), }, ); @@ -10669,6 +10767,7 @@ mod error_outcome_emission_tests { recoverable_batch: None, control_tx: None, steer_tx: None, + permission_tx: None, successful_steer_deliveries: HashSet::new(), }, ); @@ -10803,6 +10902,7 @@ mod error_outcome_emission_tests { recoverable_batch: None, control_tx: None, steer_tx: None, + permission_tx: None, successful_steer_deliveries: HashSet::new(), }, ); @@ -10958,6 +11058,7 @@ mod error_outcome_emission_tests { recoverable_batch: None, control_tx: None, steer_tx: None, + permission_tx: None, successful_steer_deliveries: HashSet::new(), }, ); @@ -11061,6 +11162,7 @@ mod error_outcome_emission_tests { recoverable_batch: None, control_tx: None, steer_tx: None, + permission_tx: None, successful_steer_deliveries: HashSet::new(), }, ); @@ -11222,6 +11324,7 @@ mod error_outcome_emission_tests { recoverable_batch: None, control_tx: None, steer_tx: None, + permission_tx: None, successful_steer_deliveries: HashSet::new(), }, ); diff --git a/crates/buzz-acp/src/permission_ledger.rs b/crates/buzz-acp/src/permission_ledger.rs new file mode 100644 index 00000000000..7b5cddd67bc --- /dev/null +++ b/crates/buzz-acp/src/permission_ledger.rs @@ -0,0 +1,603 @@ +//! Durable, host-owned lifecycle state for ACP permission requests. +//! +//! The relay observer is a best-effort presentation stream. It cannot be the +//! source of truth for a request that can cause an adapter action. This module +//! writes the lifecycle to the harness configuration directory before a +//! permission is shown to an owner, and every write is atomically replaced and +//! synced before the caller may advance the ACP protocol. + +use std::{ + collections::BTreeMap, + fs, + io::Write, + path::{Path, PathBuf}, + sync::{Arc, Mutex, OnceLock}, +}; + +use serde::{Deserialize, Serialize}; + +const LEDGER_VERSION: u8 = 1; +const MAX_RECORDS: usize = 256; + +#[derive(Clone, Debug, Deserialize, Serialize, PartialEq)] +#[serde(rename_all = "camelCase")] +pub struct PermissionRecord { + pub key: String, + pub channel_id: Option, + pub start_nonce: String, + pub turn_id: String, + pub session_id: String, + pub request_id: serde_json::Value, + pub action_digest: String, + pub request: serde_json::Value, + pub options: Vec, + pub expires_at: String, + pub updated_at: String, + pub state: PermissionState, +} + +#[derive(Clone, Debug, Deserialize, Serialize, PartialEq)] +#[serde(rename_all = "snake_case", tag = "kind", content = "reason")] +pub enum PermissionState { + Pending, + DecisionConsumed(String), + DeliveryAttempted(String), + Selected(String), + Cancelled, + Expired, + Abandoned, + DeliveryUnknown, +} + +#[derive(Default, Deserialize, Serialize)] +struct LedgerFile { + version: u8, + records: BTreeMap, +} + +/// A process-shared ledger. The enclosing mutex is deliberately held by the +/// caller only for a small atomic file transaction; no observer or relay work +/// runs while the ledger lock is held. +pub struct PermissionLedger { + path: PathBuf, + start_nonce: String, + records: BTreeMap, + poisoned: bool, + #[cfg(test)] + fail_after_replace: bool, +} + +static SHARED_LEDGER: OnceLock>> = OnceLock::new(); +static SHARED_LEDGER_INIT: Mutex<()> = Mutex::new(()); + +struct PersistError { + replaced: bool, + message: String, +} + +impl PermissionLedger { + /// Return the one ledger shared by every ACP pool client in this harness + /// generation. A second path or generation is rejected instead of opening + /// another in-memory view of the same action authority. + pub fn shared(path: PathBuf, start_nonce: String) -> Result>, String> { + // `OnceLock::set` alone does not serialize `open()`: two first callers + // could otherwise both reconcile the on-disk file before one wins. + let _initializing = SHARED_LEDGER_INIT + .lock() + .map_err(|_| "permission lifecycle initialization lock poisoned".to_string())?; + if let Some(existing) = SHARED_LEDGER.get() { + let ledger = existing + .lock() + .map_err(|_| "permission lifecycle ledger lock poisoned".to_string())?; + if ledger.path != path || ledger.start_nonce != start_nonce { + return Err( + "permission lifecycle ledger already belongs to another runtime generation" + .into(), + ); + } + return Ok(Arc::clone(existing)); + } + let opened = Arc::new(Mutex::new(Self::open(path, start_nonce)?)); + SHARED_LEDGER + .set(Arc::clone(&opened)) + .map_err(|_| "permission lifecycle initialization race".to_string())?; + Ok(opened) + } + + pub fn open(path: PathBuf, start_nonce: String) -> Result { + if start_nonce.is_empty() { + return Err("permission lifecycle ledger requires a managed-agent start nonce".into()); + } + let mut ledger = Self { + path, + start_nonce, + records: BTreeMap::new(), + poisoned: false, + #[cfg(test)] + fail_after_replace: false, + }; + if ledger.path.exists() { + let bytes = fs::read(&ledger.path) + .map_err(|error| format!("read permission lifecycle ledger: {error}"))?; + let file: LedgerFile = serde_json::from_slice(&bytes) + .map_err(|error| format!("parse permission lifecycle ledger: {error}"))?; + if file.version != LEDGER_VERSION { + return Err(format!( + "unsupported permission lifecycle ledger version {}", + file.version + )); + } + ledger.records = file.records; + } + // An ACP process cannot resume a stdout request after it has died. Do + // not reconstruct a pending UI card that would imply otherwise. A + // write attempt may have reached the adapter pipe before a crash, so it + // is explicitly reported as delivery-unknown rather than completed. + if ledger.records.values().any(|record| { + record.start_nonce != ledger.start_nonce + && matches!( + record.state, + PermissionState::Pending + | PermissionState::DecisionConsumed(_) + | PermissionState::DeliveryAttempted(_) + ) + }) { + for record in ledger.records.values_mut() { + if record.start_nonce == ledger.start_nonce { + continue; + } + record.state = match &record.state { + PermissionState::Pending | PermissionState::DecisionConsumed(_) => { + PermissionState::Abandoned + } + PermissionState::DeliveryAttempted(_) => PermissionState::DeliveryUnknown, + state => state.clone(), + }; + } + ledger.persist().map_err(|error| error.message)?; + } + Ok(ledger) + } + + pub fn record_pending(&mut self, record: PermissionRecord) -> Result<(), String> { + self.ensure_healthy()?; + if record.start_nonce != self.start_nonce { + return Err("permission lifecycle record has the wrong harness generation".into()); + } + if self.records.contains_key(&record.key) { + return Err("permission lifecycle key already exists".into()); + } + let previous = self.records.clone(); + self.records.insert(record.key.clone(), record); + self.prune_terminal_history(); + if let Err(error) = self.persist() { + self.handle_persist_error(previous, error)?; + } + Ok(()) + } + + pub fn transition_pending(&mut self, key: &str, state: PermissionState) -> Result<(), String> { + self.ensure_healthy()?; + let previous = self.records.clone(); + let record = self + .records + .get_mut(key) + .ok_or_else(|| "permission lifecycle record is missing".to_string())?; + if record.state != PermissionState::Pending { + return Err("permission lifecycle record was already consumed".into()); + } + record.state = state; + record.updated_at = chrono::Utc::now().to_rfc3339(); + if let Err(error) = self.persist() { + self.handle_persist_error(previous, error)?; + } + Ok(()) + } + + pub fn transition_consumed_to_delivery_attempt(&mut self, key: &str) -> Result<(), String> { + self.ensure_healthy()?; + let previous = self.records.clone(); + let record = self + .records + .get_mut(key) + .ok_or_else(|| "permission lifecycle record is missing".to_string())?; + let PermissionState::DecisionConsumed(option_id) = &record.state else { + return Err("permission lifecycle record is not a consumed decision".into()); + }; + record.state = PermissionState::DeliveryAttempted(option_id.clone()); + record.updated_at = chrono::Utc::now().to_rfc3339(); + if let Err(error) = self.persist() { + self.handle_persist_error(previous, error)?; + } + Ok(()) + } + + /// Record ambiguous delivery of a cancelled ACP response without reviving owner authority. + pub fn transition_cancelled_to_delivery_unknown(&mut self, key: &str) -> Result<(), String> { + self.ensure_healthy()?; + let previous = self.records.clone(); + let record = self + .records + .get_mut(key) + .ok_or_else(|| "permission lifecycle record is missing".to_string())?; + if record.state != PermissionState::Cancelled { + return Err("permission lifecycle record is not cancelled".into()); + } + record.state = PermissionState::DeliveryUnknown; + record.updated_at = chrono::Utc::now().to_rfc3339(); + if let Err(error) = self.persist() { + self.handle_persist_error(previous, error)?; + } + Ok(()) + } + + pub fn transition_delivery_attempt_to_selected(&mut self, key: &str) -> Result<(), String> { + self.ensure_healthy()?; + let previous = self.records.clone(); + let record = self + .records + .get_mut(key) + .ok_or_else(|| "permission lifecycle record is missing".to_string())?; + let PermissionState::DeliveryAttempted(option_id) = &record.state else { + return Err("permission lifecycle record is not a delivery attempt".into()); + }; + record.state = PermissionState::Selected(option_id.clone()); + record.updated_at = chrono::Utc::now().to_rfc3339(); + if let Err(error) = self.persist() { + self.handle_persist_error(previous, error)?; + } + Ok(()) + } + + #[cfg(test)] + pub fn get(&self, key: &str) -> Option<&PermissionRecord> { + self.records.get(key) + } + + pub fn start_nonce(&self) -> &str { + &self.start_nonce + } + + fn prune_terminal_history(&mut self) { + while self.records.len() > MAX_RECORDS { + let Some(key) = self.records.iter().find_map(|(key, record)| { + (!matches!( + record.state, + PermissionState::Pending + | PermissionState::DecisionConsumed(_) + | PermissionState::DeliveryAttempted(_) + )) + .then(|| key.clone()) + }) else { + break; + }; + self.records.remove(&key); + } + } + + fn ensure_healthy(&self) -> Result<(), String> { + if self.poisoned { + Err( + "permission lifecycle ledger needs authoritative reopen after ambiguous write" + .into(), + ) + } else { + Ok(()) + } + } + + fn handle_persist_error( + &mut self, + previous: BTreeMap, + error: PersistError, + ) -> Result<(), String> { + if error.replaced { + // The replacement reached disk but its directory durability is + // unknown. Keep the new in-memory state, poison this authority, + // and stop ACP advancement until a new generation reopens it. + self.poisoned = true; + } else { + self.records = previous; + } + Err(error.message) + } + + fn persist(&self) -> Result<(), PersistError> { + let parent = self.path.parent().ok_or_else(|| PersistError { + replaced: false, + message: "permission lifecycle ledger path has no parent".into(), + })?; + fs::create_dir_all(parent).map_err(|error| PersistError { + replaced: false, + message: format!("create permission lifecycle directory: {error}"), + })?; + let body = serde_json::to_vec(&LedgerFile { + version: LEDGER_VERSION, + records: self.records.clone(), + }) + .map_err(|error| PersistError { + replaced: false, + message: format!("serialize permission lifecycle ledger: {error}"), + })?; + let temporary = self.path.with_extension("json.tmp"); + let mut file = fs::OpenOptions::new() + .create(true) + .truncate(true) + .write(true) + .open(&temporary) + .map_err(|error| PersistError { + replaced: false, + message: format!("open permission lifecycle temporary file: {error}"), + })?; + #[cfg(unix)] + { + use std::os::unix::fs::PermissionsExt; + file.set_permissions(fs::Permissions::from_mode(0o600)) + .map_err(|error| PersistError { + replaced: false, + message: format!("restrict permission lifecycle temporary file: {error}"), + })?; + } + file.write_all(&body).map_err(|error| PersistError { + replaced: false, + message: format!("write permission lifecycle temporary file: {error}"), + })?; + file.sync_all().map_err(|error| PersistError { + replaced: false, + message: format!("sync permission lifecycle temporary file: {error}"), + })?; + drop(file); + fs::rename(&temporary, &self.path).map_err(|error| PersistError { + replaced: false, + message: format!("replace permission lifecycle ledger: {error}"), + })?; + #[cfg(test)] + if self.fail_after_replace { + return Err(PersistError { + replaced: true, + message: "injected failure after permission lifecycle replacement".into(), + }); + } + sync_directory(parent).map_err(|message| PersistError { + replaced: true, + message, + })?; + Ok(()) + } +} + +fn sync_directory(path: &Path) -> Result<(), String> { + #[cfg(unix)] + { + fs::File::open(path) + .and_then(|directory| directory.sync_all()) + .map_err(|error| format!("sync permission lifecycle directory: {error}"))?; + } + #[cfg(not(unix))] + let _ = path; + Ok(()) +} + +pub fn record_key( + start_nonce: &str, + turn_id: &str, + session_id: &str, + request_id: &serde_json::Value, + action_digest: &str, +) -> String { + let typed_id = serde_json::to_string(request_id).unwrap_or_default(); + format!("{start_nonce}:{turn_id}:{session_id}:{typed_id}:{action_digest}") +} + +#[cfg(test)] +mod tests { + use super::*; + + fn record(key: &str) -> PermissionRecord { + PermissionRecord { + key: key.into(), + channel_id: Some("channel".into()), + start_nonce: "generation-a".into(), + turn_id: "turn".into(), + session_id: "session".into(), + request_id: serde_json::json!(7), + action_digest: "digest".into(), + request: serde_json::json!({"method": "session/request_permission"}), + options: vec![serde_json::json!({"optionId": "opaque-allow"})], + expires_at: "2030-01-01T00:00:00Z".into(), + updated_at: "2029-01-01T00:00:00Z".into(), + state: PermissionState::Pending, + } + } + + #[test] + fn restart_reconciles_unfinished_request_to_abandoned() { + let temp = std::env::temp_dir().join(format!("buzz-acp-ledger-{}", uuid::Uuid::new_v4())); + fs::create_dir_all(&temp).unwrap(); + let path = temp.join("permission-lifecycle.json"); + let mut ledger = PermissionLedger::open(path.clone(), "generation-a".into()).unwrap(); + ledger.record_pending(record("record")).unwrap(); + drop(ledger); + + let restarted = PermissionLedger::open(path, "generation-b".into()).unwrap(); + assert_eq!( + restarted.get("record").unwrap().state, + PermissionState::Abandoned + ); + if temp.exists() { + fs::remove_dir_all(temp).unwrap(); + } + } + + #[test] + fn terminal_transition_is_durable_and_single_use() { + let temp = std::env::temp_dir().join(format!("buzz-acp-ledger-{}", uuid::Uuid::new_v4())); + fs::create_dir_all(&temp).unwrap(); + let path = temp.join("permission-lifecycle.json"); + let mut ledger = PermissionLedger::open(path.clone(), "generation-a".into()).unwrap(); + ledger.record_pending(record("record")).unwrap(); + ledger + .transition_pending( + "record", + PermissionState::DecisionConsumed("opaque-allow".into()), + ) + .unwrap(); + assert!(ledger + .transition_pending("record", PermissionState::Cancelled) + .is_err()); + ledger + .transition_consumed_to_delivery_attempt("record") + .unwrap(); + ledger + .transition_delivery_attempt_to_selected("record") + .unwrap(); + drop(ledger); + + let restarted = PermissionLedger::open(path, "generation-b".into()).unwrap(); + assert_eq!( + restarted.get("record").unwrap().state, + PermissionState::Selected("opaque-allow".into()) + ); + fs::remove_dir_all(temp).unwrap(); + } + + #[test] + fn cancelled_response_delivery_unknown_is_durable_and_never_reopens() { + let temp = std::env::temp_dir().join(format!("buzz-acp-ledger-{}", uuid::Uuid::new_v4())); + fs::create_dir_all(&temp).unwrap(); + let path = temp.join("permission-lifecycle.json"); + let mut ledger = PermissionLedger::open(path.clone(), "generation-a".into()).unwrap(); + ledger.record_pending(record("record")).unwrap(); + assert!(ledger + .transition_cancelled_to_delivery_unknown("record") + .is_err()); + ledger + .transition_pending("record", PermissionState::Cancelled) + .unwrap(); + ledger + .transition_cancelled_to_delivery_unknown("record") + .unwrap(); + assert!(ledger + .transition_cancelled_to_delivery_unknown("record") + .is_err()); + assert!(ledger + .transition_pending("record", PermissionState::DecisionConsumed("allow".into())) + .is_err()); + drop(ledger); + + let restarted = PermissionLedger::open(path, "generation-b".into()).unwrap(); + assert_eq!( + restarted.get("record").unwrap().state, + PermissionState::DeliveryUnknown + ); + fs::remove_dir_all(temp).unwrap(); + } + + #[test] + fn restart_does_not_replay_consumed_or_uncertain_delivery() { + let temp = std::env::temp_dir().join(format!("buzz-acp-ledger-{}", uuid::Uuid::new_v4())); + fs::create_dir_all(&temp).unwrap(); + let consumed_path = temp.join("consumed.json"); + let mut consumed = + PermissionLedger::open(consumed_path.clone(), "generation-a".into()).unwrap(); + consumed.record_pending(record("consumed")).unwrap(); + consumed + .transition_pending( + "consumed", + PermissionState::DecisionConsumed("opaque-allow".into()), + ) + .unwrap(); + drop(consumed); + assert_eq!( + PermissionLedger::open(consumed_path, "generation-b".into()) + .unwrap() + .get("consumed") + .unwrap() + .state, + PermissionState::Abandoned + ); + + let attempted_path = temp.join("attempted.json"); + let mut attempted = + PermissionLedger::open(attempted_path.clone(), "generation-a".into()).unwrap(); + attempted.record_pending(record("attempted")).unwrap(); + attempted + .transition_pending( + "attempted", + PermissionState::DecisionConsumed("opaque-allow".into()), + ) + .unwrap(); + attempted + .transition_consumed_to_delivery_attempt("attempted") + .unwrap(); + drop(attempted); + assert_eq!( + PermissionLedger::open(attempted_path, "generation-b".into()) + .unwrap() + .get("attempted") + .unwrap() + .state, + PermissionState::DeliveryUnknown + ); + fs::remove_dir_all(temp).unwrap(); + } + + #[test] + fn shared_initialization_serializes_first_open() { + let temp = + std::env::temp_dir().join(format!("buzz-acp-ledger-shared-{}", uuid::Uuid::new_v4())); + let path = temp.join("permission-lifecycle.json"); + let barrier = Arc::new(std::sync::Barrier::new(3)); + let mut joins = Vec::new(); + for _ in 0..2 { + let path = path.clone(); + let barrier = Arc::clone(&barrier); + joins.push(std::thread::spawn(move || { + barrier.wait(); + PermissionLedger::shared(path, "shared-generation".into()).unwrap() + })); + } + barrier.wait(); + let first = joins.remove(0).join().unwrap(); + let second = joins.remove(0).join().unwrap(); + assert!(Arc::ptr_eq(&first, &second)); + if temp.exists() { + fs::remove_dir_all(temp).unwrap(); + } + } + + #[test] + fn pre_replace_rolls_back_but_post_replace_poison_stops_advancement() { + let temp = + std::env::temp_dir().join(format!("buzz-acp-ledger-failure-{}", uuid::Uuid::new_v4())); + fs::create_dir_all(&temp).unwrap(); + let mut pre_replace = PermissionLedger { + path: temp.clone(), + start_nonce: "generation-a".into(), + records: BTreeMap::new(), + poisoned: false, + fail_after_replace: false, + }; + assert!(pre_replace.record_pending(record("pre-replace")).is_err()); + assert!(pre_replace.get("pre-replace").is_none()); + + let path = temp.join("post-replace.json"); + let mut post_replace = PermissionLedger::open(path.clone(), "generation-a".into()).unwrap(); + post_replace.fail_after_replace = true; + assert!(post_replace.record_pending(record("post-replace")).is_err()); + assert!(post_replace.get("post-replace").is_some()); + assert!(post_replace + .transition_pending("post-replace", PermissionState::Cancelled) + .is_err()); + drop(post_replace); + assert_eq!( + PermissionLedger::open(path, "generation-b".into()) + .unwrap() + .get("post-replace") + .unwrap() + .state, + PermissionState::Abandoned + ); + fs::remove_dir_all(temp).unwrap(); + } +} diff --git a/crates/buzz-acp/src/pool.rs b/crates/buzz-acp/src/pool.rs index c188633bef7..a57685a6129 100644 --- a/crates/buzz-acp/src/pool.rs +++ b/crates/buzz-acp/src/pool.rs @@ -86,6 +86,9 @@ pub struct TaskMeta { /// tasks only — all prompt tasks install a steer channel regardless /// of the agent's name. pub steer_tx: Option>, + /// Signed owner decisions for a pending ACP permission request. Unlike the + /// one-shot turn control, this never cancels or replays the live turn. + pub permission_tx: Option>, /// Successful non-cancelling steers acknowledged while this task owned the /// live session. The session ID prevents a late ack from contaminating a /// replacement session after task return. @@ -2265,9 +2268,13 @@ pub async fn run_prompt_task( prompt_text: Option, ctx: Arc, result_tx: mpsc::UnboundedSender, - control_rx: Option>, + turn_controls: ( + Option>, + Option>, + ), turn_id: String, ) { + let (control_rx, permission_rx) = turn_controls; // Is this a channel prompt or a heartbeat? let source = match &batch { Some(b) => PromptSource::Channel(b.scope.clone()), @@ -2281,6 +2288,9 @@ pub async fn run_prompt_task( turn_id.clone(), turn_started_at.clone(), )); + if let Some(rx) = permission_rx { + agent.acp.install_permission_rx(rx); + } let triggering_event_ids: Vec = batch .as_ref() .map(|b| b.events.iter().map(|be| be.event.id.to_hex()).collect()) @@ -6934,7 +6944,7 @@ done"# Some(format!("heartbeat-{turn}")), Arc::clone(&ctx), result_tx.clone(), - None, + (None, None), format!("turn-{turn}"), ) .await; @@ -7056,7 +7066,7 @@ done"# None, Arc::clone(&ctx), result_tx.clone(), - None, + (None, None), format!("turn-{turn}"), ) .await; @@ -7267,7 +7277,7 @@ done"# None, Arc::new(make_context(10)), result_tx.clone(), - None, + (None, None), "first-turn".into(), ) .await; @@ -7286,7 +7296,7 @@ done"# None, Arc::new(make_context(10)), result_tx, - None, + (None, None), "follow-up-turn".into(), ) .await; @@ -7453,7 +7463,7 @@ done"# None, Arc::clone(&ctx), result_tx.clone(), - None, + (None, None), turn_id.into(), ) .await; @@ -7616,7 +7626,7 @@ printf '%s\n' '{{"jsonrpc":"2.0","id":0,"result":{{"stopReason":"end_turn"}}}}'" None, Arc::new(ctx), result_tx, - None, + (None, None), "next-turn".into(), ) .await; @@ -8335,6 +8345,7 @@ printf '%s\n' '{{"jsonrpc":"2.0","id":0,"result":{{"stopReason":"end_turn"}}}}'" recoverable_batch: None, control_tx: None, steer_tx: None, + permission_tx: None, successful_steer_deliveries: HashSet::new(), }, ); @@ -10643,7 +10654,7 @@ done"# None, Arc::new(ctx), result_tx, - None, + (None, None), "indeterminate-project-turn".into(), ) .await; diff --git a/desktop/src-tauri/src/commands/mod.rs b/desktop/src-tauri/src/commands/mod.rs index c8184a01031..4ee9018cc55 100644 --- a/desktop/src-tauri/src/commands/mod.rs +++ b/desktop/src-tauri/src/commands/mod.rs @@ -48,6 +48,7 @@ mod notifications; mod observer_archive; mod os_idle; pub mod pairing; +mod permission_lifecycle; mod personas; mod prevent_sleep; mod profile; @@ -111,6 +112,7 @@ pub use notifications::*; pub use observer_archive::*; pub use os_idle::*; pub use pairing::*; +pub use permission_lifecycle::*; pub use personas::*; pub use prevent_sleep::*; pub use profile::*; diff --git a/desktop/src-tauri/src/commands/permission_lifecycle.rs b/desktop/src-tauri/src/commands/permission_lifecycle.rs new file mode 100644 index 00000000000..7460b6a3871 --- /dev/null +++ b/desktop/src-tauri/src/commands/permission_lifecycle.rs @@ -0,0 +1,37 @@ +//! Desktop-owned recovery reader for managed ACP permission decisions. + +use serde::Serialize; +use tauri::AppHandle; + +#[derive(Serialize)] +#[serde(rename_all = "camelCase")] +pub struct PermissionLifecycleSnapshot { + records: Vec, +} + +#[tauri::command] +pub fn get_managed_agent_permission_lifecycle( + app: AppHandle, + pubkey: String, + relay_url: String, +) -> Result { + let key = crate::managed_agents::ManagedAgentRuntimeKey::new(pubkey, &relay_url)?; + let path = crate::managed_agents::managed_agent_permission_ledger_path(&app, &key)?; + let bytes = match std::fs::read(path) { + Ok(bytes) => bytes, + Err(error) if error.kind() == std::io::ErrorKind::NotFound => { + return Ok(PermissionLifecycleSnapshot { + records: Vec::new(), + }); + } + Err(error) => return Err(format!("read permission lifecycle ledger: {error}")), + }; + let file: serde_json::Value = serde_json::from_slice(&bytes) + .map_err(|error| format!("parse permission lifecycle ledger: {error}"))?; + let records = file + .get("records") + .and_then(serde_json::Value::as_object) + .map(|records| records.values().cloned().collect()) + .unwrap_or_default(); + Ok(PermissionLifecycleSnapshot { records }) +} diff --git a/desktop/src-tauri/src/lib.rs b/desktop/src-tauri/src/lib.rs index e8b767e7c56..aa8a39dc45b 100644 --- a/desktop/src-tauri/src/lib.rs +++ b/desktop/src-tauri/src/lib.rs @@ -712,6 +712,7 @@ pub fn run() { set_managed_agent_auto_restart, delete_managed_agent, get_managed_agent_log, + get_managed_agent_permission_lifecycle, get_agent_models, discover_agent_models, agent_access_owner_only, diff --git a/desktop/src-tauri/src/managed_agents/reserved_env_keys.rs b/desktop/src-tauri/src/managed_agents/reserved_env_keys.rs index 07fcf8d2592..2f91b5b1a35 100644 --- a/desktop/src-tauri/src/managed_agents/reserved_env_keys.rs +++ b/desktop/src-tauri/src/managed_agents/reserved_env_keys.rs @@ -79,6 +79,10 @@ pub(crate) const RESERVED_ENV_KEYS: &[&str] = &[ // for same-session sweep decisions. "BUZZ_MANAGED_AGENT", "BUZZ_MANAGED_AGENT_START_NONCE", + // Desktop chooses the build- and runtime-scoped permission authority path. + // A definition-provided replacement could point the harness at another + // agent's decisions or an unprotected directory. + "BUZZ_ACP_PERMISSION_LEDGER_PATH", ]; pub(crate) fn is_reserved_env_key(key: &str) -> bool { diff --git a/desktop/src-tauri/src/managed_agents/runtime.rs b/desktop/src-tauri/src/managed_agents/runtime.rs index 5b44f95de92..0c39d76f434 100644 --- a/desktop/src-tauri/src/managed_agents/runtime.rs +++ b/desktop/src-tauri/src/managed_agents/runtime.rs @@ -792,9 +792,11 @@ pub fn spawn_agent_child( // Stamp desktop ownership and an unpredictable harness-generation identity. let start_nonce = uuid::Uuid::new_v4().simple().to_string(); + let permission_ledger_path = super::managed_agent_permission_ledger_path(app, &runtime_key)?; command .env("BUZZ_MANAGED_AGENT", current_instance_id(app)) - .env("BUZZ_MANAGED_AGENT_START_NONCE", &start_nonce); + .env("BUZZ_MANAGED_AGENT_START_NONCE", &start_nonce) + .env("BUZZ_ACP_PERMISSION_LEDGER_PATH", permission_ledger_path); // Stamp spawn config from values above, BEFORE spawning — a post-spawn // re-resolve races config edits and would stamp the wrong values. diff --git a/desktop/src-tauri/src/managed_agents/storage.rs b/desktop/src-tauri/src/managed_agents/storage.rs index f8a2c1039a8..049950cd120 100644 --- a/desktop/src-tauri/src/managed_agents/storage.rs +++ b/desktop/src-tauri/src/managed_agents/storage.rs @@ -95,6 +95,19 @@ pub fn managed_agent_runtime_log_path( Ok(managed_agents_logs_dir(app)?.join(format!("{}.log", key.runtime_id()))) } +/// Owner-only permission lifecycle ledger for exactly one managed runtime. +/// The stable filename is pair-scoped and lives under the build-scoped app data +/// root, so demo and production builds cannot read each other's authority. +pub fn managed_agent_permission_ledger_path( + app: &AppHandle, + key: &ManagedAgentRuntimeKey, +) -> Result { + let dir = managed_agents_base_dir(app)?.join("permission-lifecycle"); + fs::create_dir_all(&dir) + .map_err(|error| format!("failed to create permission lifecycle directory: {error}"))?; + Ok(dir.join(format!("{}.json", key.runtime_id()))) +} + /// Log path to surface for an agent whose runtime is not tracked in memory: /// the most recently written of its pair-scoped logs, falling back to the /// legacy single-runtime path when the agent has not run since harnesses diff --git a/desktop/src/features/agents/ui/ManagedAgentSessionPanel.tsx b/desktop/src/features/agents/ui/ManagedAgentSessionPanel.tsx index 4e526692246..9bc98d411df 100644 --- a/desktop/src/features/agents/ui/ManagedAgentSessionPanel.tsx +++ b/desktop/src/features/agents/ui/ManagedAgentSessionPanel.tsx @@ -36,12 +36,17 @@ import { useObserverEvents, useArchivedChannelEvents, } from "./useObserverEvents"; -import { buildTranscriptState } from "./agentSessionTranscript"; +import { + buildTranscriptState, + projectPermissionLedgerEvents, +} from "./agentSessionTranscript"; +import { getManagedAgentPermissionLifecycle } from "@/shared/api/tauriManagedAgents"; type ManagedAgentSessionPanelProps = { agent: Pick & { status: ManagedAgent["status"] | "unknown"; avatarUrl?: string | null; + relayUrl?: string | null; }; autoTail?: boolean; channelId?: string | null; @@ -85,6 +90,36 @@ export function ManagedAgentSessionPanel({ hasObserver, agent.pubkey, ); + const [ledgerRecords, setLedgerRecords] = React.useState([]); + React.useEffect(() => { + // A live runtime may reconnect its observer stream. Wait for that stream + // to be open, then refresh the authoritative Desktop ledger; idle + // runtimes still read it immediately for historical recovery. + if (!agent.relayUrl || (hasObserver && connectionState !== "open")) return; + let current = true; + void getManagedAgentPermissionLifecycle(agent.pubkey, agent.relayUrl) + .then((snapshot) => { + if (current) setLedgerRecords(snapshot.records); + }) + .catch(() => { + if (current) setLedgerRecords([]); + }); + return () => { + current = false; + }; + }, [agent.pubkey, agent.relayUrl, connectionState, hasObserver]); + + const ledgerEvents = React.useMemo( + () => + scopeByChannel( + projectPermissionLedgerEvents( + ledgerRecords, + hasObserver && connectionState === "open", + ), + channelId, + ), + [channelId, connectionState, hasObserver, ledgerRecords], + ); // Channel-scoped live events (capped at MAX_OBSERVER_EVENTS) and uncapped // archived events from SQLite paging. Both are raw ObserverEvent[] — we merge @@ -105,8 +140,12 @@ export function ManagedAgentSessionPanel({ // sorted ascending. Used as the single source for both the transcript and the // raw event rail / header count. const combinedEvents = React.useMemo( - () => mergeObserverEventWindows(scopedLiveEvents, archivedChannelEvents), - [scopedLiveEvents, archivedChannelEvents], + () => + mergeObserverEventWindows( + [...scopedLiveEvents, ...ledgerEvents], + archivedChannelEvents, + ), + [scopedLiveEvents, archivedChannelEvents, ledgerEvents], ); // Derive transcript once from the combined raw window. When transcriptOverride diff --git a/desktop/src/features/agents/ui/activityRenderClasses/LifecycleActivity.permission.test.mjs b/desktop/src/features/agents/ui/activityRenderClasses/LifecycleActivity.permission.test.mjs new file mode 100644 index 00000000000..6cd45296ae2 --- /dev/null +++ b/desktop/src/features/agents/ui/activityRenderClasses/LifecycleActivity.permission.test.mjs @@ -0,0 +1,42 @@ +import assert from "node:assert/strict"; +import test from "node:test"; +import { createElement } from "react"; +import { renderToStaticMarkup } from "react-dom/server"; + +import { LifecycleActivity } from "./LifecycleActivity.tsx"; + +test("owner permission card offers only recorded ACP options and keeps native retry unavailable", () => { + const html = renderToStaticMarkup( + createElement(LifecycleActivity, { + agentAvatarUrl: null, + agentName: "Fixture agent", + agentPubkey: "agent-pubkey", + item: { + id: "permission:fixture", + type: "lifecycle", + renderClass: "permission", + title: "Permission requested", + text: "Fixture action\nOptions: Allow orchid, Deny fern", + timestamp: "2026-09-22T10:00:00Z", + channelId: "channel-fixture", + turnId: "turn-fixture", + sessionId: "session-fixture", + pendingResolution: { + turnId: "turn-fixture", + sessionId: "session-fixture", + requestId: "opaque-request-9", + actionDigest: "digest-fixture", + options: [ + { optionId: "allow-orchid-6", label: "Allow orchid" }, + { optionId: "reject-fern-3", label: "Deny fern" }, + ], + }, + }, + }), + ); + + assert.match(html, /Allow orchid/); + assert.match(html, /Deny fern/); + assert.match(html, /Auto-review retry unavailable for this adapter/); + assert.doesNotMatch(html, /approveGuardianDeniedAction/); +}); diff --git a/desktop/src/features/agents/ui/activityRenderClasses/LifecycleActivity.tsx b/desktop/src/features/agents/ui/activityRenderClasses/LifecycleActivity.tsx index 53a00e64638..b2f091a11d2 100644 --- a/desktop/src/features/agents/ui/activityRenderClasses/LifecycleActivity.tsx +++ b/desktop/src/features/agents/ui/activityRenderClasses/LifecycleActivity.tsx @@ -1,5 +1,7 @@ import { AlertCircle, CheckCircle2, ShieldCheck, XCircle } from "lucide-react"; +import * as React from "react"; +import { resolveManagedAgentPermission } from "@/shared/api/agentControl"; import { formatTranscriptTimestampTitle } from "../agentSessionUtils"; import { ActivityRow, ActivityRowLabel } from "./ActivityRow"; import { ToolActivity } from "./ToolActivity"; @@ -38,6 +40,9 @@ function permissionOutcomeTone(outcome: string): "approve" | "deny" | "cancel" { } export function LifecycleActivity(props: ActivityRenderClassItemProps) { + const [submitting, setSubmitting] = React.useState(null); + const [submitError, setSubmitError] = React.useState(null); + if (props.item.type === "tool") { return ; } @@ -54,6 +59,7 @@ export function LifecycleActivity(props: ActivityRenderClassItemProps) { if (isPermission) { const { requestLines, optionsLine } = splitPermissionText(props.item.text); const outcome = props.item.outcome; + const pendingResolution = props.item.pendingResolution; const tone = outcome ? permissionOutcomeTone(outcome) : null; return (
{optionsLine}
) : null} + {pendingResolution && !outcome ? ( +
+
+ Waiting for owner decision +
+
+ {pendingResolution.options.map((option) => ( + + ))} +
+ {submitError ? ( +
{submitError}
+ ) : null} +
+ Auto-review retry unavailable for this adapter. +
+
+ ) : null} {/* Row 3: decision — only when outcome is resolved */} {outcome && tone ? ( <> diff --git a/desktop/src/features/agents/ui/agentSessionTranscript.test.mjs b/desktop/src/features/agents/ui/agentSessionTranscript.test.mjs index c8cfd30088e..e7c9215b2db 100644 --- a/desktop/src/features/agents/ui/agentSessionTranscript.test.mjs +++ b/desktop/src/features/agents/ui/agentSessionTranscript.test.mjs @@ -1,7 +1,10 @@ import assert from "node:assert/strict"; import test from "node:test"; -import { buildTranscript } from "./agentSessionTranscript.ts"; +import { + buildTranscript, + projectPermissionLedgerEvents, +} from "./agentSessionTranscript.ts"; import { buildTranscriptDisplayBlocks, flattenDisplayBlocks, @@ -640,6 +643,206 @@ test("buildTranscript surfaces session/request_permission as a permission lifecy assert.match(transcript[0].text, /Confirm force-with-lease push/); }); +test("buildTranscript shows bound adapter tool-call action and target before owner approval", () => { + const transcript = buildTranscript([ + { + ...baseEvent, + seq: 1, + kind: "acp_read", + payload: { + id: "adapter-shaped-8", + method: "session/request_permission", + params: { + sessionId: "sess-1", + toolCall: { + toolCallId: "call-adapter-8", + title: "Run shell command", + rawInput: { + command: "git push --force-with-lease origin review", + cwd: "/workspace/buzz", + }, + content: [{ type: "text", text: "Force push requires approval." }], + }, + options: [{ optionId: "approve-orchid", name: "Approve once" }], + }, + }, + }, + { + ...baseEvent, + seq: 2, + kind: "permission_pending", + payload: { + requestId: "adapter-shaped-8", + sessionId: "sess-1", + actionDigest: "adapter-digest-8", + options: [{ optionId: "approve-orchid", name: "Approve once" }], + }, + }, + ]); + assert.equal(transcript.length, 1); + assert.match(transcript[0].text, /Run shell command/); + assert.match(transcript[0].text, /git push --force-with-lease origin review/); + assert.match(transcript[0].text, /Working directory: \/workspace\/buzz/); + assert.match(transcript[0].text, /Force push requires approval/); + assert.equal( + transcript[0].pendingResolution.options[0].optionId, + "approve-orchid", + ); +}); + +test("buildTranscript restores an owner decision card only for the exact archived permission request", () => { + const transcript = buildTranscript([ + { + ...baseEvent, + seq: 1, + kind: "acp_read", + payload: { + jsonrpc: "2.0", + id: "opaque-rpc-54", + method: "session/request_permission", + params: { + sessionId: "sess-1", + title: "Fixture request", + options: [{ optionId: "grant-orchid-8", name: "Grant orchid" }], + }, + }, + }, + { + ...baseEvent, + seq: 2, + kind: "permission_pending", + payload: { + requestId: "opaque-rpc-54", + sessionId: "sess-1", + actionDigest: "digest-54", + options: [{ optionId: "grant-orchid-8", name: "Grant orchid" }], + }, + }, + ]); + + assert.equal(transcript.length, 1); + assert.deepEqual(transcript[0].pendingResolution, { + turnId: "turn-1", + sessionId: "sess-1", + requestId: "opaque-rpc-54", + actionDigest: "digest-54", + options: [{ optionId: "grant-orchid-8", label: "Grant orchid" }], + }); +}); + +test("buildTranscript preserves sequential approvals and concurrent same-ID permissions as distinct cards", () => { + const request = (seq, sessionId, requestId, optionId) => ({ + ...baseEvent, + seq, + kind: "acp_read", + sessionId, + payload: { + id: requestId, + method: "session/request_permission", + params: { + sessionId, + options: [{ optionId, kind: optionId, name: optionId }], + }, + }, + }); + const pending = (seq, sessionId, requestId, optionId) => ({ + ...baseEvent, + seq, + kind: "permission_pending", + sessionId, + payload: { + requestId, + sessionId, + actionDigest: `digest-${sessionId}-${requestId}`, + options: [{ optionId, name: optionId }], + }, + }); + const response = (seq, sessionId, requestId, optionId) => ({ + ...baseEvent, + seq, + kind: "acp_write", + sessionId, + payload: { + id: requestId, + result: { outcome: { outcome: "selected", optionId } }, + }, + }); + const transcript = buildTranscript([ + request(1, "session-one", 11, "allow-one"), + pending(2, "session-one", 11, "allow-one"), + response(3, "session-one", 11, "allow-one"), + request(4, "session-one", 12, "reject-two"), + pending(5, "session-one", 12, "reject-two"), + request(6, "session-two", 12, "allow-three"), + pending(7, "session-two", 12, "allow-three"), + ]); + const permissions = transcript.filter( + (item) => item.renderClass === "permission", + ); + assert.equal(permissions.length, 3); + assert.equal( + permissions.filter((item) => item.outcome === "Approved (allow-one)") + .length, + 1, + ); + assert.deepEqual( + permissions + .filter((item) => item.pendingResolution) + .map((item) => [ + item.pendingResolution.sessionId, + item.pendingResolution.requestId, + ]) + .sort(), + [ + ["session-one", 12], + ["session-two", 12], + ], + ); +}); + +test("buildTranscript turns an abandoned archived permission into a terminal state without a synthetic resume", () => { + const transcript = buildTranscript([ + { + ...baseEvent, + seq: 1, + kind: "acp_read", + payload: { + id: "request-ended", + method: "session/request_permission", + params: { + sessionId: "sess-1", + options: [{ optionId: "allow-ended", name: "Allow" }], + }, + }, + }, + { + ...baseEvent, + seq: 2, + kind: "permission_pending", + payload: { + requestId: "request-ended", + sessionId: "sess-1", + actionDigest: "digest-ended", + options: [{ optionId: "allow-ended", name: "Allow" }], + }, + }, + { + ...baseEvent, + seq: 3, + kind: "permission_abandoned", + payload: { + requestId: "request-ended", + sessionId: "sess-1", + actionDigest: "digest-ended", + outcome: "agent_session_ended", + }, + }, + ]); + assert.equal(transcript.length, 1); + assert.equal(transcript[0].outcome, "Unavailable (agent session ended)"); + assert.equal(transcript[0].pendingResolution, undefined); +}); + test("buildTranscript stamps completedAt when a terminal tool update is inserted first", () => { const transcript = buildTranscript([ { @@ -898,6 +1101,185 @@ test("buildTranscript appends Cancelled outcome for a numeric JSON-RPC id (cance assert.doesNotMatch(item.text ?? "", /Cancelled/); }); +function makePermissionLedgerEvent({ + seq = -1, + channelId = "channel-1", + sessionId = "session-1", + turnId = "turn-1", + requestId = "req-17", + state = { kind: "pending" }, + timestamp = "2026-06-30T09:59:59.000Z", + options = [ + { optionId: "allow-17", kind: "allow_once", name: "Approve once" }, + { optionId: "deny-17", kind: "reject_once", name: "Deny" }, + ], +} = {}) { + return { + seq, + timestamp, + kind: "permission_ledger", + agentIndex: null, + channelId, + sessionId, + turnId, + payload: { + channelId, + sessionId, + turnId, + requestId, + actionDigest: "digest-17", + request: { + method: "session/request_permission", + params: { + sessionId, + toolCall: { + title: "Run protected command", + rawInput: { command: "git push", cwd: "/workspace/repo" }, + content: [{ type: "text", text: "Push the reviewed branch" }], + }, + options, + }, + }, + options, + state, + }, + }; +} + +test("buildTranscript merges ledger and live permission state into one terminal card", () => { + const pendingSnapshot = makePermissionLedgerEvent(); + const liveRequest = makePermissionRequest(1, "req-17"); + liveRequest.payload.params.sessionId = "session-1"; + liveRequest.payload.params.options = [ + { optionId: "allow-17", kind: "allow_once", name: "Approve once" }, + { optionId: "deny-17", kind: "reject_once", name: "Deny" }, + ]; + const selected = makePermissionResponse(2, "req-17", "selected", "allow-17"); + const stalePendingSnapshot = makePermissionLedgerEvent({ + seq: -2, + timestamp: "2026-06-30T10:00:02.000Z", + }); + + const transcript = buildTranscript([ + pendingSnapshot, + liveRequest, + selected, + stalePendingSnapshot, + ]); + const permissions = transcript.filter( + (item) => item.renderClass === "permission", + ); + + assert.equal(permissions.length, 1, "one canonical request card remains"); + assert.equal(permissions[0].outcome, "Approved (allow_once)"); + assert.equal( + permissions[0].pendingResolution, + undefined, + "a retained pending snapshot cannot restore controls after selection", + ); + assert.match(permissions[0].text, /Command: git push/); + assert.match(permissions[0].text, /Working directory: \/workspace\/repo/); +}); + +test("SDK cancellation makes the live card inert and persists an inert reconnect card", () => { + const pendingSnapshot = makePermissionLedgerEvent(); + const liveRequest = makePermissionRequest(1, "req-17"); + liveRequest.payload.params.sessionId = "session-1"; + const cancelledResponse = makePermissionResponse(2, "req-17", "cancelled"); + const stalePendingSnapshot = makePermissionLedgerEvent({ seq: -2 }); + const [liveCard] = buildTranscript([ + pendingSnapshot, + liveRequest, + cancelledResponse, + stalePendingSnapshot, + ]).filter((item) => item.renderClass === "permission"); + assert.equal(liveCard.outcome, "Cancelled"); + assert.equal(liveCard.pendingResolution, undefined); + + const [reconnectCard] = buildTranscript([ + makePermissionLedgerEvent({ state: { kind: "cancelled" } }), + ]).filter((item) => item.renderClass === "permission"); + assert.equal(reconnectCard.outcome, "Cancelled"); + assert.equal(reconnectCard.pendingResolution, undefined); +}); + +test("failed cancellation response delivery makes the live card inert", () => { + const pendingSnapshot = makePermissionLedgerEvent(); + const deliveryUnknown = { + ...pendingSnapshot, + seq: 2, + kind: "permission_delivery_unknown", + payload: { requestId: "req-17", sessionId: "session-1" }, + }; + const [card] = buildTranscript([pendingSnapshot, deliveryUnknown]).filter( + (item) => item.renderClass === "permission", + ); + assert.equal(card.outcome, "Delivery unknown (not replayed)"); + assert.equal(card.pendingResolution, undefined); +}); + +test("buildTranscript keeps same request ids in different channels separate", () => { + const channelOne = makePermissionLedgerEvent({ channelId: "channel-1" }); + const channelTwo = makePermissionLedgerEvent({ channelId: "channel-2" }); + const transcript = buildTranscript([channelOne, channelTwo]); + const permissions = transcript.filter( + (item) => item.renderClass === "permission", + ); + + assert.equal(permissions.length, 2); + assert.deepEqual(permissions.map((item) => item.channelId).sort(), [ + "channel-1", + "channel-2", + ]); +}); + +test("buildTranscript preserves typed JSON-RPC ids in reconnect-only owner cards", () => { + for (const requestId of [17, "request-17"]) { + const [permission] = buildTranscript([ + makePermissionLedgerEvent({ requestId }), + ]).filter((item) => item.renderClass === "permission"); + + assert.equal( + permission.pendingResolution?.requestId, + requestId, + "the owner broker receives the original JSON-RPC value, not the map key", + ); + } +}); + +test("buildTranscript labels an opaque reconnect deny selection as denied", () => { + const [permission] = buildTranscript([ + makePermissionLedgerEvent({ + state: { kind: "selected", reason: "opaque-choice-9" }, + options: [ + { optionId: "opaque-choice-9", kind: "reject_once", name: "Continue" }, + ], + }), + ]).filter((item) => item.renderClass === "permission"); + + assert.equal(permission.outcome, "Denied (reject_once)"); + assert.equal(permission.pendingResolution, undefined); +}); + +test("permission ledger is actionable only for a live reconnect, never a dead runtime", () => { + const record = makePermissionLedgerEvent().payload; + const [livePermission] = buildTranscript( + projectPermissionLedgerEvents([record], true), + ).filter((item) => item.renderClass === "permission"); + const [deadPermission] = buildTranscript( + projectPermissionLedgerEvents([record], false), + ).filter((item) => item.renderClass === "permission"); + + assert.ok(livePermission.pendingResolution, "open runtime restores its card"); + assert.equal(deadPermission.pendingResolution, undefined); + assert.equal(deadPermission.outcome, "Unavailable (agent session ended)"); + assert.equal( + record.state.kind, + "pending", + "Desktop projection does not rewrite the harness-owned ledger record", + ); +}); + test('buildTranscript does not collide between numeric id 1 and string id "1"', () => { // JSON.stringify produces "1" for number 1 and "\"1\"" for string "1", // so requests with these two different id types must NOT cross-attach. diff --git a/desktop/src/features/agents/ui/agentSessionTranscript.ts b/desktop/src/features/agents/ui/agentSessionTranscript.ts index ffc7f522037..fe5f168d044 100644 --- a/desktop/src/features/agents/ui/agentSessionTranscript.ts +++ b/desktop/src/features/agents/ui/agentSessionTranscript.ts @@ -28,75 +28,28 @@ import { parseSystemPromptSections, } from "./agentSessionTranscriptHelpers"; import { friendlyTurnErrorCopy } from "../lib/friendlyAgentLastError"; +import { + describePermissionOutcome, + describePermissionRequest, + jsonRpcId, + ownerResolutionFromPending, + permissionIdentity, + permissionTerminalObserverOutcomes, + rawJsonRpcId, +} from "./agentSessionTranscriptPermissions"; +import { + createEmptyTranscriptState, + createTranscriptDraft, + type TranscriptDraft, + type TranscriptState, +} from "./agentSessionTranscriptState"; export { describeRawEvent } from "./agentSessionTranscriptHelpers"; - -export type TranscriptState = { - items: TranscriptItem[]; - itemsById: Map; - activeMessageKey: Map; - sealedKeys: Set; - triggeringEventIdsByTurn: Map; - /** - * Maps JSON-RPC request id → { itemId, optionNames }. - * Populated when a `session/request_permission` request is ingested so the - * matching response (which carries the same JSON-RPC id, no `method`) can - * correlate and append the outcome to the lifecycle item. - */ - pendingPermissions: Map< - string, - { itemId: string; optionNames: Map } - >; - continuationSeq: number; - latestSessionId: string | null; -}; - -export function createEmptyTranscriptState(): TranscriptState { - return { - items: [], - itemsById: new Map(), - activeMessageKey: new Map(), - sealedKeys: new Set(), - triggeringEventIdsByTurn: new Map(), - pendingPermissions: new Map(), - continuationSeq: 0, - latestSessionId: null, - }; -} - -/** - * Mutable draft that collects changes during a single processTranscriptEvent - * call. Replaces the previous pattern of nested closures capturing bare `let` - * bindings — all mutation now targets this explicit object. - */ -type TranscriptDraft = { - items: TranscriptItem[]; - itemsById: Map; - activeMessageKey: Map; - sealedKeys: Set; - triggeringEventIdsByTurn: Map; - pendingPermissions: Map< - string, - { itemId: string; optionNames: Map } - >; - continuationSeq: number; - latestSessionId: string | null; - changed: boolean; -}; - -function draftFrom(state: TranscriptState): TranscriptDraft { - return { - items: state.items, - itemsById: state.itemsById, - activeMessageKey: state.activeMessageKey, - sealedKeys: state.sealedKeys, - triggeringEventIdsByTurn: state.triggeringEventIdsByTurn, - pendingPermissions: state.pendingPermissions, - continuationSeq: state.continuationSeq, - latestSessionId: state.latestSessionId, - changed: false, - }; -} +export { + createEmptyTranscriptState, + type TranscriptState, +} from "./agentSessionTranscriptState"; +export { projectPermissionLedgerEvents } from "./agentSessionTranscriptPermissions"; /** Lazily copy items + itemsById on first mutation so callers get new refs. */ function ensureMutable(d: TranscriptDraft) { @@ -171,98 +124,6 @@ function stringifyPayload(value: unknown) { } } -function describePermissionRequest(payload: Record) { - const params = asRecord(payload.params); - const title = - asString(params.title) ?? - asString(params.message) ?? - asString(params.reason) ?? - "Permission requested"; - const toolCallId = - asString(params.toolCallId) ?? asString(params.tool_call_id); - const options = Array.isArray(params.options) - ? params.options - .map((option) => { - const record = asRecord(option); - return ( - asString(record.name) ?? - asString(record.kind) ?? - asString(record.optionId) - ); - }) - .filter((option): option is string => Boolean(option)) - : []; - const detail: string[] = []; - if (title !== "Permission requested") detail.push(title); - if (toolCallId) detail.push(`Tool call: ${toolCallId}`); - if (options.length > 0) detail.push(`Options: ${options.join(", ")}`); - - // Build optionId → kind map for outcome labeling on the response. - const optionNames = new Map(); - if (Array.isArray(params.options)) { - for (const option of params.options) { - const record = asRecord(option); - const optionId = asString(record.optionId); - const kind = asString(record.kind); - if (optionId && kind) { - optionNames.set(optionId, kind); - } - } - } - - return { - title, - text: detail.join("\n"), - optionNames, - descriptor: { - renderClass: "permission" as const, - label: "Permission requested", - preview: title, - action: { verb: "Requested", object: title }, - tone: "admin" as const, - operation: "session/request_permission", - object: title, - source: "acp" as const, - groupKey: "permission:request", - }, - }; -} - -/** - * Format a human-readable outcome label from a permission response. - * kind values from ACP: allow_once, allow_always, reject_once, reject_always. - * "reject_*" kinds are denials; anything else that is selected is an approval. - */ -function describePermissionOutcome( - outcome: string, - optionId: string | null, - optionNames: Map, -): string { - if (outcome === "cancelled") { - return "Cancelled"; - } - if (outcome === "selected" && optionId) { - const kind = optionNames.get(optionId) ?? optionId; - const isDenial = kind.startsWith("reject"); - const verb = isDenial ? "Denied" : "Approved"; - return `${verb} (${kind})`; - } - return outcome; -} - -/** - * Stable map key for a JSON-RPC id, which may be a string or a finite number - * per the spec. Using JSON.stringify avoids collisions between the number 1 and - * the string "1". Returns null for null, undefined, or non-id values (objects, - * booleans) so callers can gate on presence without a separate type check. - */ -function jsonRpcId(value: unknown): string | null { - if (typeof value === "string") return JSON.stringify(value); - if (typeof value === "number" && Number.isFinite(value)) - return JSON.stringify(value); - return null; -} - function describeFreeformStatus(payload: Record) { const statusType = asString(payload.type) ?? asString(payload.status); const title = @@ -701,7 +562,7 @@ export function processTranscriptEvent( state: TranscriptState, event: ObserverEvent, ): TranscriptState { - const d = draftFrom(state); + const d = createTranscriptDraft(state); if (event.sessionId && event.sessionId !== d.latestSessionId) { d.latestSessionId = event.sessionId; @@ -786,13 +647,162 @@ export function processTranscriptEvent( ctx, event.kind, ); + } else if (event.kind === "permission_pending") { + const payload = asRecord(event.payload); + const requestId = jsonRpcId(payload.requestId); + const resolution = ownerResolutionFromPending(payload, event); + const pending = + requestId && resolution + ? d.pendingPermissions.get( + permissionIdentity( + channelId, + resolution.sessionId, + resolution.turnId, + requestId, + ), + ) + : null; + const existing = pending ? d.itemsById.get(pending.itemId) : null; + // A request is visible before this lifecycle marker. Keep the card inert + // unless the exact request, session, and turn all agree. + if ( + pending && + resolution && + existing?.type === "lifecycle" && + existing.turnId === resolution.turnId && + existing.sessionId === resolution.sessionId + ) { + replaceItem(d, pending.itemId, { + ...existing, + pendingResolution: resolution, + }); + } + } else if (event.kind === "permission_ledger") { + const payload = asRecord(event.payload); + const request = asRecord(payload.request); + const rawRequestId = rawJsonRpcId(payload.requestId); + const requestId = jsonRpcId(rawRequestId); + const sessionId = asString(payload.sessionId); + const turnId = asString(payload.turnId); + const digest = asString(payload.actionDigest); + const options = Array.isArray(payload.options) ? payload.options : []; + const lifecycleState = asRecord(payload.state); + const stateKind = asString(lifecycleState.kind); + if ( + Object.keys(request).length === 0 || + !requestId || + rawRequestId === null || + !sessionId || + !turnId || + !digest + ) + return state; + const description = describePermissionRequest(request); + const canonicalId = permissionIdentity( + channelId, + sessionId, + turnId, + requestId, + ); + const existing = d.itemsById.get(`permission:${canonicalId}`); + const existingPermission = existing?.type === "lifecycle" ? existing : null; + const pendingResolution = + stateKind === "pending" && !existingPermission?.outcome + ? { + turnId, + sessionId, + requestId: rawRequestId, + actionDigest: digest, + options: options.flatMap((option) => { + const value = asRecord(option); + const optionId = asString(value.optionId); + return optionId + ? [{ optionId, label: asString(value.name) ?? optionId }] + : []; + }), + } + : undefined; + const outcome = + stateKind === "selected" + ? describePermissionOutcome( + "selected", + asString(lifecycleState.reason) ?? null, + description.optionNames, + ) + : stateKind === "cancelled" || stateKind === "expired" + ? "Cancelled" + : stateKind === "delivery_unknown" + ? "Delivery unknown (not replayed)" + : stateKind === "abandoned" || + stateKind === "decision_consumed" || + stateKind === "delivery_attempted" + ? "Unavailable (agent session ended)" + : undefined; + const item = { + id: `permission:${canonicalId}`, + type: "lifecycle", + renderClass: "permission", + title: "Permission request", + text: description.text, + timestamp: event.timestamp, + channelId, + turnId, + sessionId, + pendingResolution, + outcome, + acpSource: event.kind, + } as const; + if (existing?.type === "lifecycle") { + // A stale fetched Pending row must never re-enable a terminal live card. + replaceItem( + d, + item.id, + existingPermission?.outcome && stateKind === "pending" + ? existingPermission + : item, + ); + } else { + pushItem(d, item); + } + if (pendingResolution) { + d.pendingPermissions = new Map(d.pendingPermissions); + d.pendingPermissions.set(canonicalId, { + itemId: item.id, + optionNames: description.optionNames, + }); + } + } else if (Object.hasOwn(permissionTerminalObserverOutcomes, event.kind)) { + const payload = asRecord(event.payload); + const requestId = jsonRpcId(payload.requestId); + const sessionId = asString(payload.sessionId); + const key = + requestId && sessionId && event.turnId + ? permissionIdentity(channelId, sessionId, event.turnId, requestId) + : null; + const pending = key ? d.pendingPermissions.get(key) : null; + const existing = pending ? d.itemsById.get(pending.itemId) : null; + if (key && pending && existing?.type === "lifecycle") { + replaceItem(d, pending.itemId, { + ...existing, + outcome: permissionTerminalObserverOutcomes[event.kind], + pendingResolution: undefined, + }); + d.pendingPermissions = new Map(d.pendingPermissions); + d.pendingPermissions.delete(key); + } } else if (event.kind === "acp_read" || event.kind === "acp_write") { const payload = asRecord(event.payload); const method = asString(payload.method); if (method === "session/request_permission") { const request = describePermissionRequest(payload); - const itemId = `permission:${ch}:${event.turnId ?? event.seq}`; + const requestId = jsonRpcId(payload.id); + const requestSessionId = + asString(asRecord(payload.params).sessionId) ?? ctx.sessionId; + const identity = requestId + ? permissionIdentity(channelId, requestSessionId, ctx.turnId, requestId) + : `${ch}:${requestSessionId ?? "unknown-session"}:${event.turnId ?? event.seq}:${event.seq}`; + const itemId = `permission:${identity}`; upsertLifecycleItem( d, itemId, @@ -806,21 +816,33 @@ export function processTranscriptEvent( ); // Index by JSON-RPC id so the response (acp_write with result.outcome, // no method) can correlate by id rather than by turn/seq. - const requestId = jsonRpcId(payload.id); if (requestId) { d.pendingPermissions = new Map(d.pendingPermissions); - d.pendingPermissions.set(requestId, { - itemId, - optionNames: request.optionNames, - }); + d.pendingPermissions.set( + permissionIdentity( + channelId, + requestSessionId, + ctx.turnId, + requestId, + ), + { + itemId, + optionNames: request.optionNames, + }, + ); } } else if (event.kind === "acp_write" && !method) { // Permission response: {"id": , "result": {"outcome": {...}}} const responseId = jsonRpcId(payload.id); const result = asRecord(asRecord(payload.result).outcome); const outcomeKind = asString(result.outcome); - const pending = responseId ? d.pendingPermissions.get(responseId) : null; - if (pending && outcomeKind && responseId) { + const responseKey = responseId + ? permissionIdentity(channelId, ctx.sessionId, ctx.turnId, responseId) + : null; + const pending = responseKey + ? d.pendingPermissions.get(responseKey) + : null; + if (pending && outcomeKind && responseKey) { const optionId = asString(result.optionId) ?? null; const outcomeText = describePermissionOutcome( outcomeKind, @@ -832,10 +854,11 @@ export function processTranscriptEvent( replaceItem(d, pending.itemId, { ...existing, outcome: outcomeText, + pendingResolution: undefined, }); // Remove from pending map — the outcome is now recorded. d.pendingPermissions = new Map(d.pendingPermissions); - d.pendingPermissions.delete(responseId); + d.pendingPermissions.delete(responseKey); } } } else if (event.kind === "acp_write" && method === "session/prompt") { diff --git a/desktop/src/features/agents/ui/agentSessionTranscriptPermissions.ts b/desktop/src/features/agents/ui/agentSessionTranscriptPermissions.ts new file mode 100644 index 00000000000..7e6fa629f8a --- /dev/null +++ b/desktop/src/features/agents/ui/agentSessionTranscriptPermissions.ts @@ -0,0 +1,213 @@ +import type { + ObserverEvent, + PendingPermissionResolution, +} from "./agentSessionTypes"; +import { asRecord, asString } from "./agentSessionUtils"; + +export const permissionTerminalObserverOutcomes: Record = { + permission_abandoned: "Unavailable (agent session ended)", + permission_delivery_unknown: "Delivery unknown (not replayed)", +}; + +export function describePermissionRequest(payload: Record) { + const params = asRecord(payload.params); + const toolCall = asRecord(params.toolCall); + const rawInput = asRecord(toolCall.rawInput); + const title = + asString(params.title) ?? + asString(params.message) ?? + asString(params.reason) ?? + asString(toolCall.title) ?? + "Permission requested"; + const toolCallId = + asString(params.toolCallId) ?? + asString(params.tool_call_id) ?? + asString(toolCall.toolCallId) ?? + asString(toolCall.id); + const command = asString(rawInput.command); + const cwd = asString(rawInput.cwd); + const toolText = Array.isArray(toolCall.content) + ? toolCall.content + .map((content) => asString(asRecord(content).text)) + .filter((text): text is string => Boolean(text)) + .join("\n") + : null; + const options = Array.isArray(params.options) + ? params.options + .map((option) => { + const record = asRecord(option); + return ( + asString(record.name) ?? + asString(record.kind) ?? + asString(record.optionId) + ); + }) + .filter((option): option is string => Boolean(option)) + : []; + const detail: string[] = []; + if (title !== "Permission requested") detail.push(title); + if (toolCallId) detail.push(`Tool call: ${toolCallId}`); + if (command) detail.push(`Command: ${command}`); + if (cwd) detail.push(`Working directory: ${cwd}`); + if (toolText) detail.push(toolText); + if (options.length > 0) detail.push(`Options: ${options.join(", ")}`); + + const optionNames = new Map(); + if (Array.isArray(params.options)) { + for (const option of params.options) { + const record = asRecord(option); + const optionId = asString(record.optionId); + const kind = asString(record.kind); + if (optionId && kind) { + optionNames.set(optionId, kind); + } + } + } + + return { + title, + text: detail.join("\n"), + optionNames, + descriptor: { + renderClass: "permission" as const, + label: "Permission requested", + preview: title, + action: { verb: "Requested", object: title }, + tone: "admin" as const, + operation: "session/request_permission", + object: title, + source: "acp" as const, + groupKey: "permission:request", + }, + }; +} + +/** + * Format a human-readable outcome label from a permission response. + * kind values from ACP: allow_once, allow_always, reject_once, reject_always. + * "reject_*" kinds are denials; anything else that is selected is an approval. + */ +export function describePermissionOutcome( + outcome: string, + optionId: string | null, + optionNames: Map, +): string { + if (outcome === "cancelled") { + return "Cancelled"; + } + if (outcome === "selected" && optionId) { + const kind = optionNames.get(optionId) ?? optionId; + const isDenial = kind.startsWith("reject"); + const verb = isDenial ? "Denied" : "Approved"; + return `${verb} (${kind})`; + } + return outcome; +} + +/** + * Stable map key for a JSON-RPC id, which may be a string or a finite number + * per the spec. Using JSON.stringify avoids collisions between the number 1 and + * the string "1". Returns null for null, undefined, or non-id values (objects, + * booleans) so callers can gate on presence without a separate type check. + */ +export function jsonRpcId(value: unknown): string | null { + const rawId = rawJsonRpcId(value); + return rawId === null ? null : JSON.stringify(rawId); +} + +/** Preserve the wire value for a later authenticated owner-resolution call. */ +export function rawJsonRpcId(value: unknown): string | number | null { + if (typeof value === "string") return value; + if (typeof value === "number" && Number.isFinite(value)) return value; + return null; +} + +/** + * Project durable records into the transcript only after the matching runtime + * has an open observer connection. A dead harness cannot resume the ACP read, + * so its persisted Pending row remains visible as an unavailable history row + * without Desktop mutating the single-writer ledger. + */ +export function projectPermissionLedgerEvents( + records: unknown[], + runtimeCanResolve: boolean, +): ObserverEvent[] { + return records.flatMap((record, index) => { + if (!record || typeof record !== "object") return []; + const value = record as Record; + const channelId = + typeof value.channelId === "string" ? value.channelId : null; + const sessionId = + typeof value.sessionId === "string" ? value.sessionId : null; + const turnId = typeof value.turnId === "string" ? value.turnId : null; + const timestamp = + typeof value.updatedAt === "string" + ? value.updatedAt + : new Date(0).toISOString(); + const state = asRecord(value.state); + const payload = + runtimeCanResolve || asString(state.kind) !== "pending" + ? value + : { ...value, state: { ...state, kind: "abandoned" } }; + return [ + { + seq: -1 - index, + timestamp, + kind: "permission_ledger", + agentIndex: null, + channelId, + sessionId, + turnId, + payload, + }, + ]; + }); +} + +/** + * ACP permits request IDs to repeat in different sessions. A channel can also + * hold multiple thread sessions, so no one of those values is a safe owner + * decision correlator on its own. Preserve the JSON type in `requestId` so + * numeric 1 and string "1" remain distinct. + */ +export function permissionIdentity( + channelId: string | null, + sessionId: string | null, + turnId: string | null, + requestId: string, +) { + return `${channelId ?? "global"}:${sessionId ?? "unknown-session"}:${turnId ?? "unknown-turn"}:${requestId}`; +} + +export function ownerResolutionFromPending( + payload: Record, + event: ObserverEvent, +): PendingPermissionResolution | null { + const requestId = payload.requestId; + const actionDigest = asString(payload.actionDigest); + const sessionId = asString(payload.sessionId); + const turnId = event.turnId; + if ( + (typeof requestId !== "string" && typeof requestId !== "number") || + !actionDigest || + !sessionId || + !turnId || + !Array.isArray(payload.options) + ) { + return null; + } + const options = payload.options.flatMap((value) => { + const option = asRecord(value); + const optionId = asString(option.optionId); + if (!optionId) return []; + return [ + { + optionId, + label: asString(option.name) ?? asString(option.kind) ?? optionId, + }, + ]; + }); + return options.length > 0 + ? { turnId, sessionId, requestId, actionDigest, options } + : null; +} diff --git a/desktop/src/features/agents/ui/agentSessionTranscriptState.ts b/desktop/src/features/agents/ui/agentSessionTranscriptState.ts new file mode 100644 index 00000000000..6ecbd0f6d97 --- /dev/null +++ b/desktop/src/features/agents/ui/agentSessionTranscriptState.ts @@ -0,0 +1,67 @@ +import type { TranscriptItem } from "./agentSessionTypes"; + +type PendingPermission = { + itemId: string; + optionNames: Map; +}; + +export type TranscriptState = { + items: TranscriptItem[]; + itemsById: Map; + activeMessageKey: Map; + sealedKeys: Set; + triggeringEventIdsByTurn: Map; + /** + * Maps JSON-RPC request id → { itemId, optionNames }. + * Populated when a `session/request_permission` request is ingested so the + * matching response (which carries the same JSON-RPC id, no `method`) can + * correlate and append the outcome to the lifecycle item. + */ + pendingPermissions: Map; + continuationSeq: number; + latestSessionId: string | null; +}; + +export function createEmptyTranscriptState(): TranscriptState { + return { + items: [], + itemsById: new Map(), + activeMessageKey: new Map(), + sealedKeys: new Set(), + triggeringEventIdsByTurn: new Map(), + pendingPermissions: new Map(), + continuationSeq: 0, + latestSessionId: null, + }; +} + +/** + * Mutable draft that collects changes during a single processTranscriptEvent + * call. Replaces the previous pattern of nested closures capturing bare `let` + * bindings — all mutation now targets this explicit object. + */ +export type TranscriptDraft = { + items: TranscriptItem[]; + itemsById: Map; + activeMessageKey: Map; + sealedKeys: Set; + triggeringEventIdsByTurn: Map; + pendingPermissions: Map; + continuationSeq: number; + latestSessionId: string | null; + changed: boolean; +}; + +export function createTranscriptDraft(state: TranscriptState): TranscriptDraft { + return { + items: state.items, + itemsById: state.itemsById, + activeMessageKey: state.activeMessageKey, + sealedKeys: state.sealedKeys, + triggeringEventIdsByTurn: state.triggeringEventIdsByTurn, + pendingPermissions: state.pendingPermissions, + continuationSeq: state.continuationSeq, + latestSessionId: state.latestSessionId, + changed: false, + }; +} diff --git a/desktop/src/features/agents/ui/agentSessionTypes.ts b/desktop/src/features/agents/ui/agentSessionTypes.ts index 578f98076cd..5dc32a721f3 100644 --- a/desktop/src/features/agents/ui/agentSessionTypes.ts +++ b/desktop/src/features/agents/ui/agentSessionTypes.ts @@ -68,6 +68,15 @@ export type TranscriptItemIdentity = { channelId?: string | null; }; +/** A signed-owner decision may only select one option from this ACP request. */ +export type PendingPermissionResolution = { + turnId: string; + sessionId: string; + requestId: string | number; + actionDigest: string; + options: Array<{ optionId: string; label: string }>; +}; + export type TranscriptItem = | ({ id: string; @@ -109,6 +118,8 @@ export type TranscriptItem = text: string; /** Resolved outcome for permission items (e.g. "Approved (allow_once)", "Denied (reject_once)", "Cancelled"). */ outcome?: string; + /** Present only after the harness has created an authenticated owner wait. */ + pendingResolution?: PendingPermissionResolution; timestamp: string; descriptor?: AgentActivityDescriptor; acpSource?: TranscriptAcpSource; diff --git a/desktop/src/shared/api/agentControl.ts b/desktop/src/shared/api/agentControl.ts index 67a6d377d28..9cc2ca4ce7a 100644 --- a/desktop/src/shared/api/agentControl.ts +++ b/desktop/src/shared/api/agentControl.ts @@ -36,3 +36,24 @@ export async function switchManagedAgentModel( requestId, }); } + +/** + * Resolve exactly one ACP permission request through the authenticated owner + * control path. The harness independently rechecks every binding before it + * writes the adapter's selected option response. + */ +export async function resolveManagedAgentPermission( + pubkey: string, + decision: { + turnId: string; + sessionId: string; + requestId: string | number; + actionDigest: string; + optionId: string; + }, +): Promise { + await sendAgentObserverControl(pubkey, { + type: "resolve_permission", + ...decision, + }); +} diff --git a/desktop/src/shared/api/tauriManagedAgents.ts b/desktop/src/shared/api/tauriManagedAgents.ts index ed7e053f259..cb112cc6e7c 100644 --- a/desktop/src/shared/api/tauriManagedAgents.ts +++ b/desktop/src/shared/api/tauriManagedAgents.ts @@ -114,3 +114,14 @@ export async function reconcileManagedAgentRuntimes( ): Promise { return invokeTauri("reconcile_managed_agent_runtimes", { communities }); } + +/** Read the Desktop-owned permission ledger for reconnect recovery. */ +export async function getManagedAgentPermissionLifecycle( + pubkey: string, + relayUrl: string, +): Promise<{ records: unknown[] }> { + return invokeTauri("get_managed_agent_permission_lifecycle", { + pubkey, + relayUrl, + }); +}