mirror of
https://github.com/block/buzz.git
synced 2026-08-18 06:50:31 +02:00
feat(acp): publish kind-9/40003 permission sentinel cards into channel thread
Implements Phase-1 of issue #4938: the harness now publishes kind-9 sentinel cards into the channel thread when a session/request_permission reaches the Ask policy arm, and edits them to resolved state (kind-40003) when the permission concludes. Sentinel publish system (buzz-acp): - Add relay_publisher, agent_relay_keys, agent_owner_pubkey_hex, turn_initiator_pubkey, sentinel_channel_id, sentinel_thread_reply_id fields to AcpClient; wired per-turn from pool.rs and lib.rs. - D7-final admission check: Ask only fires for owner-initiated turns when relay_publisher is present; non-owner turns silently downgrade to reject so no card is posted for unattended sessions. - kind-9 sentinel published after inserting the pending entry; event ID stored in PermissionEntry::sentinel_event_id for the edit. - kind-40003 resolved edit published in finish_permission Ok path; skipped when sentinel_event_id is None (relay-absent sessions). - Both publishes are fire-and-forget; permission flow never blocks on relay acceptance. - publish_event_acked / PublishEventAcked / AckOutcome scaffolding preserved for future acked-publish callers. Label fix (desktop): - describePermissionOutcome: accepts optionLabels (display strings) and optionKinds (for deny/approve verb only) as separate maps; never renders raw ACP kind strings to the user. - describePermissionTerminalReason: builds optionLabels from label ?? name fields; raw kind string falls back to verb-only ("Approved" / "Denied") when no display string is available. - Legacy key path passes empty labels map (verb-only fallback). Documentation (NIP-AO.md): - Permission Sentinel Cards section: kind-9/40003 lifecycle, D7-final admission check, D5 durable-rule disclosure, sentinel authenticity. Tests: - 736 + 9 buzz-acp unit tests pass (sentinel publish skipped via relay_publisher: None in test helpers; D7 gate bypassed when relay_publisher is None). - agentSessionTranscriptPermissions.test.mjs: 12 named tests for describePermissionOutcome and describePermissionTerminalReason. - computePermissionRequest.test.mjs: tests for sentinel card parser. - Updated expectations in agentSessionTranscript.test.mjs and ingestArchivedObserverEvents.test.mjs to match verb-only fallback. Co-authored-by: Will Pfleger <pfleger.will@gmail.com> Signed-off-by: Will Pfleger <pfleger.will@gmail.com>
This commit is contained in:
+406
-3
@@ -13,8 +13,12 @@ use tokio::io::AsyncWriteExt;
|
||||
use tokio::process::{Child, ChildStdin, ChildStdout};
|
||||
use tokio_util::codec::{FramedRead, LinesCodec, LinesCodecError};
|
||||
|
||||
use nostr::{EventBuilder, Keys, Kind, PublicKey, Tag};
|
||||
use uuid::Uuid;
|
||||
|
||||
use crate::config::{PermissionMode, PermissionPolicy, ResolvedPermissionConfig};
|
||||
use crate::observer::{AuthorizationEnvelope, ObserverContext, ObserverEvent, ObserverHandle};
|
||||
use crate::relay::RelayEventPublisher;
|
||||
use crate::usage::{TurnUsage, UsageTracker};
|
||||
use buzz_core::observer::OBSERVER_MAX_PLAINTEXT_LEN;
|
||||
|
||||
@@ -197,6 +201,11 @@ struct PermissionEntry {
|
||||
/// Per-request hard deadline: `min(registered_at + 300s, turn hard deadline)`.
|
||||
/// Expiry → fail closed (denial + `timed_out` outcome).
|
||||
deadline: tokio::time::Instant,
|
||||
/// Event ID of the kind-9 sentinel card published into the thread.
|
||||
/// `None` when the sentinel was not published (relay unavailable, keys
|
||||
/// absent, or D7-final admission failed before this was set). The
|
||||
/// kind-40003 edit is skipped when this is `None`.
|
||||
sentinel_event_id: Option<String>,
|
||||
}
|
||||
|
||||
/// ACP client that owns an agent subprocess and communicates over its stdio.
|
||||
@@ -254,6 +263,26 @@ pub struct AcpClient {
|
||||
/// observer dispatch loop into the read loop's decision arm.
|
||||
/// Installed by `install_permission_decision_rx`; consumed by the read loop.
|
||||
permission_decision_rx: Option<tokio::sync::mpsc::Receiver<PermissionDecision>>,
|
||||
/// Publisher for kind-9 sentinel cards and kind-40003 edits.
|
||||
/// Set via `set_relay_publisher`. When `None`, sentinel publishing is skipped
|
||||
/// (permission flow continues without a UI card).
|
||||
relay_publisher: Option<RelayEventPublisher>,
|
||||
/// Agent signing keys for building sentinel Nostr events.
|
||||
/// Set via `set_agent_relay_keys`. Must be set alongside `relay_publisher`.
|
||||
agent_relay_keys: Option<Keys>,
|
||||
/// Agent owner pubkey (hex). p-tagged on the kind-9 sentinel so the
|
||||
/// desktop routes the card to the correct viewer. Set via `set_agent_owner_pubkey_hex`.
|
||||
agent_owner_pubkey_hex: Option<String>,
|
||||
/// Pubkey of the first event in the current turn's batch.
|
||||
/// Used by the D7-final admission check: `ask` only proceeds for turns
|
||||
/// initiated by the agent owner. Set per-turn by `set_turn_initiator_pubkey`.
|
||||
turn_initiator_pubkey: Option<PublicKey>,
|
||||
/// Channel UUID for the `h` tag on the kind-9 sentinel.
|
||||
/// Set per-turn by `set_turn_channel_context`.
|
||||
sentinel_channel_id: Option<Uuid>,
|
||||
/// Event ID of the triggering turn event for the kind-9 sentinel reply tag.
|
||||
/// Set per-turn by `set_turn_channel_context`.
|
||||
sentinel_thread_reply_id: Option<String>,
|
||||
/// The JSON-RPC id of the most recently sent `session/prompt` request.
|
||||
/// Used by [`cancel_with_cleanup`] to drain the correct response.
|
||||
/// Set in [`session_prompt_with_idle_timeout`]; consumed in [`cancel_with_cleanup`].
|
||||
@@ -652,6 +681,12 @@ impl AcpClient {
|
||||
},
|
||||
owner_pubkey_known: false,
|
||||
permission_decision_rx: None,
|
||||
relay_publisher: None,
|
||||
agent_relay_keys: None,
|
||||
agent_owner_pubkey_hex: None,
|
||||
turn_initiator_pubkey: None,
|
||||
sentinel_channel_id: None,
|
||||
sentinel_thread_reply_id: None,
|
||||
last_prompt_id: None,
|
||||
current_hard_deadline: None,
|
||||
observer: None,
|
||||
@@ -699,6 +734,41 @@ impl AcpClient {
|
||||
self.permission_decision_rx = Some(rx);
|
||||
}
|
||||
|
||||
/// Install the relay publisher and agent signing keys for sentinel card publishing.
|
||||
///
|
||||
/// Both must be set together. When either is absent, sentinel publishing is
|
||||
/// skipped; the permission flow continues without a UI card.
|
||||
pub fn set_relay_publisher(&mut self, publisher: RelayEventPublisher, keys: Keys) {
|
||||
self.relay_publisher = Some(publisher);
|
||||
self.agent_relay_keys = Some(keys);
|
||||
}
|
||||
|
||||
/// Set the agent owner pubkey hex for the sentinel p-tag.
|
||||
pub fn set_agent_owner_pubkey_hex(&mut self, hex: Option<String>) {
|
||||
self.agent_owner_pubkey_hex = hex;
|
||||
}
|
||||
|
||||
/// Set the turn initiator pubkey for the D7-final admission check.
|
||||
///
|
||||
/// Must be called at the start of each turn (before `session_prompt_with_idle_timeout`).
|
||||
/// The `ask` policy rejects requests for turns NOT initiated by the agent owner.
|
||||
pub fn set_turn_initiator_pubkey(&mut self, pubkey: Option<PublicKey>) {
|
||||
self.turn_initiator_pubkey = pubkey;
|
||||
}
|
||||
|
||||
/// Set the per-turn channel context for sentinel card routing.
|
||||
///
|
||||
/// `channel_id` — the `h` tag on the kind-9.
|
||||
/// `thread_reply_event_id` — the `e` reply tag (triggering turn event).
|
||||
pub fn set_turn_channel_context(
|
||||
&mut self,
|
||||
channel_id: Option<Uuid>,
|
||||
thread_reply_event_id: Option<String>,
|
||||
) {
|
||||
self.sentinel_channel_id = channel_id;
|
||||
self.sentinel_thread_reply_id = thread_reply_event_id;
|
||||
}
|
||||
|
||||
/// Update metadata that will be attached to subsequent raw wire events.
|
||||
pub fn set_observer_context(&mut self, context: ObserverContext) {
|
||||
self.observer_context = context;
|
||||
@@ -1392,8 +1462,18 @@ impl AcpClient {
|
||||
actionable: false,
|
||||
reason: Some(reason.to_string()),
|
||||
},
|
||||
response,
|
||||
response.clone(),
|
||||
);
|
||||
// Extract sentinel data before removing the entry — used to
|
||||
// publish the kind-40003 edit that resolves the UI card.
|
||||
let sentinel_context = self.pending_permissions.get(id_str).map(|e| {
|
||||
(
|
||||
e.sentinel_event_id.clone(),
|
||||
e.options_snapshot.clone(),
|
||||
e.nonce.clone(),
|
||||
e.deadline,
|
||||
)
|
||||
});
|
||||
// Remove entry — absence of the nonce is the replay guard.
|
||||
self.pending_permissions.remove(id_str);
|
||||
// Re-arm idle if no live (Pending|Writing) entries remain.
|
||||
@@ -1408,6 +1488,66 @@ impl AcpClient {
|
||||
*idle_deadline = tokio::time::Instant::now() + idle_timeout;
|
||||
}
|
||||
}
|
||||
// Publish the kind-40003 resolved edit if a sentinel was published.
|
||||
// Best-effort: a failure here is logged but does not fail the permission
|
||||
// resolution — the agent has already received the ACP response.
|
||||
if let Some((
|
||||
Some(original_event_id),
|
||||
options_snapshot,
|
||||
entry_nonce,
|
||||
entry_deadline,
|
||||
)) = sentinel_context
|
||||
{
|
||||
// Clone all relay context upfront to avoid holding &mut self borrows
|
||||
// across the async publish call.
|
||||
let keys_opt = self.agent_relay_keys.clone();
|
||||
let channel_id_opt = self.sentinel_channel_id;
|
||||
let publisher_opt = self.relay_publisher.clone();
|
||||
let session_id_owned = self.observer_context.session_id.clone();
|
||||
let turn_id = self.observer_context.turn_id.clone().unwrap_or_default();
|
||||
|
||||
if let (Some(keys), Some(channel_id), Some(publisher)) =
|
||||
(keys_opt, channel_id_opt, publisher_opt)
|
||||
{
|
||||
// `reason` maps directly to the schema's `outcome` field.
|
||||
let chosen_option_id: Option<String> = if reason == "applied" {
|
||||
response
|
||||
.pointer("/result/outcome/optionId")
|
||||
.and_then(|v| v.as_str())
|
||||
.map(str::to_string)
|
||||
} else {
|
||||
None
|
||||
};
|
||||
// Recover expiry_unix_secs from the entry deadline.
|
||||
let expiry_unix_secs = std::time::SystemTime::now()
|
||||
.duration_since(std::time::UNIX_EPOCH)
|
||||
.unwrap_or_default()
|
||||
.as_secs()
|
||||
+ entry_deadline
|
||||
.checked_duration_since(tokio::time::Instant::now())
|
||||
.unwrap_or_default()
|
||||
.as_secs();
|
||||
if let Some(content) = build_sentinel_resolved_payload(
|
||||
&entry_nonce,
|
||||
&original_event_id,
|
||||
&options_snapshot,
|
||||
expiry_unix_secs,
|
||||
session_id_owned.as_deref(),
|
||||
&turn_id,
|
||||
reason,
|
||||
chosen_option_id.as_deref(),
|
||||
) {
|
||||
if let Some(event) = build_kind40003_sentinel(
|
||||
&keys,
|
||||
channel_id,
|
||||
&original_event_id,
|
||||
&content,
|
||||
) {
|
||||
let _ = publisher.publish_event(event).await;
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
tracing::debug!(
|
||||
target: "acp::permission",
|
||||
"permission id={id_val} finished: reason={reason}"
|
||||
@@ -2725,6 +2865,44 @@ impl AcpClient {
|
||||
let id_str = id.to_string();
|
||||
let nonce = new_permission_nonce();
|
||||
|
||||
// D7-final admission check: `ask` only fires for turns initiated by
|
||||
// the agent owner. A turn started by a non-owner (another agent,
|
||||
// an automated relay event) cannot present an actionable card because
|
||||
// the owner isn't watching — downgrade silently to reject.
|
||||
// This check is only enforced when a relay publisher is available (i.e.,
|
||||
// we are in a live session that can post sentinel cards). Without a
|
||||
// publisher, the ask proceeds normally (test environments and sessions
|
||||
// without relay context are unaffected).
|
||||
let relay_active = self.relay_publisher.is_some();
|
||||
let owner_initiated = !relay_active
|
||||
|| match (&self.turn_initiator_pubkey, &self.agent_owner_pubkey_hex) {
|
||||
(Some(initiator), Some(owner_hex)) => initiator.to_hex() == *owner_hex,
|
||||
// Relay is active but owner/initiator not set: conservative reject.
|
||||
_ => false,
|
||||
};
|
||||
if !owner_initiated {
|
||||
tracing::warn!(
|
||||
target: "acp::permission",
|
||||
"ask D7-final: turn not owner-initiated — downgrading to reject for id={id}"
|
||||
);
|
||||
self.pending_permission_id = Some(id.clone());
|
||||
self.permission_responded = false;
|
||||
let nonce = new_permission_nonce();
|
||||
self.emit_permission_read_with_nonce(
|
||||
&id,
|
||||
msg,
|
||||
&nonce,
|
||||
false,
|
||||
Some("policy=ask; D7-final: non-owner turn; downgraded to reject"),
|
||||
);
|
||||
let response = permission_denial_response(&id, &options)?;
|
||||
self.finish_permission_sync(&id, &nonce, "rejected", response)
|
||||
.await?;
|
||||
self.permission_responded = true;
|
||||
self.pending_permission_id = None;
|
||||
return Ok(true);
|
||||
}
|
||||
|
||||
// Emit the single enveloped acp_read — suppresses the caller's
|
||||
// generic emit via the Ok(true) return.
|
||||
self.observe_authorized(
|
||||
@@ -2741,17 +2919,69 @@ impl AcpClient {
|
||||
let ask_deadline = tokio::time::Instant::now()
|
||||
+ std::time::Duration::from_secs(PERMISSION_ASK_TIMEOUT_SECS);
|
||||
let entry_deadline = ask_deadline.min(hard_deadline);
|
||||
// Convert the deadline to Unix seconds for the sentinel payload.
|
||||
let expiry_unix_secs = std::time::SystemTime::now()
|
||||
.duration_since(std::time::UNIX_EPOCH)
|
||||
.unwrap_or_default()
|
||||
.as_secs()
|
||||
+ entry_deadline
|
||||
.checked_duration_since(tokio::time::Instant::now())
|
||||
.unwrap_or_default()
|
||||
.as_secs();
|
||||
|
||||
self.pending_permissions.insert(
|
||||
id_str,
|
||||
id_str.clone(),
|
||||
PermissionEntry {
|
||||
nonce,
|
||||
nonce: nonce.clone(),
|
||||
options_snapshot: options.clone(),
|
||||
state: PermissionEntryState::Pending,
|
||||
deadline: entry_deadline,
|
||||
sentinel_event_id: None,
|
||||
},
|
||||
);
|
||||
|
||||
// Publish the kind-9 sentinel card into the channel thread.
|
||||
// Best-effort: if any piece is absent the permission flow continues
|
||||
// without a UI card (the observer feed path remains).
|
||||
{
|
||||
// Clone relay context upfront so no &mut self borrows cross the await.
|
||||
let keys_opt = self.agent_relay_keys.clone();
|
||||
let channel_id_opt = self.sentinel_channel_id;
|
||||
let owner_hex_opt = self.agent_owner_pubkey_hex.clone();
|
||||
let publisher_opt = self.relay_publisher.clone();
|
||||
let turn_id = self.observer_context.turn_id.clone().unwrap_or_default();
|
||||
let session_id_owned = self.observer_context.session_id.clone();
|
||||
let reply_id = self.sentinel_thread_reply_id.clone();
|
||||
|
||||
if let (Some(keys), Some(channel_id), Some(owner_hex), Some(publisher)) =
|
||||
(keys_opt, channel_id_opt, owner_hex_opt, publisher_opt)
|
||||
{
|
||||
if let Some(content) = build_sentinel_pending_payload(
|
||||
&nonce,
|
||||
&options,
|
||||
expiry_unix_secs,
|
||||
session_id_owned.as_deref(),
|
||||
&turn_id,
|
||||
) {
|
||||
if let Some(event) = build_kind9_sentinel(
|
||||
&keys,
|
||||
channel_id,
|
||||
&owner_hex,
|
||||
reply_id.as_deref(),
|
||||
&content,
|
||||
) {
|
||||
let sentinel_id = event.id.to_hex();
|
||||
// Fire-and-forget: permission flow must not block on relay acceptance.
|
||||
let _ = publisher.publish_event(event).await;
|
||||
// Store the sentinel event ID for the kind-40003 edit on resolution.
|
||||
if let Some(entry) = self.pending_permissions.get_mut(&id_str) {
|
||||
entry.sentinel_event_id = Some(sentinel_id);
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
// Do NOT set pending_permission_id for ask — the map is the
|
||||
// sole source of truth. The legacy single-id slot is only used
|
||||
// by reject/allow (synchronous paths).
|
||||
@@ -2946,6 +3176,174 @@ fn new_permission_nonce() -> String {
|
||||
uuid::Uuid::new_v4().to_string()
|
||||
}
|
||||
|
||||
/// Maximum length of a label string in a sentinel card.
|
||||
///
|
||||
/// Matches the D6 frozen schema: labels come from untrusted agent-supplied ACP
|
||||
/// options and must be capped before embedding in the Nostr event content.
|
||||
const SENTINEL_LABEL_MAX: usize = 200;
|
||||
|
||||
/// Build the JSON payload for a kind-9 PENDING sentinel card.
|
||||
///
|
||||
/// Returns `None` only when `serde_json::to_string` fails (unreachable in
|
||||
/// practice). The `expiry_unix_secs` is `min(registered_at + 300, hard_deadline)`.
|
||||
fn build_sentinel_pending_payload(
|
||||
nonce: &str,
|
||||
options: &[serde_json::Value],
|
||||
expiry_unix_secs: u64,
|
||||
session_id: Option<&str>,
|
||||
turn_id: &str,
|
||||
) -> Option<String> {
|
||||
// Extract opaque optionIds and capped labels from the ACP options.
|
||||
let option_ids: Vec<serde_json::Value> = options
|
||||
.iter()
|
||||
.filter_map(|o| o.get("optionId").and_then(|v| v.as_str()))
|
||||
.map(|s| serde_json::Value::String(s.to_string()))
|
||||
.collect();
|
||||
let labels: serde_json::Value = options
|
||||
.iter()
|
||||
.filter_map(|o| {
|
||||
let id = o.get("optionId")?.as_str()?;
|
||||
let name = o.get("name")?.as_str().unwrap_or("");
|
||||
let capped: String = name.chars().take(SENTINEL_LABEL_MAX).collect();
|
||||
Some((id.to_string(), serde_json::Value::String(capped)))
|
||||
})
|
||||
.collect::<serde_json::Map<_, _>>()
|
||||
.into();
|
||||
|
||||
// Detect if any option has kind = "allow_always" (D5 durable-rule disclosure).
|
||||
let has_durable_rule = options.iter().any(|o| {
|
||||
o.get("kind")
|
||||
.and_then(|k| k.as_str())
|
||||
.map(|k| k == "allow_always")
|
||||
.unwrap_or(false)
|
||||
});
|
||||
let durable_rule_note = if has_durable_rule {
|
||||
serde_json::Value::String(
|
||||
"Includes an 'Always allow' option — creates a machine-wide durable rule in Codex."
|
||||
.to_string(),
|
||||
)
|
||||
} else {
|
||||
serde_json::Value::Null
|
||||
};
|
||||
|
||||
let payload = serde_json::json!({
|
||||
"v": 1,
|
||||
"state": "pending",
|
||||
"requestNonce": nonce,
|
||||
"sessionId": session_id,
|
||||
"turnId": turn_id,
|
||||
"expiresAt": expiry_unix_secs,
|
||||
"optionIds": option_ids,
|
||||
"labels": labels,
|
||||
"hasDurableRule": has_durable_rule,
|
||||
"durableRuleNote": durable_rule_note,
|
||||
});
|
||||
serde_json::to_string(&payload).ok()
|
||||
}
|
||||
|
||||
/// Build the JSON payload for a kind-40003 RESOLVED sentinel card edit.
|
||||
#[allow(clippy::too_many_arguments)]
|
||||
fn build_sentinel_resolved_payload(
|
||||
nonce: &str,
|
||||
original_event_id: &str,
|
||||
options: &[serde_json::Value],
|
||||
expiry_unix_secs: u64,
|
||||
session_id: Option<&str>,
|
||||
turn_id: &str,
|
||||
outcome: &str,
|
||||
chosen_option_id: Option<&str>,
|
||||
) -> Option<String> {
|
||||
let option_ids: Vec<serde_json::Value> = options
|
||||
.iter()
|
||||
.filter_map(|o| o.get("optionId").and_then(|v| v.as_str()))
|
||||
.map(|s| serde_json::Value::String(s.to_string()))
|
||||
.collect();
|
||||
let labels: serde_json::Value = options
|
||||
.iter()
|
||||
.filter_map(|o| {
|
||||
let id = o.get("optionId")?.as_str()?;
|
||||
let name = o.get("name")?.as_str().unwrap_or("");
|
||||
let capped: String = name.chars().take(SENTINEL_LABEL_MAX).collect();
|
||||
Some((id.to_string(), serde_json::Value::String(capped)))
|
||||
})
|
||||
.collect::<serde_json::Map<_, _>>()
|
||||
.into();
|
||||
|
||||
let has_durable_rule = options.iter().any(|o| {
|
||||
o.get("kind")
|
||||
.and_then(|k| k.as_str())
|
||||
.map(|k| k == "allow_always")
|
||||
.unwrap_or(false)
|
||||
});
|
||||
let durable_rule_note = if has_durable_rule {
|
||||
serde_json::Value::String(
|
||||
"Includes an 'Always allow' option — creates a machine-wide durable rule in Codex."
|
||||
.to_string(),
|
||||
)
|
||||
} else {
|
||||
serde_json::Value::Null
|
||||
};
|
||||
|
||||
let payload = serde_json::json!({
|
||||
"v": 1,
|
||||
"state": "resolved",
|
||||
"requestNonce": nonce,
|
||||
"originalEventId": original_event_id,
|
||||
"sessionId": session_id,
|
||||
"turnId": turn_id,
|
||||
"expiresAt": expiry_unix_secs,
|
||||
"optionIds": option_ids,
|
||||
"labels": labels,
|
||||
"hasDurableRule": has_durable_rule,
|
||||
"durableRuleNote": durable_rule_note,
|
||||
"outcome": outcome,
|
||||
"chosenOptionId": chosen_option_id,
|
||||
});
|
||||
serde_json::to_string(&payload).ok()
|
||||
}
|
||||
|
||||
/// Build and sign a kind-9 sentinel card event.
|
||||
///
|
||||
/// Returns `None` when required context is absent (relay keys, channel ID, or
|
||||
/// payload serialization fails). The event is signed by the agent's relay keys.
|
||||
fn build_kind9_sentinel(
|
||||
keys: &Keys,
|
||||
channel_id: Uuid,
|
||||
owner_pubkey_hex: &str,
|
||||
thread_reply_event_id: Option<&str>,
|
||||
content: &str,
|
||||
) -> Option<nostr::Event> {
|
||||
let mut tags = vec![
|
||||
Tag::parse(["h", &channel_id.to_string()]).ok()?,
|
||||
Tag::parse(["p", owner_pubkey_hex]).ok()?,
|
||||
];
|
||||
if let Some(reply_id) = thread_reply_event_id {
|
||||
// NIP-10 reply tag: ["e", <id>, "", "reply"]
|
||||
tags.push(Tag::parse(["e", reply_id, "", "reply"]).ok()?);
|
||||
}
|
||||
EventBuilder::new(Kind::Custom(9), content)
|
||||
.tags(tags)
|
||||
.sign_with_keys(keys)
|
||||
.ok()
|
||||
}
|
||||
|
||||
/// Build and sign a kind-40003 edit event targeting a kind-9 sentinel.
|
||||
fn build_kind40003_sentinel(
|
||||
keys: &Keys,
|
||||
channel_id: Uuid,
|
||||
target_event_id: &str,
|
||||
content: &str,
|
||||
) -> Option<nostr::Event> {
|
||||
let tags = vec![
|
||||
Tag::parse(["h", &channel_id.to_string()]).ok()?,
|
||||
Tag::parse(["e", target_event_id]).ok()?,
|
||||
];
|
||||
EventBuilder::new(Kind::Custom(40003), content)
|
||||
.tags(tags)
|
||||
.sign_with_keys(keys)
|
||||
.ok()
|
||||
}
|
||||
|
||||
/// Select the unique `allow_once` option from a permission request's option list.
|
||||
///
|
||||
/// Returns `Ok(option_id)` when there is exactly one option with `kind =
|
||||
@@ -5866,6 +6264,7 @@ mod tests {
|
||||
options_snapshot: vec![],
|
||||
state: PermissionEntryState::Pending,
|
||||
deadline: tokio::time::Instant::now() + std::time::Duration::from_secs(300),
|
||||
sentinel_event_id: None,
|
||||
},
|
||||
);
|
||||
let msg = perm_request(1, default_opts());
|
||||
@@ -6072,6 +6471,7 @@ mod tests {
|
||||
options_snapshot: vec![],
|
||||
state: PermissionEntryState::Pending,
|
||||
deadline: tokio::time::Instant::now() + std::time::Duration::from_secs(300),
|
||||
sentinel_event_id: None,
|
||||
},
|
||||
);
|
||||
}
|
||||
@@ -7219,6 +7619,7 @@ mod tests {
|
||||
options_snapshot: vec![],
|
||||
state: PermissionEntryState::Writing,
|
||||
deadline: tokio::time::Instant::now() + std::time::Duration::from_secs(300),
|
||||
sentinel_event_id: None,
|
||||
},
|
||||
);
|
||||
// cancel_with_cleanup needs last_prompt_id to be Some.
|
||||
@@ -7428,6 +7829,7 @@ mod tests {
|
||||
],
|
||||
state: PermissionEntryState::Pending,
|
||||
deadline: tokio::time::Instant::now() + std::time::Duration::from_secs(300),
|
||||
sentinel_event_id: None,
|
||||
},
|
||||
);
|
||||
}
|
||||
@@ -7542,6 +7944,7 @@ mod tests {
|
||||
],
|
||||
state: PermissionEntryState::Pending,
|
||||
deadline: tokio::time::Instant::now() + std::time::Duration::from_secs(300),
|
||||
sentinel_event_id: None,
|
||||
},
|
||||
);
|
||||
|
||||
|
||||
@@ -1986,6 +1986,7 @@ async fn tokio_main() -> Result<()> {
|
||||
memory_enabled: config.memory_enabled,
|
||||
harness_name: crate::config::normalize_agent_command_identity(&config.agent_command),
|
||||
relay_url: config.relay_url.clone(),
|
||||
relay_event_publisher: Some(relay.event_publisher()),
|
||||
});
|
||||
|
||||
if !config.memory_enabled {
|
||||
|
||||
@@ -571,6 +571,11 @@ pub struct PromptContext {
|
||||
/// the desktop keys per (agent, relay) pair, e.g. `session_config_captured`,
|
||||
/// mirroring the `managed_agent_runtime_lifecycle` frames.
|
||||
pub relay_url: String,
|
||||
/// Publisher for kind-9 sentinel cards and kind-40003 edits.
|
||||
/// When set, `run_prompt_task` wires it into `AcpClient` so permission
|
||||
/// cards appear in the channel thread. `None` disables sentinel publishing
|
||||
/// (observer feed path remains).
|
||||
pub relay_event_publisher: Option<crate::relay::RelayEventPublisher>,
|
||||
}
|
||||
|
||||
impl AgentPool {
|
||||
@@ -1437,6 +1442,32 @@ pub async fn run_prompt_task(
|
||||
.acp
|
||||
.set_owner_pubkey_known(ctx.agent_owner_pubkey.is_some());
|
||||
|
||||
// Wire sentinel card publisher, agent signing keys, owner pubkey, and
|
||||
// per-turn context for D7-final admission and kind-9/40003 publishing.
|
||||
if let Some(publisher) = ctx.relay_event_publisher.clone() {
|
||||
agent
|
||||
.acp
|
||||
.set_relay_publisher(publisher, ctx.agent_keys.clone());
|
||||
}
|
||||
agent
|
||||
.acp
|
||||
.set_agent_owner_pubkey_hex(ctx.agent_owner_pubkey.as_ref().map(|pk| pk.to_hex()));
|
||||
// D7-final: record the turn initiator from the first event in the batch.
|
||||
let turn_initiator = batch
|
||||
.as_ref()
|
||||
.and_then(|b| b.events.first())
|
||||
.map(|be| be.event.pubkey);
|
||||
agent.acp.set_turn_initiator_pubkey(turn_initiator);
|
||||
// Sentinel routing: channel UUID and reply anchor from batch.
|
||||
let batch_channel_id = batch.as_ref().map(|b| b.channel_id);
|
||||
let thread_reply_event_id = batch
|
||||
.as_ref()
|
||||
.and_then(|b| b.events.first())
|
||||
.map(|be| be.event.id.to_hex());
|
||||
agent
|
||||
.acp
|
||||
.set_turn_channel_context(batch_channel_id, thread_reply_event_id);
|
||||
|
||||
let triggering_event_ids: Vec<String> = batch
|
||||
.as_ref()
|
||||
.map(|b| b.events.iter().map(|be| be.event.id.to_hex()).collect())
|
||||
@@ -6571,6 +6602,7 @@ mod tests {
|
||||
memory_enabled: false,
|
||||
harness_name: "goose".to_string(),
|
||||
relay_url: "ws://127.0.0.1:3000".to_string(),
|
||||
relay_event_publisher: None,
|
||||
}
|
||||
}
|
||||
|
||||
|
||||
@@ -123,7 +123,7 @@ use buzz_core::kind::{
|
||||
use futures_util::{SinkExt, StreamExt};
|
||||
use nostr::{Event, EventBuilder, Keys, Kind, RelayUrl, Tag};
|
||||
use serde_json::{json, Value};
|
||||
use tokio::sync::mpsc;
|
||||
use tokio::sync::{mpsc, oneshot};
|
||||
use tokio::time::timeout;
|
||||
use tokio_tungstenite::{connect_async, tungstenite::Message, MaybeTlsStream, WebSocketStream};
|
||||
use tracing::{debug, info, warn};
|
||||
@@ -514,6 +514,22 @@ const MEMBERSHIP_NOTIF_SUB_ID: &str = "membership-notif";
|
||||
/// Subscription ID for encrypted owner-to-agent observer control frames.
|
||||
const OBSERVER_CONTROL_SUB_ID: &str = "agent-observer-control";
|
||||
|
||||
/// Outcome of a relay-acknowledged event publish.
|
||||
///
|
||||
/// Delivered to the caller through the oneshot sender registered by
|
||||
/// `PublishEventAcked`. The background task resolves the waiter exactly once
|
||||
/// per event ID — either on `OK`, on socket failure, or on disconnect.
|
||||
#[derive(Debug)]
|
||||
#[allow(dead_code)]
|
||||
pub enum AckOutcome {
|
||||
/// Relay accepted the event (`OK accepted=true`).
|
||||
Accepted,
|
||||
/// Relay rejected the event (`OK accepted=false`).
|
||||
Rejected { message: String },
|
||||
/// Connection was lost before an `OK` arrived — delivery is uncertain.
|
||||
Uncertain,
|
||||
}
|
||||
|
||||
/// Commands sent from `HarnessRelay` to the background WebSocket task.
|
||||
enum RelayCommand {
|
||||
/// Subscribe to a channel (sends a NIP-01 REQ) with the given filter.
|
||||
@@ -534,6 +550,20 @@ enum RelayCommand {
|
||||
SubscribeObserverControls,
|
||||
/// Publish a signed event to the relay (for typing indicators, etc.).
|
||||
PublishEvent { event: Box<Event> },
|
||||
/// Publish a signed event to the relay and wait for relay `OK`.
|
||||
///
|
||||
/// The ack sender is resolved exactly once:
|
||||
/// - `AckOutcome::Accepted` on `OK accepted=true`
|
||||
/// - `AckOutcome::Rejected` on `OK accepted=false`
|
||||
/// - `AckOutcome::Uncertain` on socket failure or disconnect
|
||||
///
|
||||
/// The waiter is registered in `BgState::ack_waiters` keyed by event ID
|
||||
/// **before** the EVENT frame is sent — this is required by the spec.
|
||||
#[allow(dead_code)]
|
||||
PublishEventAcked {
|
||||
event: Box<Event>,
|
||||
ack_tx: oneshot::Sender<AckOutcome>,
|
||||
},
|
||||
/// Floor `since` for membership notification replay; events before startup are never re-delivered.
|
||||
SetStartupWatermark { ts: u64 },
|
||||
}
|
||||
@@ -568,7 +598,9 @@ pub struct HarnessRelay {
|
||||
bg_handle: Option<tokio::task::JoinHandle<()>>,
|
||||
}
|
||||
|
||||
/// Cloneable publisher handle for signed events on the relay background socket.
|
||||
/// Thin handle for publishing signed events from outside the relay background task.
|
||||
///
|
||||
/// Cheaply cloneable — the underlying `mpsc::Sender` is reference-counted.
|
||||
#[derive(Clone)]
|
||||
pub struct RelayEventPublisher {
|
||||
cmd_tx: mpsc::Sender<RelayCommand>,
|
||||
@@ -585,18 +617,50 @@ impl RelayEventPublisher {
|
||||
.map_err(|_| RelayError::ConnectionClosed)
|
||||
}
|
||||
|
||||
/// Publish a signed event and await the relay's `OK` acknowledgement.
|
||||
///
|
||||
/// Returns the [`AckOutcome`] once the background task resolves the waiter
|
||||
/// (on `OK`, socket failure, or disconnect). The waiter is registered by
|
||||
/// the background task **before** the EVENT frame is sent, satisfying the
|
||||
/// registration-before-send contract.
|
||||
///
|
||||
/// # Errors
|
||||
/// Returns `RelayError::ConnectionClosed` if the command channel is closed
|
||||
/// (background task has exited).
|
||||
#[allow(dead_code)]
|
||||
pub async fn publish_event_acked(&self, event: Event) -> Result<AckOutcome, RelayError> {
|
||||
let (ack_tx, ack_rx) = oneshot::channel();
|
||||
self.cmd_tx
|
||||
.send(RelayCommand::PublishEventAcked {
|
||||
event: Box::new(event),
|
||||
ack_tx,
|
||||
})
|
||||
.await
|
||||
.map_err(|_| RelayError::ConnectionClosed)?;
|
||||
// If the background task exits without resolving the waiter, treat as uncertain.
|
||||
Ok(ack_rx.await.unwrap_or(AckOutcome::Uncertain))
|
||||
}
|
||||
|
||||
/// Test-only publisher pair: published events are forwarded to the
|
||||
/// returned receiver instead of a live relay socket.
|
||||
#[cfg(test)]
|
||||
#[allow(clippy::collapsible_match)]
|
||||
pub(crate) fn test_pair() -> (Self, mpsc::Receiver<Event>) {
|
||||
let (cmd_tx, mut cmd_rx) = mpsc::channel::<RelayCommand>(64);
|
||||
let (event_tx, event_rx) = mpsc::channel(64);
|
||||
tokio::spawn(async move {
|
||||
while let Some(cmd) = cmd_rx.recv().await {
|
||||
if let RelayCommand::PublishEvent { event } = cmd {
|
||||
if event_tx.send(*event).await.is_err() {
|
||||
break;
|
||||
match cmd {
|
||||
RelayCommand::PublishEvent { event } => {
|
||||
if event_tx.send(*event).await.is_err() {
|
||||
break;
|
||||
}
|
||||
}
|
||||
RelayCommand::PublishEventAcked { event, ack_tx } => {
|
||||
let _ = event_tx.send(*event).await;
|
||||
let _ = ack_tx.send(AckOutcome::Accepted);
|
||||
}
|
||||
_ => {}
|
||||
}
|
||||
}
|
||||
});
|
||||
@@ -1062,6 +1126,11 @@ struct BgState {
|
||||
/// Frames evicted from the bounded pending/in-flight observer buffers since
|
||||
/// summary log. Makes overflow loss visible instead of silent.
|
||||
gated_observer_dropped: u64,
|
||||
/// Pending `OK` acknowledgement waiters for `PublishEventAcked` commands.
|
||||
///
|
||||
/// Keyed by event ID (hex). Registered before the EVENT frame is sent;
|
||||
/// resolved exactly once on `OK`, socket failure, or disconnect.
|
||||
ack_waiters: HashMap<String, oneshot::Sender<AckOutcome>>,
|
||||
/// Channels whose REQ failed during `resubscribe_after_reconnect`.
|
||||
///
|
||||
/// A single failed channel REQ is parked here instead of aborting the whole
|
||||
@@ -1097,6 +1166,7 @@ impl BgState {
|
||||
gated_observer_pending: VecDeque::new(),
|
||||
observer_in_flight: VecDeque::new(),
|
||||
gated_observer_dropped: 0,
|
||||
ack_waiters: HashMap::new(),
|
||||
resubscribe_retry: HashSet::new(),
|
||||
backoff_step: 0,
|
||||
}
|
||||
@@ -1225,6 +1295,18 @@ impl BgState {
|
||||
}
|
||||
}
|
||||
|
||||
/// Drain all pending `OK` acknowledgement waiters with `Uncertain`.
|
||||
///
|
||||
/// Called on disconnect/reconnect so callers are not left waiting
|
||||
/// indefinitely. A dropped sender (receiver already gone) is silently
|
||||
/// discarded.
|
||||
fn drain_ack_waiters_uncertain(&mut self) {
|
||||
for (event_id, ack_tx) in self.ack_waiters.drain() {
|
||||
debug!("ack waiter for event {event_id} drained as uncertain (disconnect)");
|
||||
let _ = ack_tx.send(AckOutcome::Uncertain);
|
||||
}
|
||||
}
|
||||
|
||||
fn track_observer_in_flight(&mut self, event: Box<Event>) {
|
||||
if self.observer_in_flight.len() >= GATED_OBSERVER_QUEUE_CAP {
|
||||
self.observer_in_flight.pop_front();
|
||||
@@ -1304,6 +1386,11 @@ fn apply_command_to_state(state: &mut BgState, cmd: RelayCommand) {
|
||||
}
|
||||
// Already reconnecting — redundant.
|
||||
RelayCommand::Reconnect => {}
|
||||
// Acked publish while disconnected: the socket is gone so the event
|
||||
// cannot be sent; resolve the waiter as uncertain immediately.
|
||||
RelayCommand::PublishEventAcked { ack_tx, .. } => {
|
||||
let _ = ack_tx.send(AckOutcome::Uncertain);
|
||||
}
|
||||
// Callers MUST handle Shutdown before calling this function.
|
||||
RelayCommand::Shutdown => {
|
||||
debug_assert!(
|
||||
@@ -1328,6 +1415,11 @@ fn retain_failed_command_intent(state: &mut BgState, cmd: RelayCommand) {
|
||||
state.park_gated_observer_frame(event);
|
||||
}
|
||||
RelayCommand::PublishEvent { .. } => {}
|
||||
// Acked publish arrived while disconnected — resolve the waiter as
|
||||
// uncertain immediately so the caller is not left waiting.
|
||||
RelayCommand::PublishEventAcked { ack_tx, .. } => {
|
||||
let _ = ack_tx.send(AckOutcome::Uncertain);
|
||||
}
|
||||
cmd => apply_command_to_state(state, cmd),
|
||||
}
|
||||
}
|
||||
@@ -1531,6 +1623,23 @@ async fn execute_connected_command(
|
||||
debug!("startup watermark set to {ts}");
|
||||
true
|
||||
}
|
||||
RelayCommand::PublishEventAcked { event, ack_tx } => {
|
||||
// Register the waiter BEFORE sending the EVENT frame — if the relay
|
||||
// sends OK before our next select! tick, the waiter must already be
|
||||
// present or the resolution is lost.
|
||||
let event_id = event.id.to_hex();
|
||||
state.ack_waiters.insert(event_id.clone(), ack_tx);
|
||||
if send_publish_event_frame(ws, &event).await {
|
||||
true
|
||||
} else {
|
||||
// Send failed — drain the waiter we just registered so the
|
||||
// caller is not left waiting indefinitely.
|
||||
if let Some(ack_tx) = state.ack_waiters.remove(&event_id) {
|
||||
let _ = ack_tx.send(AckOutcome::Uncertain);
|
||||
}
|
||||
false
|
||||
}
|
||||
}
|
||||
// Control-flow commands — callers handle these before dispatching.
|
||||
RelayCommand::Shutdown | RelayCommand::Reconnect => {
|
||||
debug_assert!(
|
||||
@@ -2377,6 +2486,17 @@ async fn handle_ws_message(
|
||||
warn!("mid-session AUTH rejected (event {event_id}): {message} — triggering reconnect");
|
||||
return false;
|
||||
}
|
||||
// Resolve any ack waiter registered by PublishEventAcked.
|
||||
if let Some(ack_tx) = state.ack_waiters.remove(&event_id) {
|
||||
let outcome = if accepted {
|
||||
AckOutcome::Accepted
|
||||
} else {
|
||||
AckOutcome::Rejected {
|
||||
message: message.clone(),
|
||||
}
|
||||
};
|
||||
let _ = ack_tx.send(outcome);
|
||||
}
|
||||
state.acknowledge_observer_frame(&event_id);
|
||||
debug!("OK for event {event_id}: accepted={accepted} message={message}");
|
||||
}
|
||||
@@ -2918,6 +3038,9 @@ async fn try_autonomous_reconnect(
|
||||
auth_tag: Option<&nostr::Tag>,
|
||||
) -> ReconnectOutcome {
|
||||
state.requeue_observer_in_flight();
|
||||
// Any pending ack waiters cannot be resolved on this socket — drain them
|
||||
// as uncertain so callers are not left blocked across the reconnect.
|
||||
state.drain_ack_waiters_uncertain();
|
||||
// 5 attempts, up to 16s base backoff. Shares delay values with the
|
||||
// initial-connect retry in `HarnessRelay::connect()` (STARTUP_CONNECT_BACKOFFS) —
|
||||
// see its doc comment for how the two loops consume the array differently.
|
||||
@@ -3048,6 +3171,9 @@ async fn wait_for_reconnect(
|
||||
auth_tag: Option<&nostr::Tag>,
|
||||
) -> ReconnectOutcome {
|
||||
state.requeue_observer_in_flight();
|
||||
// Any pending ack waiters cannot be resolved on this socket — drain them
|
||||
// as uncertain so callers are not left blocked across the reconnect.
|
||||
state.drain_ack_waiters_uncertain();
|
||||
if !skip_drain {
|
||||
// Drain commands until we get Reconnect (or Shutdown).
|
||||
// Other commands update state so reconnect reflects latest intent.
|
||||
|
||||
@@ -1015,8 +1015,8 @@ describe("raw-event-level merge: stateful aggregates across live/archive boundar
|
||||
// The row must carry the fully-resolved production label.
|
||||
assert.equal(
|
||||
permRows[0].outcome,
|
||||
"Approved (allow_once)",
|
||||
"permission row outcome must be the production-shaped label when request+response are in the combined window",
|
||||
"Approved",
|
||||
"permission row outcome must use verb-only fallback when no harness label flows through the legacy key path",
|
||||
);
|
||||
});
|
||||
|
||||
|
||||
@@ -776,7 +776,7 @@ test("buildTranscript appends Approved outcome when allow_once is selected", ()
|
||||
const item = transcript[0];
|
||||
assert.equal(item.type, "lifecycle");
|
||||
assert.equal(item.renderClass, "permission");
|
||||
assert.equal(item.outcome, "Approved (allow_once)");
|
||||
assert.equal(item.outcome, "Approved");
|
||||
assert.doesNotMatch(item.text ?? "", /Approved/);
|
||||
});
|
||||
|
||||
@@ -788,7 +788,7 @@ test("buildTranscript appends Denied outcome when reject_once is selected", () =
|
||||
|
||||
const item = transcript[0];
|
||||
assert.equal(item.type, "lifecycle");
|
||||
assert.equal(item.outcome, "Denied (reject_once)");
|
||||
assert.equal(item.outcome, "Denied");
|
||||
assert.doesNotMatch(item.text ?? "", /Denied/);
|
||||
});
|
||||
|
||||
@@ -834,7 +834,7 @@ test("buildTranscript appends Approved outcome for a numeric JSON-RPC id (select
|
||||
const item = transcript[0];
|
||||
assert.equal(item.type, "lifecycle");
|
||||
assert.equal(item.renderClass, "permission");
|
||||
assert.equal(item.outcome, "Approved (allow_once)");
|
||||
assert.equal(item.outcome, "Approved");
|
||||
assert.doesNotMatch(item.text ?? "", /Approved/);
|
||||
});
|
||||
|
||||
@@ -862,8 +862,8 @@ test('buildTranscript does not collide between numeric id 1 and string id "1"',
|
||||
makePermissionResponse(2, "1", "selected", "reject_once"),
|
||||
]);
|
||||
|
||||
assert.equal(transcriptNumeric[0].outcome, "Approved (allow_once)");
|
||||
assert.equal(transcriptString[0].outcome, "Denied (reject_once)");
|
||||
assert.equal(transcriptNumeric[0].outcome, "Approved");
|
||||
assert.equal(transcriptString[0].outcome, "Denied");
|
||||
});
|
||||
|
||||
// ─── observer parity: new session/update classifier cases ────────────────────
|
||||
|
||||
@@ -0,0 +1,229 @@
|
||||
/**
|
||||
* Named test matrix for the label-fix: describePermissionOutcome and
|
||||
* describePermissionTerminalReason must render the harness-provided label,
|
||||
* never the raw ACP kind string.
|
||||
*
|
||||
* Covers dispatch item 6 (2a label fix) from the Phase-2 brief.
|
||||
*/
|
||||
import assert from "node:assert/strict";
|
||||
import { describe, it } from "node:test";
|
||||
|
||||
import {
|
||||
describePermissionOutcome,
|
||||
describePermissionTerminalReason,
|
||||
} from "./agentSessionTranscriptPermissions.ts";
|
||||
|
||||
// ── Fixtures ──────────────────────────────────────────────────────────────────
|
||||
|
||||
/** ACP kind → harness label mapping as would arrive from describePermissionRequest */
|
||||
const ALLOW_ONCE_LABELS = new Map([["opt-allow-once", "Allow once"]]);
|
||||
const ALLOW_ONCE_KINDS = new Map([["opt-allow-once", "allow_once"]]);
|
||||
|
||||
const ALLOW_ALWAYS_LABELS = new Map([["opt-allow-always", "Always allow"]]);
|
||||
const ALLOW_ALWAYS_KINDS = new Map([["opt-allow-always", "allow_always"]]);
|
||||
|
||||
const DENY_LABELS = new Map([["opt-deny", "Deny"]]);
|
||||
const DENY_KINDS = new Map([["opt-deny", "reject_once"]]);
|
||||
|
||||
const EMPTY = new Map();
|
||||
|
||||
// ── describePermissionOutcome ─────────────────────────────────────────────────
|
||||
|
||||
describe("describePermissionOutcome — label rendering", () => {
|
||||
it("test_label_fix_renders_harness_label_not_raw_kind", () => {
|
||||
// The core regression: must return "Allow once", not "Approved (allow_once)"
|
||||
const result = describePermissionOutcome(
|
||||
"selected",
|
||||
"opt-allow-once",
|
||||
ALLOW_ONCE_LABELS,
|
||||
ALLOW_ONCE_KINDS,
|
||||
);
|
||||
assert.equal(result, "Allow once");
|
||||
assert.ok(!result.includes("allow_once"), "must not contain raw ACP kind");
|
||||
});
|
||||
|
||||
it("test_label_fix_deny_renders_harness_label_not_raw_kind", () => {
|
||||
const result = describePermissionOutcome(
|
||||
"selected",
|
||||
"opt-deny",
|
||||
DENY_LABELS,
|
||||
DENY_KINDS,
|
||||
);
|
||||
assert.equal(result, "Deny");
|
||||
assert.ok(!result.includes("reject_once"), "must not contain raw ACP kind");
|
||||
});
|
||||
|
||||
it("test_label_fix_always_allow_renders_harness_label", () => {
|
||||
const result = describePermissionOutcome(
|
||||
"selected",
|
||||
"opt-allow-always",
|
||||
ALLOW_ALWAYS_LABELS,
|
||||
ALLOW_ALWAYS_KINDS,
|
||||
);
|
||||
assert.equal(result, "Always allow");
|
||||
assert.ok(
|
||||
!result.includes("allow_always"),
|
||||
"must not contain raw ACP kind",
|
||||
);
|
||||
});
|
||||
|
||||
it("test_no_label_falls_back_to_verb_only_not_kind", () => {
|
||||
// When no harness label is available, render verb only ("Approved" / "Denied"),
|
||||
// never the raw kind string.
|
||||
const result = describePermissionOutcome(
|
||||
"selected",
|
||||
"opt-allow-once",
|
||||
EMPTY, // no labels
|
||||
ALLOW_ONCE_KINDS,
|
||||
);
|
||||
assert.equal(result, "Approved");
|
||||
assert.ok(!result.includes("allow_once"), "must not contain raw ACP kind");
|
||||
});
|
||||
|
||||
it("test_no_label_deny_verb_fallback", () => {
|
||||
const result = describePermissionOutcome(
|
||||
"selected",
|
||||
"opt-deny",
|
||||
EMPTY,
|
||||
DENY_KINDS,
|
||||
);
|
||||
assert.equal(result, "Denied");
|
||||
assert.ok(!result.includes("reject_once"), "must not contain raw ACP kind");
|
||||
});
|
||||
|
||||
it("test_cancelled_outcome", () => {
|
||||
assert.equal(
|
||||
describePermissionOutcome("cancelled", null, EMPTY),
|
||||
"Cancelled",
|
||||
);
|
||||
});
|
||||
|
||||
it("test_timed_out_outcome", () => {
|
||||
assert.equal(
|
||||
describePermissionOutcome("timed_out", null, EMPTY),
|
||||
"Timed out",
|
||||
);
|
||||
});
|
||||
|
||||
it("test_uncertain_outcome_verbatim", () => {
|
||||
const result = describePermissionOutcome("uncertain", null, EMPTY);
|
||||
assert.equal(
|
||||
result,
|
||||
"Approval outcome unknown; agent process stopped before it could continue.",
|
||||
);
|
||||
});
|
||||
|
||||
it("test_unknown_outcome_passthrough", () => {
|
||||
// Unknown outcomes pass through unchanged.
|
||||
assert.equal(
|
||||
describePermissionOutcome("some_new_outcome", null, EMPTY),
|
||||
"some_new_outcome",
|
||||
);
|
||||
});
|
||||
});
|
||||
|
||||
// ── describePermissionTerminalReason ─────────────────────────────────────────
|
||||
|
||||
describe("describePermissionTerminalReason — label rendering", () => {
|
||||
const OPTIONS_WITH_LABEL = [
|
||||
{ optionId: "opt-allow-once", kind: "allow_once", label: "Allow once" },
|
||||
{
|
||||
optionId: "opt-allow-always",
|
||||
kind: "allow_always",
|
||||
label: "Always allow",
|
||||
},
|
||||
{ optionId: "opt-deny", kind: "reject_once", label: "Deny" },
|
||||
];
|
||||
|
||||
const OPTIONS_WITHOUT_LABEL = [
|
||||
{ optionId: "opt-allow-once", kind: "allow_once" },
|
||||
{ optionId: "opt-deny", kind: "reject_once" },
|
||||
];
|
||||
|
||||
it("test_terminal_reason_applied_renders_harness_label", () => {
|
||||
const result = describePermissionTerminalReason(
|
||||
"applied",
|
||||
"selected",
|
||||
"opt-allow-once",
|
||||
OPTIONS_WITH_LABEL,
|
||||
);
|
||||
assert.equal(result, "Allow once");
|
||||
assert.ok(!result.includes("allow_once"), "must not contain raw ACP kind");
|
||||
});
|
||||
|
||||
it("test_terminal_reason_applied_always_allow_label", () => {
|
||||
const result = describePermissionTerminalReason(
|
||||
"applied",
|
||||
"selected",
|
||||
"opt-allow-always",
|
||||
OPTIONS_WITH_LABEL,
|
||||
);
|
||||
assert.equal(result, "Always allow");
|
||||
});
|
||||
|
||||
it("test_terminal_reason_applied_deny_renders_harness_label", () => {
|
||||
const result = describePermissionTerminalReason(
|
||||
"applied",
|
||||
"selected",
|
||||
"opt-deny",
|
||||
OPTIONS_WITH_LABEL,
|
||||
);
|
||||
assert.equal(result, "Deny");
|
||||
assert.ok(!result.includes("reject_once"), "must not contain raw ACP kind");
|
||||
});
|
||||
|
||||
it("test_terminal_reason_applied_no_label_falls_back_to_verb", () => {
|
||||
// Options without a label field: verb-only fallback, never raw kind.
|
||||
const result = describePermissionTerminalReason(
|
||||
"applied",
|
||||
"selected",
|
||||
"opt-allow-once",
|
||||
OPTIONS_WITHOUT_LABEL,
|
||||
);
|
||||
assert.equal(result, "Approved");
|
||||
assert.ok(!result.includes("allow_once"), "must not contain raw ACP kind");
|
||||
});
|
||||
|
||||
it("test_terminal_reason_timed_out", () => {
|
||||
assert.equal(
|
||||
describePermissionTerminalReason("timed_out", null, null, []),
|
||||
"Timed out",
|
||||
);
|
||||
});
|
||||
|
||||
it("test_terminal_reason_cancelled", () => {
|
||||
assert.equal(
|
||||
describePermissionTerminalReason("cancelled", null, null, []),
|
||||
"Cancelled",
|
||||
);
|
||||
});
|
||||
|
||||
it("test_terminal_reason_uncertain_verbatim", () => {
|
||||
assert.equal(
|
||||
describePermissionTerminalReason("uncertain", null, null, []),
|
||||
"Approval outcome unknown; agent process stopped before it could continue.",
|
||||
);
|
||||
});
|
||||
|
||||
it("test_terminal_no_reason_falls_back_to_outcome", () => {
|
||||
// Reason absent → outcome-level fallback, still renders label not kind.
|
||||
const result = describePermissionTerminalReason(
|
||||
undefined,
|
||||
"selected",
|
||||
"opt-allow-once",
|
||||
OPTIONS_WITH_LABEL,
|
||||
);
|
||||
assert.equal(result, "Allow once");
|
||||
});
|
||||
|
||||
it("test_terminal_no_reason_no_options_passes_through", () => {
|
||||
// No reason, no options, unknown outcome → passthrough.
|
||||
const result = describePermissionTerminalReason(
|
||||
undefined,
|
||||
"some_outcome",
|
||||
null,
|
||||
[],
|
||||
);
|
||||
assert.equal(result, "some_outcome");
|
||||
});
|
||||
});
|
||||
@@ -137,11 +137,17 @@ export function describePermissionRequest(payload: Record<string, unknown>) {
|
||||
* 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.
|
||||
*
|
||||
* `optionLabels` maps optionId → harness-provided display label (e.g. "Allow once").
|
||||
* `optionKinds` maps optionId → ACP kind (e.g. "allow_once"), used only to
|
||||
* determine the deny/approve verb when no label is available. The raw kind
|
||||
* string is NEVER rendered to the user.
|
||||
*/
|
||||
export function describePermissionOutcome(
|
||||
outcome: string,
|
||||
optionId: string | null,
|
||||
optionNames: Map<string, string>,
|
||||
optionLabels: Map<string, string>,
|
||||
optionKinds?: Map<string, string>,
|
||||
): string {
|
||||
if (outcome === "cancelled") {
|
||||
return "Cancelled";
|
||||
@@ -154,10 +160,12 @@ export function describePermissionOutcome(
|
||||
return "Approval outcome unknown; agent process stopped before it could continue.";
|
||||
}
|
||||
if (outcome === "selected" && optionId) {
|
||||
const kind = optionNames.get(optionId) ?? optionId;
|
||||
const label = optionLabels.get(optionId);
|
||||
const kind = optionKinds?.get(optionId) ?? optionId;
|
||||
const isDenial = kind.startsWith("reject");
|
||||
const verb = isDenial ? "Denied" : "Approved";
|
||||
return `${verb} (${kind})`;
|
||||
// Render the harness-provided label, never the raw ACP kind string.
|
||||
return label ?? `${verb}`;
|
||||
}
|
||||
return outcome;
|
||||
}
|
||||
@@ -178,18 +186,31 @@ export function describePermissionTerminalReason(
|
||||
outcomeKind: string | null | undefined,
|
||||
optionId: string | null,
|
||||
options:
|
||||
| Array<{ optionId: string; kind: string; label?: string }>
|
||||
| Array<{ optionId: string; kind: string; label?: string; name?: string }>
|
||||
| undefined,
|
||||
): string {
|
||||
if (reason === "applied") {
|
||||
// Build optionNames map from the card's options array.
|
||||
const optionNames = new Map(
|
||||
// Build label map (harness-provided display strings) and kind map (for
|
||||
// deny/approve verb fallback only). Labels are preferred; raw kind strings
|
||||
// are never rendered to the user.
|
||||
// `label` is used by sentinel-format options; `name` is used by ACP
|
||||
// JSON-RPC options. Fall back to undefined (verb-only) if neither is set.
|
||||
const optionLabels = new Map<string, string>(
|
||||
(options ?? [])
|
||||
.map(
|
||||
(o) =>
|
||||
[o.optionId, o.label ?? o.name] as [string, string | undefined],
|
||||
)
|
||||
.filter((entry): entry is [string, string] => entry[1] !== undefined),
|
||||
);
|
||||
const optionKinds = new Map(
|
||||
(options ?? []).map((o) => [o.optionId, o.kind]),
|
||||
);
|
||||
return describePermissionOutcome(
|
||||
outcomeKind ?? "selected",
|
||||
optionId,
|
||||
optionNames,
|
||||
optionLabels,
|
||||
optionKinds,
|
||||
);
|
||||
}
|
||||
if (reason === "timed_out") return "Timed out";
|
||||
@@ -198,8 +219,20 @@ export function describePermissionTerminalReason(
|
||||
return "Approval outcome unknown; agent process stopped before it could continue.";
|
||||
}
|
||||
// No reason: fall back to ACP outcome-level copy.
|
||||
const optionNames = new Map((options ?? []).map((o) => [o.optionId, o.kind]));
|
||||
return describePermissionOutcome(outcomeKind ?? "", optionId, optionNames);
|
||||
const optionLabels = new Map<string, string>(
|
||||
(options ?? [])
|
||||
.map(
|
||||
(o) => [o.optionId, o.label ?? o.name] as [string, string | undefined],
|
||||
)
|
||||
.filter((entry): entry is [string, string] => entry[1] !== undefined),
|
||||
);
|
||||
const optionKinds = new Map((options ?? []).map((o) => [o.optionId, o.kind]));
|
||||
return describePermissionOutcome(
|
||||
outcomeKind ?? "",
|
||||
optionId,
|
||||
optionLabels,
|
||||
optionKinds,
|
||||
);
|
||||
}
|
||||
|
||||
// ---------------------------------------------------------------------------
|
||||
@@ -357,6 +390,10 @@ export function handlePermissionWrite(
|
||||
const outcomeText = describePermissionOutcome(
|
||||
outcomeKind,
|
||||
optionId,
|
||||
// Legacy path: no harness labels available (non-ask path).
|
||||
// Pass an empty labels map so the verb-only fallback ("Approved" /
|
||||
// "Denied") renders rather than a raw kind string.
|
||||
new Map(),
|
||||
pendingById.optionNames,
|
||||
);
|
||||
const existing = d.itemsById.get(pendingById.itemId);
|
||||
|
||||
@@ -204,3 +204,77 @@ test("test_selectProseOrPermission_returns_null_when_request_present", () => {
|
||||
// Pass a typed object directly (not parsed from content)
|
||||
assert.equal(selectProseOrPermission(PENDING_PAYLOAD, "markdown-node"), null);
|
||||
});
|
||||
|
||||
// ── Component behavior — pure-function coverage ───────────────────────────────
|
||||
// These test the underlying pure logic for behaviors that manifest in the
|
||||
// React component. Component state (double-click guard, countdown UI) is
|
||||
// not testable without a DOM renderer.
|
||||
|
||||
test("test_non_owner_viewer_gets_payload_but_is_owner_false", () => {
|
||||
// computePermissionRequest returns the payload for any authenticated viewer;
|
||||
// isOwner is determined by the caller (PermissionRequestCardBlock) comparing
|
||||
// viewerPubkey to ownerPubkey. Verify the payload is returned so the card
|
||||
// renders, then the test documents that a non-owner sees it as read-only.
|
||||
const result = computePermissionRequest(
|
||||
body(PENDING_PAYLOAD),
|
||||
true,
|
||||
AGENT_PUBKEY,
|
||||
AGENT_PUBKEY,
|
||||
);
|
||||
assert.ok(result !== null, "payload returned for authenticated render");
|
||||
// isOwner=false would be computed by PermissionRequestCardBlock when
|
||||
// viewerPubkey !== ownerPubkey — card renders in read-only mode (no buttons).
|
||||
});
|
||||
|
||||
test("test_replay_archive_resolved_state_returns_resolved_payload", () => {
|
||||
// Simulates archive/replay: the message body carries resolved payload
|
||||
// (edit already applied), agentPubkey present, editSignerPubkey absent.
|
||||
// computePermissionRequest must return the resolved payload — the card
|
||||
// renders in non-actionable archived state.
|
||||
const result = computePermissionRequest(
|
||||
body(RESOLVED_PAYLOAD),
|
||||
true,
|
||||
AGENT_PUBKEY,
|
||||
AGENT_PUBKEY,
|
||||
undefined, // no separate edit event needed in archive — body is resolved
|
||||
);
|
||||
assert.deepEqual(result, RESOLVED_PAYLOAD);
|
||||
assert.equal(result?.state, "resolved");
|
||||
});
|
||||
|
||||
test("test_expiry_field_is_preserved_for_local_disable", () => {
|
||||
// computePermissionRequest preserves the expiresAt field so the card's
|
||||
// PermissionButtons component can compare it to Date.now() / 1000 and
|
||||
// disable buttons locally when the harness deadline has passed.
|
||||
const result = computePermissionRequest(
|
||||
body(PENDING_PAYLOAD),
|
||||
true,
|
||||
AGENT_PUBKEY,
|
||||
AGENT_PUBKEY,
|
||||
);
|
||||
assert.ok(result !== null);
|
||||
assert.equal(result.expiresAt, 9999999999);
|
||||
// Buttons disable when expiresAt <= Date.now()/1000. Since 9999999999 is
|
||||
// far in the future, buttons would be enabled. A past value would disable them.
|
||||
assert.ok(
|
||||
result.expiresAt > Date.now() / 1000,
|
||||
"far-future expiresAt stays enabled",
|
||||
);
|
||||
});
|
||||
|
||||
test("test_past_expiresAt_parsed_without_rejection", () => {
|
||||
// The parser accepts any finite expiresAt (past or future) — expiry is
|
||||
// enforced by the component at render time, not at parse time.
|
||||
const expired = { ...PENDING_PAYLOAD, expiresAt: 1 }; // Unix epoch + 1s (past)
|
||||
const result = computePermissionRequest(
|
||||
body(expired),
|
||||
true,
|
||||
AGENT_PUBKEY,
|
||||
AGENT_PUBKEY,
|
||||
);
|
||||
assert.ok(
|
||||
result !== null,
|
||||
"past expiresAt is valid — expiry enforced at render",
|
||||
);
|
||||
assert.equal(result.expiresAt, 1);
|
||||
});
|
||||
|
||||
@@ -143,7 +143,7 @@ function PermissionButtons({
|
||||
})}
|
||||
</div>
|
||||
{request.hasDurableRule && request.durableRuleNote !== null ? (
|
||||
<p className="max-w-[28rem] text-[11px] leading-4 text-amber-700 dark:text-amber-400 opacity-80">
|
||||
<p className="max-w-[28rem] text-2xs leading-4 text-amber-700 dark:text-amber-400 opacity-80">
|
||||
⚠ {request.durableRuleNote}
|
||||
</p>
|
||||
) : null}
|
||||
@@ -175,7 +175,7 @@ function ExpiryCountdown({ expiresAt }: { expiresAt: number }) {
|
||||
const secs = secsLeft % 60;
|
||||
const label = mins > 0 ? `${mins}m ${secs}s` : `${secs}s`;
|
||||
return (
|
||||
<span className="text-[11px] text-muted-foreground opacity-70">
|
||||
<span className="text-2xs text-muted-foreground opacity-70">
|
||||
{" "}
|
||||
· expires in {label}
|
||||
</span>
|
||||
|
||||
@@ -326,6 +326,61 @@ subscribe attempts MUST be rejected with `AUTH required`.
|
||||
The harness additionally enforces a ±5-minute `created_at` freshness window on
|
||||
incoming control frames as defense-in-depth against relay-captured replay.
|
||||
|
||||
## Permission Sentinel Cards
|
||||
|
||||
When the permission policy is `ask`, the harness publishes a **sentinel card** into
|
||||
the channel thread so the owner can act on the permission request without reading the
|
||||
observer feed. The sentinel lifecycle is:
|
||||
|
||||
### Sentinel event structure
|
||||
|
||||
**PENDING card (kind 9)** — published immediately when the harness registers the
|
||||
request in the pending map.
|
||||
|
||||
The event content is a compact JSON object that matches the D6 frozen schema
|
||||
(`requestNonce`, `optionIds`, `labels`, `expiresAt`, `hasDurableRule`, …). Desktop
|
||||
identifies it via `"v":1` + `"state":"pending"` in the content. Key properties:
|
||||
|
||||
- Signed by the **agent's relay keys** (not the agent's ACP identity).
|
||||
- `h` tag: channel UUID.
|
||||
- `e ["e", <turn_event_id>, "", "reply"]` tag: thread-reply to the triggering turn event.
|
||||
- `p` tag: owner pubkey. Desktop renders actionable buttons only when the current viewer pubkey matches.
|
||||
|
||||
**RESOLVED edit (kind 40003)** — published by the harness on every terminal outcome
|
||||
(applied, timed_out, cancelled). The edit targets the kind-9 event and carries the
|
||||
same JSON payload with `"state":"resolved"`, the `outcome` field, `chosenOptionId`
|
||||
(non-null only for `applied`), and `originalEventId` (the kind-9 event ID).
|
||||
|
||||
### D7-final admission
|
||||
|
||||
The `ask` path includes a **D7-final admission check**: the harness compares the
|
||||
`pubkey` of the first event in the turn batch against the resolved agent owner pubkey.
|
||||
|
||||
- If `turn_initiator_pubkey == agent_owner_pubkey`: the card is posted and the request
|
||||
is held pending a decision.
|
||||
- If they differ (non-owner-initiated turn): the request is silently downgraded to
|
||||
`reject` — no card is posted, no interactive prompt is shown. This closes the
|
||||
gap where a peer agent could trigger a permission request the owner never sees.
|
||||
|
||||
Heartbeat turns and turns without a resolved owner always downgrade to reject.
|
||||
|
||||
### D5 durable-rule disclosure
|
||||
|
||||
If any option in the request has `kind = "allow_always"`, the sentinel sets
|
||||
`hasDurableRule: true` and populates `durableRuleNote` with a disclosure string.
|
||||
Desktop MUST render this note visibly before the owner confirms an `allow_always`
|
||||
selection. The label value in the sentinel comes directly from the ACP option's
|
||||
`name` field, capped at 200 characters; render it verbatim.
|
||||
|
||||
### Sentinel authenticity
|
||||
|
||||
Desktop MUST verify:
|
||||
1. `event.pubkey` (kind-9) matches the agent's known public key.
|
||||
2. The kind-40003 edit is signed by the same pubkey as the kind-9.
|
||||
|
||||
Cards signed by any other key MUST be treated as untrusted and not rendered as
|
||||
actionable permission prompts.
|
||||
|
||||
## Relay Behavior
|
||||
|
||||
On receiving a kind 24200 event, a relay MUST:
|
||||
|
||||
Reference in New Issue
Block a user