diff --git a/crates/buzz-acp/src/acp.rs b/crates/buzz-acp/src/acp.rs index 78db7ff718..6fa001be81 100644 --- a/crates/buzz-acp/src/acp.rs +++ b/crates/buzz-acp/src/acp.rs @@ -510,6 +510,12 @@ impl AcpClient { self.observer_context = context; } + /// Return the current observer routing context for terminal events emitted + /// after this client moves back to the harness loop. + pub(crate) fn observer_context(&self) -> &ObserverContext { + &self.observer_context + } + /// Return a clone of the observer handle, if attached. pub(crate) fn observer_handle(&self) -> Option { self.observer.clone() @@ -1211,14 +1217,18 @@ impl AcpClient { // so the ack_tx oneshot is never leaked silently). let mut steer_rx = self.steer_rx.take(); - // Tracks the in-flight steer write: `(request_id, ack_tx)`. While + // Tracks the in-flight steer write: + // `(request_id, ack_tx, requester_pubkey)`. While // `Some`, the steer arm is gated off so we don't stack writes, // and a response matching `id` is routed to the ack_tx instead // of being treated as the prompt result. Drained on every return // path with `PromptCompletedNeutral` so callers are never left // hanging. - let mut pending_steer: Option<(u64, tokio::sync::oneshot::Sender)> = - None; + let mut pending_steer: Option<( + u64, + tokio::sync::oneshot::Sender, + String, + )> = None; let now = Instant::now(); let mut idle_deadline = now + idle_timeout; @@ -1243,7 +1253,7 @@ impl AcpClient { // exists). Check the classified deadline here so a steady- // stream agent is still bounded. if Instant::now() >= next_deadline { - if let Some((_, ack_tx)) = pending_steer.take() { + if let Some((_, ack_tx, _)) = pending_steer.take() { // Prompt is timing out — release the withheld event via // PromptCompletedNeutral (no fallback signal: there is // no in-flight turn to signal once we return, and @@ -1320,7 +1330,8 @@ impl AcpClient { ); match self.write_ndjson(&msg).await { Ok(()) => { - pending_steer = Some((id, req.ack_tx)); + pending_steer = + Some((id, req.ack_tx, req.requester_pubkey)); } Err(e) => { tracing::warn!( @@ -1343,7 +1354,7 @@ impl AcpClient { // would catch this anyway, but firing the deadline arm // here makes the wakeup immediate (no extra reader poll // round-trip when stdout is idle). - if let Some((_, ack_tx)) = pending_steer.take() { + if let Some((_, ack_tx, _)) = pending_steer.take() { let _ = ack_tx.send(crate::pool::SteerAck::PromptCompletedNeutral); } if idle_fires_first { @@ -1367,13 +1378,13 @@ impl AcpClient { match read_result { None => { - if let Some((_, ack_tx)) = pending_steer.take() { + if let Some((_, ack_tx, _)) = pending_steer.take() { let _ = ack_tx.send(crate::pool::SteerAck::PromptCompletedNeutral); } return Err(AcpError::AgentExited); } Some(Err(LinesCodecError::MaxLineLengthExceeded)) => { - if let Some((_, ack_tx)) = pending_steer.take() { + if let Some((_, ack_tx, _)) = pending_steer.take() { let _ = ack_tx.send(crate::pool::SteerAck::PromptCompletedNeutral); } return Err(AcpError::Protocol( @@ -1381,7 +1392,7 @@ impl AcpClient { )); } Some(Err(e)) => { - if let Some((_, ack_tx)) = pending_steer.take() { + if let Some((_, ack_tx, _)) = pending_steer.take() { let _ = ack_tx.send(crate::pool::SteerAck::PromptCompletedNeutral); } return Err(AcpError::Io(std::io::Error::other(e))); @@ -1424,13 +1435,14 @@ impl AcpClient { // share the `no method` guard. if let Some(id) = msg.get("id") { if msg.get("method").is_none() { - if let Some((steer_id, _)) = pending_steer.as_ref() { + if let Some((steer_id, _, _)) = pending_steer.as_ref() { if *id == serde_json::json!(*steer_id) { // Take the ack_tx out and route the // response. We do not return — keep // reading until the prompt response // arrives. - let (_, ack_tx) = pending_steer.take().expect("just checked"); + let (_, ack_tx, requester_pubkey) = + pending_steer.take().expect("just checked"); let ack = if let Some(error) = msg.get("error") { let code = error .get("code") @@ -1441,6 +1453,12 @@ impl AcpClient { crate::pool::SteerError::AgentError { code, message }, ) } else { + self.observer_context + .add_requester_pubkey(requester_pubkey); + self.observe( + "turn_liveness", + serde_json::json!({"source": "native_steer"}), + ); let renew_now = Instant::now(); let new_deadline = renew_now + max_duration; if new_deadline > hard_deadline { @@ -1458,13 +1476,13 @@ impl AcpClient { } if *id == serde_json::json!(expected_id) { if let Some(error) = msg.get("error") { - if let Some((_, ack_tx)) = pending_steer.take() { + if let Some((_, ack_tx, _)) = pending_steer.take() { let _ = ack_tx .send(crate::pool::SteerAck::PromptCompletedNeutral); } return Err(agent_error_from_json(error)); } - if let Some((_, ack_tx)) = pending_steer.take() { + if let Some((_, ack_tx, _)) = pending_steer.take() { let _ = ack_tx.send(crate::pool::SteerAck::PromptCompletedNeutral); } @@ -3194,6 +3212,7 @@ mod tests { steer_tx .send(crate::pool::SteerRequest { prompt_blocks: vec!["test steer body".into()], + requester_pubkey: "requester-a".into(), ack_tx, }) .await @@ -3248,6 +3267,19 @@ mod tests { echo '{\"jsonrpc\":\"2.0\",\"id\":0,\"result\":{\"stopReason\":\"end_turn\"}}'; \ sleep 10"; let mut client = spawn_script(script).await; + let observer = crate::observer::ObserverHandle::in_process(); + let initial_requester = "aaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaa"; + let steered_requester = "bbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbb"; + client.set_observer(Some(observer.clone()), 0); + client.set_observer_context( + crate::observer::context_for_turn( + None, + Some("sess-test".into()), + "turn-test".into(), + "2026-07-24T10:00:00Z".into(), + ) + .with_requester_pubkeys(vec![initial_requester.into()]), + ); // Set active_run_id via a synthesized session_info_update so the // steer arm has a non-None value to read at write time. @@ -3263,6 +3295,7 @@ mod tests { steer_tx .send(crate::pool::SteerRequest { prompt_blocks: vec!["test steer body".into()], + requester_pubkey: steered_requester.into(), ack_tx, }) .await @@ -3301,6 +3334,18 @@ mod tests { crate::pool::SteerAck::Success => {} other => panic!("expected SteerAck::Success, got {other:?}"), } + let liveness = observer + .snapshot() + .into_iter() + .find(|event| { + event.kind == "turn_liveness" + && event.payload["source"] == serde_json::json!("native_steer") + }) + .expect("successful native steer should emit requester-visible liveness"); + assert_eq!( + liveness.requester_pubkeys, + [initial_requester.to_string(), steered_requester.to_string()] + ); } /// Steer-success renewal keeps the turn alive past the original hard @@ -3335,6 +3380,7 @@ mod tests { steer_tx .send(crate::pool::SteerRequest { prompt_blocks: vec!["steer body".into()], + requester_pubkey: "requester-b".into(), ack_tx, }) .await diff --git a/crates/buzz-acp/src/config.rs b/crates/buzz-acp/src/config.rs index a38d6faa14..4b1fbcfdda 100644 --- a/crates/buzz-acp/src/config.rs +++ b/crates/buzz-acp/src/config.rs @@ -111,6 +111,26 @@ impl std::fmt::Display for RespondTo { } } +/// Who receives encrypted relay observer telemetry. +/// +/// The registered owner always receives every frame. `Requester` additionally +/// sends turn-scoped activity to the authors whose events triggered that turn. +#[derive(Debug, Clone, Copy, Default, PartialEq, Eq, clap::ValueEnum)] +pub enum ObserverVisibility { + #[default] + OwnerOnly, + Requester, +} + +impl std::fmt::Display for ObserverVisibility { + fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result { + match self { + Self::OwnerOnly => f.write_str("owner-only"), + Self::Requester => f.write_str("requester"), + } + } +} + /// Permission mode for agents that support `session/set_config_option` with /// `configId: "mode"` (e.g. `claude-agent-acp`). /// @@ -468,6 +488,15 @@ pub struct CliArgs { #[arg(long, env = "BUZZ_ACP_RELAY_OBSERVER", default_value_t = false)] pub relay_observer: bool, + /// Observer visibility policy. The owner always receives every frame; + /// `requester` additionally sends turn activity to triggering authors. + #[arg( + long, + env = "BUZZ_ACP_OBSERVER_VISIBILITY", + default_value_t = ObserverVisibility::OwnerOnly + )] + pub observer_visibility: ObserverVisibility, + /// Connect and subscribe before starting the ACP/LLM subprocess pool. #[arg(long, env = "BUZZ_ACP_LAZY_POOL", default_value_t = false)] pub lazy_pool: bool, @@ -541,6 +570,8 @@ pub struct Config { pub has_generated_codex_config: bool, /// Whether to publish encrypted observer frames through the relay. pub relay_observer: bool, + /// Who receives encrypted turn-scoped observer telemetry. + pub observer_visibility: ObserverVisibility, /// Whether ACP/LLM subprocess initialization is deferred until accepted work arrives. pub lazy_pool: bool, /// Agent owner pubkey (hex). Used for `--respond-to=owner-only` gate. @@ -999,6 +1030,7 @@ impl Config { persona_env_vars, has_generated_codex_config, relay_observer: args.relay_observer, + observer_visibility: args.observer_visibility, lazy_pool: args.lazy_pool, agent_owner: args.agent_owner.map(|s| s.trim().to_ascii_lowercase()), no_base_prompt: args.no_base_prompt, @@ -1024,7 +1056,7 @@ impl Config { format!(" allowed_respond_to=[{}]", modes.join(",")) }; format!( - "relay={} pubkey={} agent_cmd={} {} mcp_cmd={} idle_timeout={}s max_turn={}s agents={} heartbeat={}s subscribe={:?} dedup={:?} meh={:?} ignore_self={} context_limit={} max_turns_per_session={} presence={} typing={} memory={} model={} permission_mode={} {}{}", + "relay={} pubkey={} agent_cmd={} {} mcp_cmd={} idle_timeout={}s max_turn={}s agents={} heartbeat={}s subscribe={:?} dedup={:?} meh={:?} ignore_self={} context_limit={} max_turns_per_session={} presence={} typing={} memory={} model={} permission_mode={} observer_visibility={} {}{}", self.relay_url, self.keys.public_key().to_hex(), self.agent_command, @@ -1045,6 +1077,7 @@ impl Config { self.memory_enabled, self.model.as_deref().unwrap_or("(agent default)"), self.permission_mode, + self.observer_visibility, respond_to_detail, allowed_respond_to_detail, ) @@ -1368,6 +1401,7 @@ mod tests { persona_env_vars: vec![], has_generated_codex_config: false, relay_observer: false, + observer_visibility: ObserverVisibility::OwnerOnly, lazy_pool: false, agent_owner: None, no_base_prompt: false, @@ -2048,6 +2082,26 @@ channels = "ALL" assert!(CliArgs::parse_from(["buzz-acp", "--private-key", &key, "--lazy-pool"]).lazy_pool); } + #[test] + fn observer_visibility_defaults_to_owner_only() { + let key = "0".repeat(64); + let args = CliArgs::parse_from(["buzz-acp", "--private-key", &key]); + assert_eq!(args.observer_visibility, ObserverVisibility::OwnerOnly); + } + + #[test] + fn observer_visibility_accepts_requester() { + let key = "0".repeat(64); + let args = CliArgs::parse_from([ + "buzz-acp", + "--private-key", + &key, + "--observer-visibility", + "requester", + ]); + assert_eq!(args.observer_visibility, ObserverVisibility::Requester); + } + #[test] fn test_summary_includes_agents_and_heartbeat() { let config = test_config(SubscribeMode::Mentions); diff --git a/crates/buzz-acp/src/lib.rs b/crates/buzz-acp/src/lib.rs index 0230ea0875..4a8fbd6f14 100644 --- a/crates/buzz-acp/src/lib.rs +++ b/crates/buzz-acp/src/lib.rs @@ -31,7 +31,7 @@ use buzz_core::observer::{ use clap::Parser; use config::{ AuthAgentArgs, AuthMethodsArgs, AuthenticateArgs, Config, DedupMode, ModelsArgs, - MultipleEventHandling, RespondTo, SubscribeMode, + MultipleEventHandling, ObserverVisibility, RespondTo, SubscribeMode, }; use filter::SubscriptionRule; use futures_util::FutureExt; @@ -408,13 +408,18 @@ impl ObserverPublishPacer { } } +struct ObserverPublishTarget { + agent_pubkey_hex: String, + owner_pubkey_hex: String, + owner_pubkey: PublicKey, + visibility: ObserverVisibility, +} + fn spawn_relay_observer_publisher( observer: observer::ObserverHandle, publisher: RelayEventPublisher, keys: nostr::Keys, - agent_pubkey_hex: String, - owner_pubkey_hex: String, - owner_pubkey: PublicKey, + target: ObserverPublishTarget, ) -> tokio::task::JoinHandle<()> { tokio::spawn(async move { // Subscribe BEFORE snapshotting so an event emitted between the two @@ -423,16 +428,7 @@ fn spawn_relay_observer_publisher( // high-water `seq` (monotonic, assigned at emit). let rx = observer.subscribe(); let snapshot = observer.snapshot(); - run_relay_observer_publisher( - snapshot, - rx, - publisher, - keys, - agent_pubkey_hex, - owner_pubkey_hex, - owner_pubkey, - ) - .await; + run_relay_observer_publisher(snapshot, rx, publisher, keys, target).await; }) } @@ -441,25 +437,14 @@ async fn run_relay_observer_publisher( mut rx: tokio::sync::broadcast::Receiver, publisher: RelayEventPublisher, keys: nostr::Keys, - agent_pubkey_hex: String, - owner_pubkey_hex: String, - owner_pubkey: PublicKey, + target: ObserverPublishTarget, ) { let mut coalescer = ObserverChunkCoalescer::default(); let mut pacer = ObserverPublishPacer::new(); let max_snapshot_seq = snapshot.iter().map(|event| event.seq).max().unwrap_or(0); for event in snapshot { for event in coalescer.ingest(event) { - publish_relay_observer_event( - &publisher, - &keys, - &agent_pubkey_hex, - &owner_pubkey_hex, - &owner_pubkey, - &mut pacer, - event, - ) - .await; + publish_relay_observer_event(&publisher, &keys, &target, &mut pacer, event).await; } } @@ -477,16 +462,14 @@ async fn run_relay_observer_publisher( } for event in coalescer.ingest(event) { publish_relay_observer_event( - &publisher, &keys, &agent_pubkey_hex, - &owner_pubkey_hex, &owner_pubkey, &mut pacer, event, + &publisher, &keys, &target, &mut pacer, event, ).await; } } Err(tokio::sync::broadcast::error::RecvError::Lagged(count)) => { for event in coalescer.flush() { publish_relay_observer_event( - &publisher, &keys, &agent_pubkey_hex, - &owner_pubkey_hex, &owner_pubkey, &mut pacer, event, + &publisher, &keys, &target, &mut pacer, event, ).await; } tracing::warn!(dropped = count, "relay observer publisher lagged"); @@ -494,8 +477,7 @@ async fn run_relay_observer_publisher( Err(tokio::sync::broadcast::error::RecvError::Closed) => { for event in coalescer.flush() { publish_relay_observer_event( - &publisher, &keys, &agent_pubkey_hex, - &owner_pubkey_hex, &owner_pubkey, &mut pacer, event, + &publisher, &keys, &target, &mut pacer, event, ).await; } break; @@ -506,15 +488,13 @@ async fn run_relay_observer_publisher( // Periodic flush ensures live streaming even during continuous chunk delivery. for event in coalescer.flush() { publish_relay_observer_event( - &publisher, &keys, &agent_pubkey_hex, - &owner_pubkey_hex, &owner_pubkey, &mut pacer, event, + &publisher, &keys, &target, &mut pacer, event, ).await; } } } } } - #[derive(Default)] struct ObserverChunkCoalescer { pending: Vec, @@ -534,6 +514,7 @@ struct ObserverChunkKey { session_id: Option, turn_id: Option, agent_index: Option, + requester_pubkeys: Vec, } /// Flush coalesced chunks before they exceed the NIP-44 plaintext limit (65,535 bytes). @@ -606,6 +587,7 @@ fn observer_chunk_key_and_text( session_id: event.session_id.clone(), turn_id: event.turn_id.clone(), agent_index: event.agent_index, + requester_pubkeys: event.requester_pubkeys.clone(), }, text, )) @@ -790,47 +772,105 @@ fn ceil_char_boundary(s: &str, mut i: usize) -> usize { async fn publish_relay_observer_event( publisher: &RelayEventPublisher, keys: &nostr::Keys, - agent_pubkey_hex: &str, - owner_pubkey_hex: &str, - owner_pubkey: &PublicKey, + target: &ObserverPublishTarget, pacer: &mut ObserverPublishPacer, mut event: observer::ObserverEvent, ) { - pacer.wait().await; // Trim oversized frames to fit the plaintext cap rather than letting // encrypt_observer_payload reject and drop them whole (silent telemetry loss). fit_observer_event_to_budget(&mut event); - let encrypted = match encrypt_observer_payload(keys, owner_pubkey, &event) { - Ok(encrypted) => encrypted, - Err(error) => { - tracing::warn!("failed to encrypt relay observer event: {error}"); - return; - } - }; - let builder = match buzz_sdk::build_agent_observer_frame( - owner_pubkey_hex, - agent_pubkey_hex, - OBSERVER_FRAME_TELEMETRY, - &encrypted, - ) { - Ok(builder) => builder, - Err(error) => { - tracing::warn!("failed to build relay observer event: {error}"); - return; + let mut recipients = vec![(target.owner_pubkey_hex.clone(), target.owner_pubkey)]; + if target.visibility == ObserverVisibility::Requester + && requester_visible_observer_event(&event) + { + for requester_hex in &event.requester_pubkeys { + if requester_hex.eq_ignore_ascii_case(&target.owner_pubkey_hex) { + continue; + } + match PublicKey::from_hex(requester_hex) { + Ok(requester_pubkey) => { + recipients.push((requester_hex.clone(), requester_pubkey)); + } + Err(error) => { + tracing::warn!( + requester = requester_hex, + "skipping invalid observer requester pubkey: {error}" + ); + } + } } - }; - let signed = match builder.sign_with_keys(keys) { - Ok(event) => event, - Err(error) => { - tracing::warn!("failed to sign relay observer event: {error}"); - return; + } + + for (recipient_hex, recipient_pubkey) in recipients { + // Every encrypted recipient copy consumes the same pacing budget. This + // keeps requester mode inside the existing 6/s and 90/min ceilings. + pacer.wait().await; + let encrypted = match encrypt_observer_payload(keys, &recipient_pubkey, &event) { + Ok(encrypted) => encrypted, + Err(error) => { + tracing::warn!( + recipient = recipient_hex, + "failed to encrypt relay observer event: {error}" + ); + continue; + } + }; + let builder = match buzz_sdk::build_agent_observer_frame( + &recipient_hex, + &target.agent_pubkey_hex, + OBSERVER_FRAME_TELEMETRY, + &encrypted, + ) { + Ok(builder) => builder, + Err(error) => { + tracing::warn!( + recipient = recipient_hex, + "failed to build relay observer event: {error}" + ); + continue; + } + }; + let signed = match builder.sign_with_keys(keys) { + Ok(event) => event, + Err(error) => { + tracing::warn!("failed to sign relay observer event: {error}"); + continue; + } + }; + if let Err(error) = publisher.publish_event(signed).await { + tracing::warn!( + recipient = recipient_hex, + "relay observer event dropped: {error}" + ); } - }; - if let Err(error) = publisher.publish_event(signed).await { - tracing::warn!("relay observer event dropped: {error}"); } } +fn requester_visible_observer_event(event: &observer::ObserverEvent) -> bool { + if event.turn_id.is_none() { + return false; + } + + match event.kind.as_str() { + "turn_started" | "turn_liveness" | "turn_completed" => true, + "acp_read" => requester_visible_session_update(&event.payload), + _ => false, + } +} + +fn requester_visible_session_update(payload: &serde_json::Value) -> bool { + if payload.get("method").and_then(serde_json::Value::as_str) != Some("session/update") { + return false; + } + + matches!( + payload + .pointer("/params/update/sessionUpdate") + .and_then(serde_json::Value::as_str), + Some("agent_message_chunk" | "tool_call" | "tool_call_update" | "plan") + ) +} + /// Maximum age (seconds) for an observer control frame to be considered fresh. const OBSERVER_CONTROL_FRESHNESS_SECS: i64 = 300; @@ -912,12 +952,7 @@ fn handle_cancel_turn_control( observer.emit( "control_result", None, - &observer::ObserverContext { - channel_id: Some(channel_id.to_string()), - session_id: None, - turn_id: None, - started_at: None, - }, + &observer::context_for(Some(channel_id), None, None), serde_json::json!({ "type": "cancel_turn", "status": status, @@ -989,12 +1024,7 @@ fn handle_switch_model_control( observer.emit( "control_result", None, - &observer::ObserverContext { - channel_id: Some(channel_id.to_string()), - session_id: None, - turn_id: None, - started_at: None, - }, + &observer::context_for(Some(channel_id), None, None), serde_json::json!({ "type": "switch_model", "status": status, @@ -1158,6 +1188,7 @@ struct RespawnResult { struct SteerAckEvent { channel_id: Uuid, event_id: String, + requester_pubkey: String, /// `Ok` if the read loop sent any of the locked `SteerAck` variants. /// `Err` if the oneshot was dropped without a send — should not happen /// under the current read-loop drains, but if it ever does the main @@ -1310,6 +1341,7 @@ async fn tokio_main() -> Result<()> { "agentArgs": config.agent_args, "parallelism": config.agents, "relayObserver": config.relay_observer, + "observerVisibility": config.observer_visibility.to_string(), }), ); } @@ -1408,6 +1440,7 @@ async fn tokio_main() -> Result<()> { pubkey_hex.clone(), owner_pubkey_hex, owner_pubkey, + config.observer_visibility, )); relay .subscribe_observer_controls() @@ -1487,16 +1520,19 @@ async fn tokio_main() -> Result<()> { } } - if let Some((observer, publisher, keys, agent_pubkey, owner_pubkey, owner)) = + if let Some((observer, publisher, keys, agent_pubkey, owner_pubkey, owner, visibility)) = relay_observer_publisher.take() { relay_observer_publisher_task = Some(spawn_relay_observer_publisher( observer, publisher, keys, - agent_pubkey, - owner_pubkey, - owner, + ObserverPublishTarget { + agent_pubkey_hex: agent_pubkey, + owner_pubkey_hex: owner_pubkey, + owner_pubkey: owner, + visibility, + }, )); } @@ -2416,6 +2452,7 @@ async fn tokio_main() -> Result<()> { Some(PoolEvent::SteerAck(SteerAckEvent { channel_id, event_id, + requester_pubkey, ack, })) => { // Goose-native steer attempt resolved. Locked semantics @@ -2506,6 +2543,13 @@ async fn tokio_main() -> Result<()> { "non-cancelling steer ack received" ); if matches!(ack, Ok(pool::SteerAck::Success)) { + if !pool.add_task_requester_pubkey(channel_id, &requester_pubkey) { + tracing::debug!( + channel = %channel_id, + requester = %requester_pubkey, + "native steer succeeded after task metadata was already removed" + ); + } queue.extend_in_flight_deadline(channel_id, config.max_turn_duration_secs); } if drop_withheld { @@ -2823,6 +2867,7 @@ fn try_native_steer( // steering (which is to inject only what's new). let (header, closing) = queue::native_steer_framing(); let event_id_hex = event.id.to_hex(); + let requester_pubkey = event.pubkey.to_hex(); let be = queue::BatchEvent { event, prompt_tag: prompt_tag.clone(), @@ -2834,6 +2879,7 @@ fn try_native_steer( let (ack_tx, ack_rx) = tokio::sync::oneshot::channel::(); let request = pool::SteerRequest { prompt_blocks: vec![body], + requester_pubkey: requester_pubkey.clone(), ack_tx, }; @@ -2867,6 +2913,7 @@ fn try_native_steer( let _ = ack_tx_clone.send(SteerAckEvent { channel_id, event_id: event_id_for_watcher, + requester_pubkey, ack, }); }); @@ -2920,6 +2967,7 @@ fn dispatch_pending( DedupMode::Queue => Some(batch.clone()), DedupMode::Drop => None, }; + let requester_pubkeys = pool::batch_requester_pubkeys(Some(&batch)); let result_tx = pool.result_tx(); let ctx_clone = Arc::clone(ctx); @@ -2963,6 +3011,7 @@ fn dispatch_pending( agent_index, channel_id: Some(channel_id), turn_id, + requester_pubkeys, recoverable_batch, control_tx: Some(control_tx), steer_tx, @@ -3203,11 +3252,12 @@ fn handle_prompt_result( .to_string(); let harness_pid = std::process::id(); - let channel_id = match &result.source { - PromptSource::Channel(ch) => Some(*ch), + let mut terminal_observer_context = result.agent.acp.observer_context().clone(); + terminal_observer_context.channel_id = match &result.source { + PromptSource::Channel(channel_id) => Some(channel_id.to_string()), PromptSource::Heartbeat => None, }; - let turn_id = result.turn_id.clone(); + terminal_observer_context.turn_id = Some(result.turn_id.clone()); let emit_turn_error = |error_msg: &str, error_code: Option| { if let Some(ref observer) = observer { let mut payload = serde_json::json!({ @@ -3220,7 +3270,7 @@ fn handle_prompt_result( observer.emit( "turn_error", Some(agent_index), - &observer::context_for(channel_id, None, Some(turn_id.clone())), + &terminal_observer_context, payload, ); } @@ -3446,10 +3496,12 @@ fn recover_panicked_agent( } if let Some(ref observer) = observer { + let context = observer::context_for(meta.channel_id, None, Some(meta.turn_id)) + .with_requester_pubkeys(meta.requester_pubkeys); observer.emit( "agent_panic", Some(i), - &observer::context_for(meta.channel_id, None, Some(meta.turn_id)), + &context, serde_json::json!({ "outcome": "panic", "error": format!("Agent task panicked: {join_error}"), @@ -3576,6 +3628,7 @@ fn dispatch_heartbeat( agent_index, channel_id: None, turn_id, + requester_pubkeys: Vec::new(), recoverable_batch: None, control_tx: None, steer_tx: None, @@ -4318,6 +4371,7 @@ mod owner_control_command_tests { agent_index: 0, channel_id: Some(channel_id), turn_id: "test-turn-id".to_string(), + requester_pubkeys: Vec::new(), recoverable_batch: None, control_tx: Some(control_tx), steer_tx: None, @@ -4780,9 +4834,12 @@ mod observer_snapshot_race_tests { rx, publisher, agent_keys.clone(), - agent_keys.public_key().to_hex(), - owner_keys.public_key().to_hex(), - owner_keys.public_key(), + ObserverPublishTarget { + agent_pubkey_hex: agent_keys.public_key().to_hex(), + owner_pubkey_hex: owner_keys.public_key().to_hex(), + owner_pubkey: owner_keys.public_key(), + visibility: ObserverVisibility::OwnerOnly, + }, ) .await; @@ -4803,6 +4860,248 @@ mod observer_snapshot_race_tests { } } +#[cfg(test)] +mod observer_visibility_publisher_tests { + use super::*; + use nostr::Keys; + + #[tokio::test(start_paused = true)] + async fn requester_visibility_encrypts_turn_frames_to_each_batch_author() { + let observer = observer::ObserverHandle::in_process(); + let agent_keys = Keys::generate(); + let owner_keys = Keys::generate(); + let requester_a = Keys::generate(); + let requester_b = Keys::generate(); + let context = observer::context_for_turn( + Some(Uuid::new_v4()), + Some("session-1".into()), + "turn-1".into(), + "2026-07-24T10:00:00Z".into(), + ) + .with_requester_pubkeys(vec![ + requester_b.public_key().to_hex(), + requester_a.public_key().to_hex(), + requester_a.public_key().to_hex(), + ]); + observer.emit("turn_started", Some(0), &context, serde_json::json!({})); + let rx = observer.subscribe(); + let snapshot = observer.snapshot(); + drop(observer); + let (publisher, mut published_rx) = RelayEventPublisher::test_pair(); + + run_relay_observer_publisher( + snapshot, + rx, + publisher, + agent_keys.clone(), + ObserverPublishTarget { + agent_pubkey_hex: agent_keys.public_key().to_hex(), + owner_pubkey_hex: owner_keys.public_key().to_hex(), + owner_pubkey: owner_keys.public_key(), + visibility: ObserverVisibility::Requester, + }, + ) + .await; + + let mut owner_frames = 0; + let mut requester_a_frames = 0; + let mut requester_b_frames = 0; + while let Some(event) = published_rx.recv().await { + if decrypt_observer_payload::(&owner_keys, &event).is_ok() { + owner_frames += 1; + } + if decrypt_observer_payload::(&requester_a, &event).is_ok() { + requester_a_frames += 1; + } + if decrypt_observer_payload::(&requester_b, &event).is_ok() { + requester_b_frames += 1; + } + } + + assert_eq!(owner_frames, 1); + assert_eq!(requester_a_frames, 1); + assert_eq!(requester_b_frames, 1); + } + + fn observer_event(kind: &str, payload: serde_json::Value) -> observer::ObserverEvent { + observer::ObserverEvent { + seq: 1, + timestamp: "2026-07-24T10:00:00Z".into(), + kind: kind.into(), + agent_index: Some(0), + channel_id: Some(Uuid::new_v4().to_string()), + session_id: Some("session-1".into()), + turn_id: Some("turn-1".into()), + started_at: Some("2026-07-24T10:00:00Z".into()), + requester_pubkeys: vec![Keys::generate().public_key().to_hex()], + payload, + } + } + + #[test] + fn requester_visibility_is_a_fail_closed_allowlist() { + for kind in ["turn_started", "turn_liveness", "turn_completed"] { + assert!( + requester_visible_observer_event(&observer_event(kind, serde_json::json!({}))), + "{kind}" + ); + } + + for update_type in [ + "agent_message_chunk", + "tool_call", + "tool_call_update", + "plan", + ] { + let event = observer_event( + "acp_read", + serde_json::json!({ + "jsonrpc": "2.0", + "method": "session/update", + "params": {"update": {"sessionUpdate": update_type}}, + }), + ); + assert!(requester_visible_observer_event(&event), "{update_type}"); + } + + let private_events = [ + observer_event( + "acp_write", + serde_json::json!({ + "method": "session/new", + "params": {"systemPrompt": "owner-private instructions"}, + }), + ), + observer_event( + "acp_read", + serde_json::json!({ + "id": 1, + "result": {"configOptions": [{"id": "system-prompt"}]}, + }), + ), + observer_event( + "acp_read", + serde_json::json!({ + "method": "session/request_permission", + "params": {"options": []}, + }), + ), + observer_event( + "acp_read", + serde_json::json!({ + "method": "session/update", + "params": { + "update": { + "sessionUpdate": "agent_thought_chunk", + "content": {"text": "private reasoning"}, + }, + }, + }), + ), + observer_event( + "acp_read", + serde_json::json!({ + "method": "session/update", + "params": { + "update": { + "sessionUpdate": "config_option_update", + "configOptions": [{"id": "system-prompt"}], + }, + }, + }), + ), + observer_event("turn_error", serde_json::json!({"error": "private path"})), + observer_event("unknown_future_kind", serde_json::json!({})), + ]; + for event in private_events { + assert!(!requester_visible_observer_event(&event), "{}", event.kind); + } + + let mut no_turn = observer_event("turn_started", serde_json::json!({})); + no_turn.turn_id = None; + assert!(!requester_visible_observer_event(&no_turn)); + } + + #[tokio::test(start_paused = true)] + async fn requester_visibility_keeps_session_setup_owner_only() { + let observer = observer::ObserverHandle::in_process(); + let agent_keys = Keys::generate(); + let owner_keys = Keys::generate(); + let requester = Keys::generate(); + let context = observer::context_for_turn( + Some(Uuid::new_v4()), + None, + "turn-1".into(), + "2026-07-24T10:00:00Z".into(), + ) + .with_requester_pubkeys(vec![requester.public_key().to_hex()]); + observer.emit( + "acp_write", + Some(0), + &context, + serde_json::json!({ + "method": "session/new", + "params": {"systemPrompt": "owner-private instructions"}, + }), + ); + let rx = observer.subscribe(); + let snapshot = observer.snapshot(); + drop(observer); + let (publisher, mut published_rx) = RelayEventPublisher::test_pair(); + + run_relay_observer_publisher( + snapshot, + rx, + publisher, + agent_keys.clone(), + ObserverPublishTarget { + agent_pubkey_hex: agent_keys.public_key().to_hex(), + owner_pubkey_hex: owner_keys.public_key().to_hex(), + owner_pubkey: owner_keys.public_key(), + visibility: ObserverVisibility::Requester, + }, + ) + .await; + + let mut owner_frames = 0; + let mut requester_frames = 0; + while let Some(event) = published_rx.recv().await { + if decrypt_observer_payload::(&owner_keys, &event).is_ok() { + owner_frames += 1; + } + if decrypt_observer_payload::(&requester, &event).is_ok() { + requester_frames += 1; + } + } + + assert_eq!(owner_frames, 1); + assert_eq!(requester_frames, 0); + } + + #[test] + fn legacy_owner_only_frame_kinds_never_fan_out_to_requesters() { + for kind in [ + "control_result", + "session_config_captured", + "managed_agent_runtime_lifecycle", + ] { + let event = observer::ObserverEvent { + seq: 1, + timestamp: "2026-07-24T10:00:00Z".into(), + kind: kind.into(), + agent_index: Some(0), + channel_id: Some(Uuid::new_v4().to_string()), + session_id: Some("session-1".into()), + turn_id: Some("turn-1".into()), + started_at: Some("2026-07-24T10:00:00Z".into()), + requester_pubkeys: vec![Keys::generate().public_key().to_hex()], + payload: serde_json::json!({}), + }; + assert!(!requester_visible_observer_event(&event), "{kind}"); + } + } +} + #[cfg(test)] mod observer_publish_pacer_tests { use super::*; @@ -4856,6 +5155,7 @@ mod observer_chunk_coalescer_tests { session_id: Some("session-1".to_string()), turn_id: Some("turn-1".to_string()), started_at: None, + requester_pubkeys: Vec::new(), payload: serde_json::json!({ "jsonrpc": "2.0", "method": "session/update", @@ -4884,6 +5184,7 @@ mod observer_chunk_coalescer_tests { session_id: Some("session-1".to_string()), turn_id: Some("turn-1".to_string()), started_at: None, + requester_pubkeys: Vec::new(), payload: serde_json::json!({ "type": "turn_started" }), } } @@ -4980,6 +5281,7 @@ mod build_mcp_servers_tests { persona_env_vars: vec![], has_generated_codex_config: false, relay_observer: false, + observer_visibility: config::ObserverVisibility::OwnerOnly, lazy_pool: false, agent_owner: None, no_base_prompt: false, @@ -5146,6 +5448,7 @@ mod error_outcome_emission_tests { persona_env_vars: vec![], has_generated_codex_config: false, relay_observer: false, + observer_visibility: config::ObserverVisibility::OwnerOnly, lazy_pool: false, agent_owner: None, no_base_prompt: false, @@ -5207,6 +5510,7 @@ mod error_outcome_emission_tests { agent_index: 0, channel_id: None, turn_id: "test-turn-id".to_string(), + requester_pubkeys: Vec::new(), recoverable_batch: None, control_tx: None, steer_tx: None, @@ -5283,6 +5587,7 @@ mod error_outcome_emission_tests { agent_index: 0, channel_id: Some(channel_id), turn_id: "panic-turn-id".to_string(), + requester_pubkeys: vec!["requester".to_string()], recoverable_batch: None, control_tx: None, steer_tx: None, @@ -5330,6 +5635,7 @@ mod error_outcome_emission_tests { Some(channel_id.to_string().as_str()) ); assert_eq!(panic.turn_id.as_deref(), Some("panic-turn-id")); + assert_eq!(panic.requester_pubkeys, ["requester"]); } #[tokio::test] @@ -5375,6 +5681,7 @@ mod error_outcome_emission_tests { agent_index: 0, channel_id: None, turn_id: "test-turn-id".to_string(), + requester_pubkeys: Vec::new(), recoverable_batch: None, control_tx: None, steer_tx: None, @@ -5466,6 +5773,7 @@ mod error_outcome_emission_tests { agent_index: 0, channel_id: None, turn_id: "test-turn-id".to_string(), + requester_pubkeys: Vec::new(), recoverable_batch: None, control_tx: None, steer_tx: None, @@ -5571,6 +5879,7 @@ mod error_outcome_emission_tests { agent_index: 0, channel_id: None, turn_id: "test-turn-id".to_string(), + requester_pubkeys: Vec::new(), recoverable_batch: None, control_tx: None, steer_tx: None, @@ -5647,6 +5956,7 @@ mod error_outcome_emission_tests { agent_index: 0, channel_id: None, turn_id: "test-turn-id".to_string(), + requester_pubkeys: Vec::new(), recoverable_batch: None, control_tx: None, steer_tx: None, @@ -5741,6 +6051,7 @@ mod error_outcome_emission_tests { agent_index: 0, channel_id: None, turn_id: "test-turn-id".to_string(), + requester_pubkeys: Vec::new(), recoverable_batch: None, control_tx: None, steer_tx: None, @@ -5857,6 +6168,7 @@ mod error_outcome_emission_tests { agent_index: 0, channel_id: None, turn_id: "test-turn-id".to_string(), + requester_pubkeys: Vec::new(), recoverable_batch: None, control_tx: None, steer_tx: None, @@ -5996,6 +6308,7 @@ mod error_outcome_emission_tests { agent_index: 0, channel_id: None, turn_id: "test-turn-id".to_string(), + requester_pubkeys: Vec::new(), recoverable_batch: None, control_tx: None, steer_tx: None, @@ -6184,6 +6497,7 @@ mod error_outcome_emission_tests { agent_index: 0, channel_id: None, turn_id: "test-turn-id".to_string(), + requester_pubkeys: Vec::new(), recoverable_batch: None, control_tx: None, steer_tx: None, @@ -6269,6 +6583,7 @@ mod error_outcome_emission_tests { agent_index: 0, channel_id: None, turn_id: "test-turn-id".to_string(), + requester_pubkeys: Vec::new(), recoverable_batch: None, control_tx: None, steer_tx: None, @@ -6334,6 +6649,7 @@ mod observer_payload_trim_tests { session_id: Some("sess-1".to_string()), turn_id: Some("turn-1".to_string()), started_at: None, + requester_pubkeys: Vec::new(), payload, } } diff --git a/crates/buzz-acp/src/observer.rs b/crates/buzz-acp/src/observer.rs index 7029e5af6d..51c580132f 100644 --- a/crates/buzz-acp/src/observer.rs +++ b/crates/buzz-acp/src/observer.rs @@ -28,6 +28,13 @@ pub struct ObserverContext { pub turn_id: Option, /// RFC3339 timestamp at which the current turn began, when known. pub started_at: Option, + /// Authors whose triggering events started this turn. + /// + /// Internal routing metadata only. It is never serialized into the + /// encrypted observer payload. Cloned contexts share this recipient set + /// so a requester accepted through a native mid-turn steer is visible to + /// liveness and terminal observers for the same turn. + requester_pubkeys: Arc>>, } /// Handle used by the harness to publish local observer events. @@ -74,6 +81,9 @@ pub struct ObserverEvent { /// RFC3339 timestamp at which the current turn began, when known. #[serde(skip_serializing_if = "Option::is_none")] pub started_at: Option, + /// Per-turn requester recipients. Internal publisher metadata only. + #[serde(skip)] + pub requester_pubkeys: Vec, /// Raw or semantic event payload. pub payload: serde_json::Value, } @@ -117,6 +127,7 @@ impl ObserverHandle { session_id: context.session_id.clone(), turn_id: context.turn_id.clone(), started_at: context.started_at.clone(), + requester_pubkeys: context.requester_pubkeys(), payload, }; @@ -147,6 +158,7 @@ pub fn context_for( session_id, turn_id, started_at: None, + requester_pubkeys: Arc::new(Mutex::new(Vec::new())), } } @@ -162,5 +174,65 @@ pub fn context_for_turn( session_id, turn_id: Some(turn_id), started_at: Some(started_at), + requester_pubkeys: Arc::new(Mutex::new(Vec::new())), + } +} + +impl ObserverContext { + /// Attach the normalized, de-duplicated authors whose events triggered + /// this turn. The list remains process-local and is used only to select + /// NIP-44 recipients. + pub fn with_requester_pubkeys(mut self, requester_pubkeys: Vec) -> Self { + let mut requester_pubkeys = requester_pubkeys + .into_iter() + .map(|pubkey| pubkey.to_ascii_lowercase()) + .collect::>(); + requester_pubkeys.sort(); + requester_pubkeys.dedup(); + self.requester_pubkeys = Arc::new(Mutex::new(requester_pubkeys)); + self + } + + /// Add a requester to this turn's shared recipient set. + pub fn add_requester_pubkey(&self, requester_pubkey: impl Into) { + let requester_pubkey = requester_pubkey.into().to_ascii_lowercase(); + let mut requester_pubkeys = match self.requester_pubkeys.lock() { + Ok(requester_pubkeys) => requester_pubkeys, + Err(poisoned) => poisoned.into_inner(), + }; + if !requester_pubkeys.contains(&requester_pubkey) { + requester_pubkeys.push(requester_pubkey); + requester_pubkeys.sort(); + } + } + + fn requester_pubkeys(&self) -> Vec { + match self.requester_pubkeys.lock() { + Ok(requester_pubkeys) => requester_pubkeys.clone(), + Err(poisoned) => poisoned.into_inner().clone(), + } + } +} + +#[cfg(test)] +mod tests { + use super::*; + + #[test] + fn requester_routing_metadata_is_normalized_but_not_serialized() { + let context = context_for_turn(None, None, "turn-1".into(), "2026-07-24T10:00:00Z".into()) + .with_requester_pubkeys(vec!["BBBB".into(), "aaaa".into(), "AAAA".into()]); + let cloned_context = context.clone(); + cloned_context.add_requester_pubkey("CCCC"); + let observer = ObserverHandle::in_process(); + observer.emit("turn_started", Some(0), &context, serde_json::json!({})); + let event = observer.snapshot().pop().expect("observer event"); + + assert_eq!(event.requester_pubkeys, ["aaaa", "bbbb", "cccc"]); + let serialized = serde_json::to_value(&event).expect("serialize observer event"); + assert!( + serialized.get("requesterPubkeys").is_none(), + "requester routing metadata must stay outside encrypted payloads" + ); } } diff --git a/crates/buzz-acp/src/pool.rs b/crates/buzz-acp/src/pool.rs index cc537f8683..7dc0c8f61e 100644 --- a/crates/buzz-acp/src/pool.rs +++ b/crates/buzz-acp/src/pool.rs @@ -53,6 +53,8 @@ pub struct TaskMeta { pub channel_id: Option, /// Identifies terminal events when the task panics before returning a result. pub turn_id: String, + /// Authors whose triggering events started the turn. + pub requester_pubkeys: Vec, /// Clone of batch for Queue mode panic recovery. pub recoverable_batch: Option, /// Control signal for the in-flight prompt task. @@ -322,6 +324,9 @@ pub struct SteerRequest { /// `queue::native_steer_framing()` + `queue::format_event_block` so /// the wording cannot drift from the cancel+merge fallback path. pub prompt_blocks: Vec, + /// Author of the event carried by this steer. Added to the active turn's + /// observer recipients only after the agent accepts the native steer. + pub requester_pubkey: String, /// Oneshot for the read loop to report the outcome. pub ack_tx: tokio::sync::oneshot::Sender, } @@ -661,6 +666,24 @@ impl AgentPool { .map_err(|e| SteerError::Transport(e.to_string())) } + /// Add a successful native-steer author to the in-flight turn's panic + /// recovery recipients. + pub fn add_task_requester_pubkey(&mut self, channel_id: Uuid, requester_pubkey: &str) -> bool { + let Some(meta) = self + .task_map + .values_mut() + .find(|meta| meta.channel_id == Some(channel_id)) + else { + return false; + }; + let requester_pubkey = requester_pubkey.to_ascii_lowercase(); + if !meta.requester_pubkeys.contains(&requester_pubkey) { + meta.requester_pubkeys.push(requester_pubkey); + meta.requester_pubkeys.sort(); + } + true + } + pub fn result_tx(&self) -> mpsc::UnboundedSender { self.result_tx.clone() } @@ -1250,6 +1273,21 @@ fn send_prompt_result( }); } +pub(crate) fn batch_requester_pubkeys(batch: Option<&FlushBatch>) -> Vec { + let Some(batch) = batch else { + return Vec::new(); + }; + let mut pubkeys: Vec = batch + .events + .iter() + .chain(batch.cancelled_events.iter()) + .map(|batch_event| batch_event.event.pubkey.to_hex()) + .collect(); + pubkeys.sort(); + pubkeys.dedup(); + pubkeys +} + /// Core async function spawned for each prompt. /// /// Lifecycle: @@ -1281,12 +1319,16 @@ pub async fn run_prompt_task( PromptSource::Heartbeat => None, }; let turn_started_at = chrono::Utc::now().to_rfc3339(); - agent.acp.set_observer_context(observer::context_for_turn( + let initial_observer_context = observer::context_for_turn( observer_channel_id, None, turn_id.clone(), turn_started_at.clone(), - )); + ) + .with_requester_pubkeys(batch_requester_pubkeys(batch.as_ref())); + agent + .acp + .set_observer_context(initial_observer_context.clone()); let triggering_event_ids: Vec = batch .as_ref() .map(|b| b.events.iter().map(|be| be.event.id.to_hex()).collect()) @@ -1309,8 +1351,7 @@ pub async fn run_prompt_task( let _turn_guard = TurnCompletionGuard::new( agent.acp.observer_handle(), agent.acp.observer_agent_index(), - observer_channel_id, - turn_id.clone(), + initial_observer_context.clone(), ); // Start liveness with `turn_started`, not the final session/prompt call: @@ -1329,12 +1370,7 @@ pub async fn run_prompt_task( let liveness = run_turn_liveness( agent.acp.observer_handle(), agent.acp.observer_agent_index(), - observer::context_for_turn( - observer_channel_id, - None, - turn_id.clone(), - turn_started_at.clone(), - ), + initial_observer_context, ctx.turn_liveness_interval, Arc::clone(&liveness_state), ); @@ -1560,12 +1596,9 @@ pub async fn run_prompt_task( } } }; - agent.acp.set_observer_context(observer::context_for_turn( - observer_channel_id, - Some(session_id.clone()), - turn_id.clone(), - turn_started_at, - )); + let mut resolved_observer_context = agent.acp.observer_context().clone(); + resolved_observer_context.session_id = Some(session_id.clone()); + agent.acp.set_observer_context(resolved_observer_context); // Backfill liveness's shared session ID so ticks after this point carry // it too, matching every other observer frame for this turn. liveness_guard.set_session_id(session_id.clone()); @@ -3267,22 +3300,19 @@ impl Drop for LivenessGuard { struct TurnCompletionGuard { observer: Option, agent_index: Option, - channel_id: Option, - turn_id: String, + context: observer::ObserverContext, } impl TurnCompletionGuard { fn new( observer: Option, agent_index: Option, - channel_id: Option, - turn_id: String, + context: observer::ObserverContext, ) -> Self { Self { observer, agent_index, - channel_id, - turn_id, + context, } } } @@ -3290,11 +3320,10 @@ impl TurnCompletionGuard { impl Drop for TurnCompletionGuard { fn drop(&mut self) { if let Some(observer) = self.observer.take() { - let context = observer::context_for(self.channel_id, None, Some(self.turn_id.clone())); observer.emit( "turn_completed", self.agent_index, - &context, + &self.context, serde_json::json!({}), ); } @@ -4455,6 +4484,36 @@ mod tests { } } + #[test] + fn batch_requesters_include_current_and_cancelled_authors_once() { + let requester_a = Keys::generate(); + let requester_b = Keys::generate(); + let make_batch_event = |keys: &Keys, content: &str| crate::queue::BatchEvent { + event: EventBuilder::new(Kind::Custom(9), content) + .sign_with_keys(keys) + .unwrap(), + prompt_tag: "test".into(), + received_at: std::time::Instant::now(), + }; + let batch = FlushBatch { + channel_id: Uuid::new_v4(), + events: vec![ + make_batch_event(&requester_b, "new"), + make_batch_event(&requester_a, "also new"), + ], + cancelled_events: vec![make_batch_event(&requester_a, "cancelled")], + cancel_reason: Some(CancelReason::Steer), + }; + let mut expected = vec![ + requester_a.public_key().to_hex(), + requester_b.public_key().to_hex(), + ]; + expected.sort(); + + assert_eq!(batch_requester_pubkeys(Some(&batch)), expected); + assert!(batch_requester_pubkeys(None).is_empty()); + } + #[test] fn test_requeue_cancelled_batch_maps_control_signal_to_cancel_reason() { let cases = [ diff --git a/crates/buzz-core/src/observer.rs b/crates/buzz-core/src/observer.rs index 8347bda054..5a2949cf64 100644 --- a/crates/buzz-core/src/observer.rs +++ b/crates/buzz-core/src/observer.rs @@ -1,8 +1,9 @@ //! Agent observer frame helpers. //! -//! Observer frames are transient, owner-scoped agent telemetry/control messages. -//! They use a Buzz ephemeral event kind and carry NIP-44 encrypted JSON in the -//! event content so relays can route frames without reading ACP internals. +//! Observer frames are transient, recipient-scoped agent telemetry/control +//! messages. They use a Buzz ephemeral event kind and carry NIP-44 encrypted +//! JSON so relays can route frames without reading ACP internals. Telemetry may +//! target an agent owner or a turn requester; controls remain owner-only. use nostr::{nips::nip44, Event, Keys, PublicKey}; use serde::{de::DeserializeOwned, Serialize}; @@ -13,7 +14,7 @@ use zeroize::Zeroize; pub const OBSERVER_AGENT_TAG: &str = "agent"; /// Tag name that identifies the cleartext frame direction. pub const OBSERVER_FRAME_TAG: &str = "frame"; -/// Frame value for agent-to-owner observer telemetry. +/// Frame value for agent-to-recipient observer telemetry. pub const OBSERVER_FRAME_TELEMETRY: &str = "telemetry"; /// Frame value for owner-to-agent observer control commands. pub const OBSERVER_FRAME_CONTROL: &str = "control"; diff --git a/crates/buzz-relay/src/api/mod.rs b/crates/buzz-relay/src/api/mod.rs index d9f829433b..d0ab20dc7b 100644 --- a/crates/buzz-relay/src/api/mod.rs +++ b/crates/buzz-relay/src/api/mod.rs @@ -227,6 +227,9 @@ pub mod relay_members { ), true, ); + state + .observer_agent_cache + .insert((tenant.community(), agent.to_bytes().to_vec()), true); } materialized } diff --git a/crates/buzz-relay/src/handlers/event.rs b/crates/buzz-relay/src/handlers/event.rs index ecb3e41fcc..b7a5cfd1f1 100644 --- a/crates/buzz-relay/src/handlers/event.rs +++ b/crates/buzz-relay/src/handlers/event.rs @@ -882,7 +882,7 @@ enum AgentObserverDirection { #[derive(Debug, Clone, Copy)] struct AgentObserverRoute { agent: PublicKey, - owner: PublicKey, + recipient: PublicKey, direction: AgentObserverDirection, } @@ -915,8 +915,10 @@ fn observer_frame_rate_limited( /// Handle encrypted agent observer frames (kind 24200). /// /// These frames bypass storage and are routed as global ephemeral events. The -/// relay gates publication by the existing `agent_owner_pubkey` mapping and -/// gates subscription in the REQ handler via the cleartext `p` tag. +/// relay requires telemetry authors to be registered NIP-OA agents and +/// requires control authors to be the registered owner. Subscription remains +/// gated by the cleartext `p` tag, so only the encrypted recipient receives +/// each frame. async fn handle_agent_observer_event( event: Event, conn_id: uuid::Uuid, @@ -972,60 +974,101 @@ async fn handle_agent_observer_event( } }; - // Fast path: if this connection authenticated via NIP-OA and the verified - // owner matches the observer frame's target owner, skip the DB lookup entirely. - let session_owner_match = { - let auth = conn.auth_state.read().await; - if let crate::connection::AuthState::Authenticated(ctx) = &*auth { - ctx.agent_owner_pubkey.as_ref() == Some(&route.owner) - } else { - false - } - }; - let agent_bytes = route.agent.to_bytes().to_vec(); - let owner_bytes = route.owner.to_bytes().to_vec(); - let cache_key = ( - conn.tenant.community(), - agent_bytes.clone(), - owner_bytes.clone(), - ); - let is_owner = if session_owner_match { - true - } else { - match state.observer_owner_cache.get(&cache_key) { - Some(cached) => cached, - None => { - let result = state - .db - .is_agent_owner(conn.tenant.community(), &agent_bytes, &owner_bytes) - .await; - match result { - Ok(v) => { - state.observer_owner_cache.insert(cache_key, v); - v + match route.direction { + AgentObserverDirection::Telemetry => { + // NIP-OA authentication already proves this event author is a + // registered agent. On cache misses (including direct membership + // auth), confirm the immutable agent-owner mapping in the DB. + let session_agent_match = { + let auth = conn.auth_state.read().await; + matches!( + &*auth, + crate::connection::AuthState::Authenticated(ctx) + if ctx.pubkey == route.agent && ctx.agent_owner_pubkey.is_some() + ) + }; + let cache_key = (conn.tenant.community(), agent_bytes.clone()); + let is_registered_agent = if session_agent_match { + true + } else { + match state.observer_agent_cache.get(&cache_key) { + Some(cached) => cached, + None => { + let result = state + .db + .get_agent_channel_policy(conn.tenant.community(), &agent_bytes) + .await; + match result { + Ok(policy) => { + let registered = policy.is_some_and(|(_, owner)| owner.is_some()); + state.observer_agent_cache.insert(cache_key, registered); + registered + } + Err(e) => { + warn!(conn_id = %conn_id, event_id = %event_id_hex, "agent observer registration check failed: {e}"); + conn.send(RelayMessage::ok( + event_id_hex, + false, + "error: internal server error", + )); + return; + } + } } - Err(e) => { - warn!(conn_id = %conn_id, event_id = %event_id_hex, "agent observer owner check failed: {e}"); - conn.send(RelayMessage::ok( - event_id_hex, - false, - "error: internal server error", - )); - return; + } + }; + if !is_registered_agent { + reject("auth"); + conn.send(RelayMessage::ok( + event_id_hex, + false, + "restricted: observer telemetry author is not a registered agent", + )); + return; + } + } + AgentObserverDirection::Control => { + let owner_bytes = event.pubkey.to_bytes().to_vec(); + let cache_key = ( + conn.tenant.community(), + agent_bytes.clone(), + owner_bytes.clone(), + ); + let is_owner = match state.observer_owner_cache.get(&cache_key) { + Some(cached) => cached, + None => { + let result = state + .db + .is_agent_owner(conn.tenant.community(), &agent_bytes, &owner_bytes) + .await; + match result { + Ok(v) => { + state.observer_owner_cache.insert(cache_key, v); + v + } + Err(e) => { + warn!(conn_id = %conn_id, event_id = %event_id_hex, "agent observer owner check failed: {e}"); + conn.send(RelayMessage::ok( + event_id_hex, + false, + "error: internal server error", + )); + return; + } } } + }; + if !is_owner { + reject("auth"); + conn.send(RelayMessage::ok( + event_id_hex, + false, + "restricted: observer control is not authorized for this agent owner", + )); + return; } } - }; - if !is_owner { - reject("auth"); - conn.send(RelayMessage::ok( - event_id_hex, - false, - "restricted: observer frame is not authorized for this agent owner", - )); - return; } // Rate limit telemetry frames only (100/sec per agent). @@ -1059,7 +1102,7 @@ async fn handle_agent_observer_event( debug!( event_id = %event_id_hex, agent = %route.agent.to_hex(), - owner = %route.owner.to_hex(), + recipient = %route.recipient.to_hex(), direction = ?route.direction, "Agent observer fan-out" ); @@ -1077,21 +1120,13 @@ fn agent_observer_route(event: &Event) -> Result, Str let agent = parse_single_pubkey_tag(event, OBSERVER_AGENT_TAG)?; let frame = single_tag_content(event, OBSERVER_FRAME_TAG)?; - let (owner, direction, expected_frame) = if event.pubkey == agent && recipient != agent { - ( - recipient, - AgentObserverDirection::Telemetry, - OBSERVER_FRAME_TELEMETRY, - ) + let (direction, expected_frame) = if event.pubkey == agent && recipient != agent { + (AgentObserverDirection::Telemetry, OBSERVER_FRAME_TELEMETRY) } else if recipient == agent && event.pubkey != agent { - ( - event.pubkey, - AgentObserverDirection::Control, - OBSERVER_FRAME_CONTROL, - ) + (AgentObserverDirection::Control, OBSERVER_FRAME_CONTROL) } else { return Err( - "invalid: observer frame must be agent-to-owner telemetry or owner-to-agent control" + "invalid: observer frame must be agent-to-recipient telemetry or owner-to-agent control" .into(), ); }; @@ -1103,7 +1138,7 @@ fn agent_observer_route(event: &Event) -> Result, Str Ok(Some(AgentObserverRoute { agent, - owner, + recipient, direction, })) } @@ -1216,18 +1251,18 @@ mod tests { } #[test] - fn agent_observer_route_accepts_agent_to_owner_telemetry() { + fn agent_observer_route_accepts_agent_to_recipient_telemetry() { let agent = Keys::generate(); - let owner = Keys::generate(); + let recipient = Keys::generate(); let encrypted = encrypt_observer_payload( &agent, - &owner.public_key(), + &recipient.public_key(), &serde_json::json!({"type": "acp_read"}), ) .expect("encrypt observer payload"); let event = EventBuilder::new(Kind::Custom(KIND_AGENT_OBSERVER_FRAME as u16), encrypted) .tags([ - Tag::parse(["p", &owner.public_key().to_hex()]).expect("p tag"), + Tag::parse(["p", &recipient.public_key().to_hex()]).expect("p tag"), Tag::parse([OBSERVER_AGENT_TAG, &agent.public_key().to_hex()]).expect("agent tag"), Tag::parse([OBSERVER_FRAME_TAG, OBSERVER_FRAME_TELEMETRY]).expect("frame tag"), ]) @@ -1238,7 +1273,7 @@ mod tests { .expect("observer route") .expect("route should be Some"); assert_eq!(route.agent, agent.public_key()); - assert_eq!(route.owner, owner.public_key()); + assert_eq!(route.recipient, recipient.public_key()); assert_eq!(route.direction, super::AgentObserverDirection::Telemetry); } @@ -1265,7 +1300,7 @@ mod tests { .expect("observer route") .expect("route should be Some"); assert_eq!(route.agent, agent.public_key()); - assert_eq!(route.owner, owner.public_key()); + assert_eq!(route.recipient, agent.public_key()); assert_eq!(route.direction, super::AgentObserverDirection::Control); } @@ -1289,6 +1324,138 @@ mod tests { assert!(err.contains("NIP-44")); } + fn observer_test_connection( + community: buzz_core::CommunityId, + pubkey: nostr::PublicKey, + agent_owner_pubkey: Option, + ) -> ( + Arc, + mpsc::Receiver, + ) { + let (send_tx, send_rx) = mpsc::channel(4); + let (ctrl_tx, _ctrl_rx) = mpsc::channel(1); + let conn = Arc::new(crate::connection::ConnectionState { + conn_id: Uuid::new_v4(), + tenant: buzz_core::TenantContext::resolved(community, "observer.example"), + remote_addr: "127.0.0.1:1234".parse().expect("socket addr"), + auth_state: RwLock::new(crate::connection::AuthState::Authenticated( + buzz_auth::AuthContext { + pubkey, + scopes: vec![], + channel_ids: None, + auth_method: buzz_auth::AuthMethod::Nip42, + agent_owner_pubkey, + }, + )), + subscriptions: Arc::new(Mutex::new(HashMap::new())), + send_tx, + ctrl_tx, + cancel: CancellationToken::new(), + backpressure_count: Arc::new(AtomicU8::new(0)), + grace_limit: 3, + }); + (conn, send_rx) + } + + #[tokio::test] + async fn registered_agent_telemetry_can_target_a_non_owner_recipient() { + let state = fanout_access::test_state().await; + let agent = Keys::generate(); + let owner = Keys::generate(); + let requester = Keys::generate(); + let community = buzz_core::CommunityId::from_uuid(Uuid::new_v4()); + state + .observer_agent_cache + .insert((community, agent.public_key().to_bytes().to_vec()), true); + let encrypted = encrypt_observer_payload( + &agent, + &requester.public_key(), + &serde_json::json!({"type": "turn_started"}), + ) + .expect("encrypt observer payload"); + let event = EventBuilder::new(Kind::Custom(KIND_AGENT_OBSERVER_FRAME as u16), encrypted) + .tags([ + Tag::parse(["p", &requester.public_key().to_hex()]).expect("p tag"), + Tag::parse([OBSERVER_AGENT_TAG, &agent.public_key().to_hex()]).expect("agent tag"), + Tag::parse([OBSERVER_FRAME_TAG, OBSERVER_FRAME_TELEMETRY]).expect("frame tag"), + ]) + .sign_with_keys(&agent) + .expect("sign event"); + let (conn, mut send_rx) = + observer_test_connection(community, agent.public_key(), Some(owner.public_key())); + + super::handle_agent_observer_event( + event.clone(), + conn.conn_id, + &event.id.to_hex(), + conn, + state, + ) + .await; + + let axum::extract::ws::Message::Text(text) = + send_rx.try_recv().expect("observer acceptance sent") + else { + panic!("expected text relay message"); + }; + let frame: serde_json::Value = serde_json::from_str(&text).expect("relay frame JSON"); + assert_eq!(frame[0], "OK"); + assert_eq!(frame[2], true); + } + + #[tokio::test] + async fn non_owner_control_stays_rejected() { + let state = fanout_access::test_state().await; + let agent = Keys::generate(); + let attacker = Keys::generate(); + let community = buzz_core::CommunityId::from_uuid(Uuid::new_v4()); + state.observer_owner_cache.insert( + ( + community, + agent.public_key().to_bytes().to_vec(), + attacker.public_key().to_bytes().to_vec(), + ), + false, + ); + let encrypted = encrypt_observer_payload( + &attacker, + &agent.public_key(), + &serde_json::json!({"type": "cancel_turn"}), + ) + .expect("encrypt observer payload"); + let event = EventBuilder::new(Kind::Custom(KIND_AGENT_OBSERVER_FRAME as u16), encrypted) + .tags([ + Tag::parse(["p", &agent.public_key().to_hex()]).expect("p tag"), + Tag::parse([OBSERVER_AGENT_TAG, &agent.public_key().to_hex()]).expect("agent tag"), + Tag::parse([OBSERVER_FRAME_TAG, OBSERVER_FRAME_CONTROL]).expect("frame tag"), + ]) + .sign_with_keys(&attacker) + .expect("sign event"); + let (conn, mut send_rx) = observer_test_connection(community, attacker.public_key(), None); + + super::handle_agent_observer_event( + event.clone(), + conn.conn_id, + &event.id.to_hex(), + conn, + state, + ) + .await; + + let axum::extract::ws::Message::Text(text) = + send_rx.try_recv().expect("observer rejection sent") + else { + panic!("expected text relay message"); + }; + let frame: serde_json::Value = serde_json::from_str(&text).expect("relay frame JSON"); + assert_eq!(frame[0], "OK"); + assert_eq!(frame[2], false); + assert_eq!( + frame[3], + "restricted: observer control is not authorized for this agent owner" + ); + } + #[tokio::test] async fn observer_frame_rate_limiter_is_scoped_by_community() { let state = fanout_access::test_state().await; @@ -1315,32 +1482,27 @@ mod tests { } #[tokio::test] - async fn observer_owner_cache_is_scoped_to_community() { + async fn observer_agent_cache_is_scoped_to_community() { let state = fanout_access::test_state().await; let agent = Keys::generate(); let owner = Keys::generate(); let agent_bytes = agent.public_key().to_bytes().to_vec(); - let owner_bytes = owner.public_key().to_bytes().to_vec(); let community_a = buzz_core::CommunityId::from_uuid(Uuid::new_v4()); let community_b = buzz_core::CommunityId::from_uuid(Uuid::new_v4()); - state.observer_owner_cache.insert( - (community_a, agent_bytes.clone(), owner_bytes.clone()), - true, - ); + state + .observer_agent_cache + .insert((community_a, agent_bytes.clone()), true); assert_eq!( - state.observer_owner_cache.get(&( - community_b, - agent_bytes.clone(), - owner_bytes.clone() - )), + state + .observer_agent_cache + .get(&(community_b, agent_bytes.clone())), None, "A cached allow must not populate B's observer authorization key" ); - state.observer_owner_cache.insert( - (community_b, agent_bytes.clone(), owner_bytes.clone()), - false, - ); + state + .observer_agent_cache + .insert((community_b, agent_bytes.clone()), false); let encrypted = encrypt_observer_payload( &agent, @@ -1399,7 +1561,7 @@ mod tests { assert_eq!(frame[2], false); assert_eq!( frame[3], - "restricted: observer frame is not authorized for this agent owner" + "restricted: observer telemetry author is not a registered agent" ); } diff --git a/crates/buzz-relay/src/state.rs b/crates/buzz-relay/src/state.rs index 758c001b96..265a1968ab 100644 --- a/crates/buzz-relay/src/state.rs +++ b/crates/buzz-relay/src/state.rs @@ -605,6 +605,13 @@ pub struct AppState { /// Prevents repeated DB lookups from bursty observer traffic. #[allow(clippy::type_complexity)] pub observer_owner_cache: Arc, Vec), bool>>, + /// Cache for registered-agent authorization on requester-addressed + /// observer telemetry. Key: (community_id, agent_pubkey_bytes). + /// + /// Telemetry recipients are intentionally not part of this decision: a + /// registered agent may encrypt its own turn activity to the requester, + /// while control authorization continues to use `observer_owner_cache`. + pub observer_agent_cache: Arc), bool>>, /// Cache for the `author_type` metric label on the ingest path. /// Key: (community_id, author pubkey bytes). Value: is_agent /// (`users.agent_owner_pubkey IS NOT NULL`). The mapping is @@ -785,6 +792,12 @@ impl AppState { .time_to_live(std::time::Duration::from_secs(300)) .build(), ), + observer_agent_cache: Arc::new( + moka::sync::Cache::builder() + .max_capacity(1_000) + .time_to_live(std::time::Duration::from_secs(300)) + .build(), + ), author_type_cache: Arc::new( moka::sync::Cache::builder() .max_capacity(10_000) diff --git a/crates/buzz-sdk/src/builders.rs b/crates/buzz-sdk/src/builders.rs index f9e54de9c5..d06f2d0e73 100644 --- a/crates/buzz-sdk/src/builders.rs +++ b/crates/buzz-sdk/src/builders.rs @@ -239,9 +239,10 @@ pub fn build_message( /// Build an encrypted agent observer frame (kind 24200). /// -/// `recipient_pubkey` is the cleartext `p` tag used by the relay for owner-only -/// routing. `agent_pubkey` identifies the managed agent whose observer stream -/// this frame belongs to. `encrypted_content` must be NIP-44 v2 ciphertext. +/// `recipient_pubkey` is the cleartext `p` tag used by the relay for encrypted +/// recipient routing. `agent_pubkey` identifies the managed agent whose +/// observer stream this frame belongs to. `encrypted_content` must be NIP-44 +/// v2 ciphertext. pub fn build_agent_observer_frame( recipient_pubkey: &str, agent_pubkey: &str, diff --git a/desktop/src/features/agents/ingestArchivedObserverEvents.test.mjs b/desktop/src/features/agents/ingestArchivedObserverEvents.test.mjs index c44c2a88e0..167dcf0bed 100644 --- a/desktop/src/features/agents/ingestArchivedObserverEvents.test.mjs +++ b/desktop/src/features/agents/ingestArchivedObserverEvents.test.mjs @@ -163,6 +163,42 @@ describe("ingestArchivedObserverEvents", () => { ); }); + it("drops owner-only frames for requester-observable agents", async () => { + _testRegisterKnownAgents(SUB_ID, [AGENT_PUBKEY], []); + const ownerOnly = makeObserverEvent({ + kind: "control_result", + payload: { type: "cancel_turn", status: "sent" }, + }); + + await ingestArchivedObserverEvents( + [makeRawEvent()], + makeDecrypt(ownerOnly), + ); + + assert.equal( + _testGetArchivedChannelEvents(AGENT_PUBKEY, "chan-1").length, + 0, + ); + }); + + it("keeps owner-only frames for owner-authorized agents", async () => { + _testRegisterKnownAgents(SUB_ID, [AGENT_PUBKEY], [AGENT_PUBKEY]); + const ownerOnly = makeObserverEvent({ + kind: "control_result", + payload: { type: "cancel_turn", status: "sent" }, + }); + + await ingestArchivedObserverEvents( + [makeRawEvent()], + makeDecrypt(ownerOnly), + ); + + assert.equal( + _testGetArchivedChannelEvents(AGENT_PUBKEY, "chan-1").length, + 1, + ); + }); + it("test_dedup_does_not_add_live_present_event", async () => { _testRegisterKnownAgents(SUB_ID, [AGENT_PUBKEY]); // Ingest the same archived event twice — the channel archive window must dedup diff --git a/desktop/src/features/agents/observerRelayStore.ts b/desktop/src/features/agents/observerRelayStore.ts index 56c69f915a..d072c6363e 100644 --- a/desktop/src/features/agents/observerRelayStore.ts +++ b/desktop/src/features/agents/observerRelayStore.ts @@ -115,6 +115,8 @@ const agentManagementListeners = new Set< // recompute the union, so co-mounted callers no longer clobber each other. const knownAgentPubkeys = new Set(); const knownAgentsBySubscription = new Map>(); +const ownerFrameAgentPubkeys = new Set(); +const ownerFrameAgentsBySubscription = new Map>(); const pendingUnknownAgentFrames: RelayEvent[] = []; // Callback invoked when session_config_captured is received, so React Query @@ -130,21 +132,32 @@ export function setSessionConfigCapturedCallback( function recomputeKnownAgentPubkeys() { knownAgentPubkeys.clear(); + ownerFrameAgentPubkeys.clear(); for (const subscriptionAgents of knownAgentsBySubscription.values()) { for (const pubkey of subscriptionAgents) { knownAgentPubkeys.add(pubkey); } } + for (const subscriptionAgents of ownerFrameAgentsBySubscription.values()) { + for (const pubkey of subscriptionAgents) { + ownerFrameAgentPubkeys.add(pubkey); + } + } } function registerKnownAgents( subscriptionId: string, pubkeys: readonly string[], + ownerFramePubkeys: readonly string[], ) { knownAgentsBySubscription.set( subscriptionId, new Set(pubkeys.map((pubkey) => normalizePubkey(pubkey))), ); + ownerFrameAgentsBySubscription.set( + subscriptionId, + new Set(ownerFramePubkeys.map((pubkey) => normalizePubkey(pubkey))), + ); recomputeKnownAgentPubkeys(); if (knownAgentPubkeys.size > 0 && pendingUnknownAgentFrames.length > 0) { const pending = pendingUnknownAgentFrames.splice(0); @@ -157,7 +170,10 @@ function registerKnownAgents( } function unregisterKnownAgents(subscriptionId: string) { - if (knownAgentsBySubscription.delete(subscriptionId)) { + const removedKnown = knownAgentsBySubscription.delete(subscriptionId); + const removedOwnerFrames = + ownerFrameAgentsBySubscription.delete(subscriptionId); + if (removedKnown || removedOwnerFrames) { recomputeKnownAgentPubkeys(); } } @@ -179,6 +195,15 @@ function invalidateSnapshot(key: string) { snapshotByAgent.delete(key); } +function isOwnerOnlyObserverFrame(parsed: ObserverEvent): boolean { + return ( + parsed.kind === "session_config_captured" || + parsed.kind === "control_result" || + parsed.kind === "managed_agent_runtime_lifecycle" || + parseAgentManagementRequest(parsed.payload) !== null + ); +} + function setConnectionState( nextState: ConnectionState, nextErrorMessage: string | null = errorMessage, @@ -372,6 +397,14 @@ async function handleRelayObserverEvent( if (activeGeneration !== generation) { return; } + const normalizedAgentPubkey = normalizePubkey(agentPubkey); + const managementRequest = parseAgentManagementRequest(parsed.payload); + if ( + isOwnerOnlyObserverFrame(parsed) && + !ownerFrameAgentPubkeys.has(normalizedAgentPubkey) + ) { + return; + } // Track the latest-live-session-id per (agent, channel) on the live path. // Only set when the parsed event carries both a sessionId and channelId, // so we never attribute a session to the wrong channel. @@ -392,7 +425,6 @@ async function handleRelayObserverEvent( } } appendAgentEvent(agentPubkey, parsed); - const managementRequest = parseAgentManagementRequest(parsed.payload); if (managementRequest) { for (const listener of agentManagementListeners) { listener(agentPubkey, managementRequest); @@ -598,6 +630,7 @@ export function shouldObserveManagedAgents( export function useManagedAgentObserverBridge( agents: readonly Pick[], + ownerFramePubkeys: readonly string[], ) { const subscriptionId = React.useId(); const hasManagedAgent = shouldObserveManagedAgents(agents); @@ -606,16 +639,25 @@ export function useManagedAgentObserverBridge( () => agents.map((agent) => agent.pubkey), [agents], ); + const normalizedOwnerFramePubkeys = React.useMemo( + () => ownerFramePubkeys.map((pubkey) => normalizePubkey(pubkey)), + [ownerFramePubkeys], + ); // Keep this subscriber's slice of the trusted-pubkey set in sync with its - // own agent list. The store recomputes the union across all subscribers, so - // a co-mounted caller no longer wipes out this caller's agents. + // own agent list and owner-authorized subset. The store recomputes the union + // across all subscribers, so a co-mounted caller no longer wipes out this + // caller's agents or widens owner-only side effects. React.useEffect(() => { - registerKnownAgents(subscriptionId, agentPubkeys); + registerKnownAgents( + subscriptionId, + agentPubkeys, + normalizedOwnerFramePubkeys, + ); return () => { unregisterKnownAgents(subscriptionId); }; - }, [subscriptionId, agentPubkeys]); + }, [subscriptionId, agentPubkeys, normalizedOwnerFramePubkeys]); React.useEffect(() => { if (!hasManagedAgent) { @@ -676,6 +718,12 @@ export async function ingestArchivedObserverEvents( } try { const parsed = (await _decryptFn(event)) as ObserverEvent; + if ( + isOwnerOnlyObserverFrame(parsed) && + !ownerFrameAgentPubkeys.has(normalizePubkey(agentPubkey)) + ) { + continue; + } // Route archived events to the channel-scoped archive window (no cap) // rather than the per-agent live-relay store (MAX_OBSERVER_EVENTS cap). // Events without a channelId fall through to the live store so they @@ -747,6 +795,8 @@ export function resetAgentObserverStore() { archiveEventsByChannel.clear(); knownAgentPubkeys.clear(); knownAgentsBySubscription.clear(); + ownerFrameAgentPubkeys.clear(); + ownerFrameAgentsBySubscription.clear(); pendingUnknownAgentFrames.length = 0; latestLiveSessionByAgentChannel.clear(); agentManagementListeners.clear(); @@ -765,8 +815,9 @@ export function resetAgentObserverStore() { export function _testRegisterKnownAgents( subscriptionId: string, pubkeys: readonly string[], + ownerFramePubkeys: readonly string[] = pubkeys, ): void { - registerKnownAgents(subscriptionId, pubkeys); + registerKnownAgents(subscriptionId, pubkeys, ownerFramePubkeys); } /** diff --git a/desktop/src/features/agents/useAgentObserverIngestion.test.mjs b/desktop/src/features/agents/useAgentObserverIngestion.test.mjs index 4a20be9be3..8e49b646cc 100644 --- a/desktop/src/features/agents/useAgentObserverIngestion.test.mjs +++ b/desktop/src/features/agents/useAgentObserverIngestion.test.mjs @@ -21,7 +21,13 @@ describe("combineObserverIngestionAgents", () => { new Map(), ME, ); - assert.deepEqual(result, [{ pubkey: AGENT_LOCAL, status: "running" }]); + assert.deepEqual(result, [ + { + pubkey: AGENT_LOCAL, + status: "running", + canProcessOwnerFrames: true, + }, + ]); }); it("adds declared-owned relay agents as deployed", () => { @@ -31,27 +37,45 @@ describe("combineObserverIngestionAgents", () => { new Map([[AGENT_REMOTE, ME]]), ME, ); - assert.deepEqual(result, [{ pubkey: AGENT_REMOTE, status: "deployed" }]); + assert.deepEqual(result, [ + { + pubkey: AGENT_REMOTE, + status: "deployed", + canProcessOwnerFrames: true, + }, + ]); }); - it("excludes relay agents owned by someone else", () => { + it("includes relay agents owned by someone else without owner-frame privileges", () => { const result = combineObserverIngestionAgents( [], [AGENT_FOREIGN], new Map([[AGENT_FOREIGN, OTHER]]), ME, ); - assert.deepEqual(result, []); + assert.deepEqual(result, [ + { + pubkey: AGENT_FOREIGN, + status: "deployed", + canProcessOwnerFrames: false, + }, + ]); }); - it("excludes relay agents with no declared owner", () => { + it("includes relay agents with no declared owner for requester telemetry", () => { const result = combineObserverIngestionAgents( [], [AGENT_REMOTE], new Map(), ME, ); - assert.deepEqual(result, []); + assert.deepEqual(result, [ + { + pubkey: AGENT_REMOTE, + status: "deployed", + canProcessOwnerFrames: false, + }, + ]); }); it("does not duplicate an agent that is both managed and on the relay", () => { @@ -61,7 +85,13 @@ describe("combineObserverIngestionAgents", () => { new Map([[AGENT_LOCAL, ME]]), ME, ); - assert.deepEqual(result, [{ pubkey: AGENT_LOCAL, status: "stopped" }]); + assert.deepEqual(result, [ + { + pubkey: AGENT_LOCAL, + status: "stopped", + canProcessOwnerFrames: true, + }, + ]); }); it("matches ownership case-insensitively", () => { @@ -72,17 +102,32 @@ describe("combineObserverIngestionAgents", () => { ME, ); assert.deepEqual(result, [ - { pubkey: AGENT_REMOTE.toUpperCase(), status: "deployed" }, + { + pubkey: AGENT_REMOTE.toUpperCase(), + status: "deployed", + canProcessOwnerFrames: true, + }, ]); }); - it("returns only managed agents when identity is not resolved yet", () => { + it("registers relay agents without owner privileges before identity resolves", () => { const result = combineObserverIngestionAgents( [{ pubkey: AGENT_LOCAL, status: "running" }], [AGENT_REMOTE], new Map([[AGENT_REMOTE, ME]]), undefined, ); - assert.deepEqual(result, [{ pubkey: AGENT_LOCAL, status: "running" }]); + assert.deepEqual(result, [ + { + pubkey: AGENT_LOCAL, + status: "running", + canProcessOwnerFrames: true, + }, + { + pubkey: AGENT_REMOTE, + status: "deployed", + canProcessOwnerFrames: false, + }, + ]); }); }); diff --git a/desktop/src/features/agents/useAgentObserverIngestion.ts b/desktop/src/features/agents/useAgentObserverIngestion.ts index 386b762142..31dd03c866 100644 --- a/desktop/src/features/agents/useAgentObserverIngestion.ts +++ b/desktop/src/features/agents/useAgentObserverIngestion.ts @@ -11,20 +11,22 @@ import { useIdentityQuery } from "@/shared/api/hooks"; import type { ManagedAgent } from "@/shared/api/types"; import { normalizePubkey } from "@/shared/lib/pubkey"; -type IngestionAgent = Pick; +type IngestionAgent = Pick & { + canProcessOwnerFrames: boolean; +}; /** - * Combine locally managed agents with relay agents the current identity - * declared-owns (NIP-OA `ownerPubkey == me`) into one ingestion list. + * Combine locally managed agents with every directory-listed relay agent into + * one ingestion list. * - * Managed agents keep their real status; owned relay agents that are not - * managed locally are treated as `deployed` so the observer subscription - * starts and their frames decrypt. Registering non-owned agents would be - * pointless — observer frames are `#p`-addressed to the owner, so frames for - * agents we do not own never arrive on our subscription in the first place. + * Managed agents keep their real status and owner-frame privileges. Relay + * agents are treated as `deployed` so requester-addressed turn frames can be + * decrypted even when the current identity is not their owner. Declared + * ownership is retained separately to gate control/config/lifecycle side + * effects in the observer store. */ export function combineObserverIngestionAgents( - managedAgents: readonly IngestionAgent[], + managedAgents: readonly Pick[], relayAgentPubkeys: readonly string[], ownerByPubkey: ReadonlyMap, currentPubkey: string | null | undefined, @@ -32,27 +34,29 @@ export function combineObserverIngestionAgents( const managed = managedAgents.map((agent) => ({ pubkey: agent.pubkey, status: agent.status, + canProcessOwnerFrames: true, })); - if (!currentPubkey) { - return managed; - } const managedSet = new Set( managed.map((agent) => normalizePubkey(agent.pubkey)), ); - const me = normalizePubkey(currentPubkey); - const owned: IngestionAgent[] = []; + const me = currentPubkey ? normalizePubkey(currentPubkey) : null; + const relay: IngestionAgent[] = []; for (const pubkey of relayAgentPubkeys) { const key = normalizePubkey(pubkey); if (managedSet.has(key)) { continue; } const owner = ownerByPubkey.get(key); - if (owner && normalizePubkey(owner) === me) { - owned.push({ pubkey, status: "deployed" as const }); - } + relay.push({ + pubkey, + status: "deployed" as const, + canProcessOwnerFrames: Boolean( + me && owner && normalizePubkey(owner) === me, + ), + }); } - return [...managed, ...owned]; + return [...managed, ...relay]; } /** @@ -63,13 +67,13 @@ export function combineObserverIngestionAgents( * which screen or panel happens to be open. Individual surfaces read from the * stores; none of them need to mount their own bridge for ingestion to work. * - * This is the product invariant: if the current identity owns an agent (local - * managed agent or declared-owned relay agent), its turn activity is ingested - * app-wide — not only while a panel that happens to mount a bridge is open. + * This is the product invariant: owner and requester-addressed turn activity + * is ingested app-wide, not only while a panel that happens to mount a bridge + * is open. Owner-only side effects remain separately gated. * * Mounts before identity resolves by design: while `currentPubkey` is still - * `undefined`, `combineObserverIngestionAgents` returns managed agents only, - * and relay-owned agents are folded in on the render after identity arrives. + * `undefined`, relay agents are already registered for requester telemetry; + * ownership privileges are folded in after identity and profiles resolve. * Do not gate this hook on identity/startup readiness — that would drop * managed-agent observer coverage during startup. */ @@ -111,6 +115,14 @@ export function useAgentObserverIngestion() { ); }, [currentPubkey, managedAgents, profiles, relayAgentPubkeys]); - useManagedAgentObserverBridge(ingestionAgents); + const ownerFrameAgentPubkeys = React.useMemo( + () => + ingestionAgents + .filter((agent) => agent.canProcessOwnerFrames) + .map((agent) => agent.pubkey), + [ingestionAgents], + ); + + useManagedAgentObserverBridge(ingestionAgents, ownerFrameAgentPubkeys); useActiveAgentTurnsBridge(ingestionAgents); } diff --git a/desktop/src/features/profile/lib/profileActivityAgent.ts b/desktop/src/features/profile/lib/profileActivityAgent.ts index 2bc4fe4465..21d78bf922 100644 --- a/desktop/src/features/profile/lib/profileActivityAgent.ts +++ b/desktop/src/features/profile/lib/profileActivityAgent.ts @@ -13,14 +13,14 @@ export function resolveProfileActivityAgent({ managedAgent, profile, relayAgent, - viewerIsOwner, + viewerCanObserve, }: { effectivePubkey: string | null; isBot: boolean; managedAgent: ManagedAgent | undefined; profile: { avatarUrl?: string | null; displayName?: string | null } | null; relayAgent: RelayAgent | undefined; - viewerIsOwner: boolean; + viewerCanObserve: boolean; }): ProfileActivityAgent | null { if (managedAgent) { return { @@ -31,7 +31,7 @@ export function resolveProfileActivityAgent({ }; } - if (!viewerIsOwner || !effectivePubkey || !isBot) { + if (!viewerCanObserve || !effectivePubkey || !isBot) { return null; } diff --git a/desktop/src/features/profile/ui/UserProfilePanel.tsx b/desktop/src/features/profile/ui/UserProfilePanel.tsx index c1f713d230..998696bc1e 100644 --- a/desktop/src/features/profile/ui/UserProfilePanel.tsx +++ b/desktop/src/features/profile/ui/UserProfilePanel.tsx @@ -6,6 +6,7 @@ import { useAgentMemoryQuery, useIsManagedAgent, } from "@/features/agent-memory/hooks"; +import { useActiveAgentTurns } from "@/features/agents/activeAgentTurnsStore"; import { type AttachManagedAgentToChannelResult, useAcpRuntimesQuery, @@ -308,6 +309,8 @@ export function UserProfilePanel({ // manage it locally (older agents may not advertise an owner pubkey). Every // real boundary is server-side, so this only controls what UI we paint. const viewerIsOwner = isCurrentUserOwner || isOwner === true; + const observableTurns = useActiveAgentTurns(isBot ? effectivePubkey : null); + const viewerCanObserve = viewerIsOwner || observableTurns.length > 0; const activityAgent = React.useMemo( () => @@ -317,13 +320,20 @@ export function UserProfilePanel({ managedAgent, profile: profile ?? null, relayAgent, - viewerIsOwner, + viewerCanObserve, }), - [effectivePubkey, isBot, managedAgent, profile, relayAgent, viewerIsOwner], + [ + effectivePubkey, + isBot, + managedAgent, + profile, + relayAgent, + viewerCanObserve, + ], ); - // Observer ingestion (frame decryption + derived active-turn liveness) is - // owner-global — mounted once in AppShell via useAgentObserverIngestion — - // covering both locally managed agents and declared-owned relay agents. + // Observer ingestion is app-global. Owners see all turns, while a requester + // can open a shared agent only while requester-addressed activity proves an + // observable turn is live. const canEditAgent = isOwner === true && (managedAgent !== undefined || resolvedPersona !== undefined); @@ -335,7 +345,7 @@ export function UserProfilePanel({ pubkeyLower.length > 0 && pubkeyLower === currentPubkey.toLowerCase(); const canViewActivity = - viewerIsOwner && + viewerCanObserve && Boolean(effectivePubkey) && canOpenAgentActivity(effectivePubkey); const canOpenAgentLogs = diff --git a/desktop/src/features/profile/ui/UserProfilePopover.tsx b/desktop/src/features/profile/ui/UserProfilePopover.tsx index f2739088a6..093538ddfd 100644 --- a/desktop/src/features/profile/ui/UserProfilePopover.tsx +++ b/desktop/src/features/profile/ui/UserProfilePopover.tsx @@ -22,6 +22,7 @@ import { useManagedAgentsQuery, } from "@/features/agents/hooks"; import { useIsManagedAgent } from "@/features/agent-memory/hooks"; +import { useActiveAgentTurns } from "@/features/agents/activeAgentTurnsStore"; import { useIdentityQuery } from "@/shared/api/hooks"; import { useAgentWorking } from "@/features/agents/agentWorkingSignal"; import { @@ -265,15 +266,18 @@ export function UserProfilePopover({ !isAgentClassificationPending && (!isBotProfile || viewerIsOwner); const showAnyProfileActions = showHumanProfileActions || showMessageAction; + const activeTurns = useAgentWorking(isBotProfile ? pubkey : null).channels; + const observableTurns = useActiveAgentTurns(isBotProfile ? pubkey : null); const canViewActivity = - isBotProfile && viewerIsOwner && canOpenAgentActivity(pubkey); + isBotProfile && + (viewerIsOwner || observableTurns.length > 0) && + canOpenAgentActivity(pubkey); const presenceStatus = presenceQuery.data?.[pubkey.toLowerCase()]; const userStatus = userStatusQuery.data?.[pubkey.toLowerCase()]; const userStatusText = userStatus?.text.trim() ?? ""; const hasUserStatus = Boolean(userStatusText || userStatus?.emoji); const profileDescription = profile?.about?.trim() ?? ""; const profileSubheader = profileDescription || profile?.nip05Handle?.trim(); - const activeTurns = useAgentWorking(isBotProfile ? pubkey : null).channels; const channelsQuery = useChannelsQuery(); const channelIdToName = React.useMemo(() => { const map: Record = {}; diff --git a/desktop/tests/e2e/channels.spec.ts b/desktop/tests/e2e/channels.spec.ts index c7dc7ecb5a..c590e95b94 100644 --- a/desktop/tests/e2e/channels.spec.ts +++ b/desktop/tests/e2e/channels.spec.ts @@ -22,6 +22,7 @@ const MOCK_IDENTITY_PUBKEY = "deadbeef".repeat(8); // view-activity gate (memberIsBot && viewerIsOwner) opens for it. const OWNED_RELAY_AGENT_PUBKEY = "a1b2c3d4e5f60718293a4b5c6d7e8f90112233445566778899aabbccddeeff00"; +const FOREIGN_AGENT_OWNER_PUBKEY = "feedface".repeat(8); const DM_RELAY_AGENT_PUBKEY = "eeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeee"; @@ -2014,6 +2015,123 @@ test("profile renders live activity for a viewer-owned relay agent", async ({ ).not.toBeVisible(); }); +test("requester can view a foreign relay agent's active turn without owner controls", async ({ + page, +}) => { + await installMockBridge(page, { + relayAgents: [ + { + pubkey: OWNED_RELAY_AGENT_PUBKEY, + name: "nadia", + agentType: "goose", + channelNames: ["agents"], + respondTo: "anyone", + }, + ], + searchProfiles: [ + { + pubkey: OWNED_RELAY_AGENT_PUBKEY, + displayName: "nadia", + ownerPubkey: FOREIGN_AGENT_OWNER_PUBKEY, + isAgent: true, + }, + ], + }); + await page.goto("/"); + + await openMembersSidebar(page, "agents"); + await page + .getByTestId(`sidebar-member-open-profile-${OWNED_RELAY_AGENT_PUBKEY}`) + .click(); + await expect( + page.getByTestId(`user-profile-view-activity-${OWNED_RELAY_AGENT_PUBKEY}`), + ).toHaveCount(0); + + await waitForMockLiveSubscription(page, "agents", KIND_TYPING_INDICATOR); + await page.evaluate((pubkey) => { + window.__BUZZ_E2E_EMIT_MOCK_TYPING__?.({ + channelName: "agents", + pubkey, + }); + }, OWNED_RELAY_AGENT_PUBKEY); + await expect( + page.getByTestId(`user-profile-view-activity-${OWNED_RELAY_AGENT_PUBKEY}`), + ).toHaveCount(0); + + await page.waitForFunction( + () => + typeof (window as MockFeedWindow).__BUZZ_E2E_SEED_ACTIVE_TURNS__ === + "function" && + typeof ( + window as MockFeedWindow & { + __BUZZ_E2E_SEED_OBSERVER_EVENTS__?: unknown; + } + ).__BUZZ_E2E_SEED_OBSERVER_EVENTS__ === "function", + ); + await page.evaluate( + ({ agentPubkey, channelId }) => { + (window as MockFeedWindow).__BUZZ_E2E_SEED_ACTIVE_TURNS__?.({ + agentPubkey, + channelId, + turnId: "requester-visible-turn", + }); + ( + window as MockFeedWindow & { + __BUZZ_E2E_SEED_OBSERVER_EVENTS__?: (input: { + agentPubkey: string; + events: unknown[]; + }) => void; + } + ).__BUZZ_E2E_SEED_OBSERVER_EVENTS__?.({ + agentPubkey, + events: [ + { + seq: 1, + timestamp: new Date().toISOString(), + kind: "acp_read", + agentIndex: 0, + channelId, + sessionId: "requester-session", + turnId: "requester-visible-turn", + payload: { + method: "session/update", + params: { + sessionId: "requester-session", + update: { + sessionUpdate: "agent_message_chunk", + messageId: "requester-progress", + content: { + type: "text", + text: "Checking the shared relay configuration now.", + }, + }, + }, + }, + }, + ], + }); + }, + { + agentPubkey: OWNED_RELAY_AGENT_PUBKEY, + channelId: AGENTS_CHANNEL_ID, + }, + ); + + const liveActivity = page.getByTestId( + `user-profile-live-activity-${OWNED_RELAY_AGENT_PUBKEY}`, + ); + await expect(liveActivity).toBeVisible(); + await liveActivity.getByRole("button").first().click(); + + const activityPanel = page.getByTestId("agent-session-thread-panel"); + await expect(activityPanel).toBeVisible(); + await expect(activityPanel).toContainText( + "Checking the shared relay configuration now.", + ); + await page.getByTestId("agent-session-settings-menu-trigger").click(); + await expect(page.getByTestId("agent-session-stop-turn")).toBeDisabled(); +}); + test("profile activity carousel switches channels via progress dots", async ({ page, }) => {