Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
230 changes: 226 additions & 4 deletions crates/buzz-acp/src/lib.rs
Original file line number Diff line number Diff line change
Expand Up @@ -257,6 +257,60 @@ async fn author_allowed(
}
}

/// Return the workflow owner represented by a relay-signed message.
///
/// Attribution is trusted only after verifying the NIP-11 relay signer, event
/// signature, message kind, and exact single-value workflow/actor tags. Any
/// ambiguity fails closed and falls back to the literal event signer.
fn trusted_workflow_actor(
event: &nostr::Event,
trusted_relay_pubkey: Option<&str>,
) -> Option<String> {
let trusted_relay_pubkey = trusted_relay_pubkey?;
if event.pubkey.to_hex() != trusted_relay_pubkey
|| event.kind.as_u16() as u32 != KIND_STREAM_MESSAGE
|| event.verify().is_err()
{
return None;
}

let workflow_tags: Vec<&[String]> = event
.tags
.iter()
.filter(|tag| tag.as_slice().first().map(String::as_str) == Some("buzz:workflow"))
.map(nostr::Tag::as_slice)
.collect();
if workflow_tags.len() != 1
|| workflow_tags[0].len() != 2
|| workflow_tags[0].get(1).map(String::as_str) != Some("true")
{
return None;
}

let actor_tags: Vec<&[String]> = event
.tags
.iter()
.filter(|tag| tag.as_slice().first().map(String::as_str) == Some("actor"))
.map(nostr::Tag::as_slice)
.collect();
if actor_tags.len() != 1 || actor_tags[0].len() != 2 {
return None;
}
let actor = actor_tags[0].get(1).map(String::as_str)?;
if actor.len() != 64 || !actor.chars().all(|c| c.is_ascii_hexdigit()) {
return None;
}

PublicKey::from_hex(actor)
.ok()
.map(|pubkey| pubkey.to_hex())
}

/// Resolve the identity used by the inbound author gate.
fn effective_inbound_author(event: &nostr::Event, trusted_relay_pubkey: Option<&str>) -> String {
trusted_workflow_actor(event, trusted_relay_pubkey).unwrap_or_else(|| event.pubkey.to_hex())
}

/// Resolve whether `channel_id` is a DM, for the inbound author gate.
///
/// Resolution order:
Expand Down Expand Up @@ -1406,6 +1460,20 @@ async fn tokio_main() -> Result<()> {

tracing::info!("connected to relay at {}", config.relay_url);

let rest_client = relay.rest_client();
let trusted_relay_pubkey = match rest_client.relay_self_pubkey().await {
Ok(pubkey) => {
tracing::info!("trusted relay workflow signer: {pubkey}");
Some(pubkey)
}
Err(error) => {
tracing::warn!(
"failed to resolve trusted relay workflow signer; workflow attribution disabled: {error}"
);
None
}
};

relay
.subscribe_membership_notifications()
.await
Expand Down Expand Up @@ -1600,8 +1668,8 @@ async fn tokio_main() -> Result<()> {
.unwrap_or_else(|_| std::path::PathBuf::from("/"))
.to_string_lossy()
.to_string(),
rest_client: relay.rest_client(),
channel_info: pool::ChannelInfoResolver::new(channel_info_map, relay.rest_client()),
rest_client: rest_client.clone(),
channel_info: pool::ChannelInfoResolver::new(channel_info_map, rest_client),
context_message_limit: config.context_message_limit,
max_turns_per_session: config.max_turns_per_session,
permission_mode: config.permission_mode,
Expand Down Expand Up @@ -2215,7 +2283,10 @@ async fn tokio_main() -> Result<()> {
// explicit pubkey list on top, for external people;
// it never revokes same-owner team bots.
{
let author = buzz_event.event.pubkey.to_hex();
let author = effective_inbound_author(
&buzz_event.event,
trusted_relay_pubkey.as_deref(),
);
// DM hardening: resolve channel type (fail-closed
// to DM) so allowlist/anyone modes cannot be
// exercised by non-owner authors inside DMs.
Expand All @@ -2233,7 +2304,8 @@ async fn tokio_main() -> Result<()> {
if !allowed {
tracing::debug!(
channel_id = %buzz_event.channel_id,
author = %buzz_event.event.pubkey.to_hex(),
signer = %buzz_event.event.pubkey.to_hex(),
effective_author = %author,
mode = %config.respond_to,
is_dm,
"inbound author gate — dropping event"
Expand Down Expand Up @@ -4889,6 +4961,156 @@ mod author_gate_tests {
}
}

#[cfg(test)]
mod workflow_author_tests {
use super::*;
use nostr::{EventBuilder, Keys, Kind, Tag};

fn workflow_event(signer: &Keys, actor: &str, kind: u32, extra_tags: Vec<Tag>) -> nostr::Event {
let mut tags = vec![
Tag::parse(["buzz:workflow", "true"]).expect("workflow tag"),
Tag::parse(["actor", actor]).expect("actor tag"),
];
tags.extend(extra_tags);
EventBuilder::new(Kind::Custom(kind as u16), "nightly work")
.tags(tags)
.sign_with_keys(signer)
.expect("signed workflow event")
}

#[test]
fn relay_signed_workflow_uses_verified_actor() {
let relay = Keys::generate();
let actor = Keys::generate();
let event = workflow_event(
&relay,
&actor.public_key().to_hex(),
KIND_STREAM_MESSAGE,
vec![],
);

assert_eq!(
trusted_workflow_actor(&event, Some(&relay.public_key().to_hex())),
Some(actor.public_key().to_hex())
);
assert_eq!(
effective_inbound_author(&event, Some(&relay.public_key().to_hex())),
actor.public_key().to_hex()
);
}

#[test]
fn forged_workflow_attribution_falls_back_to_literal_signer() {
let relay = Keys::generate();
let attacker = Keys::generate();
let actor = Keys::generate();
let event = workflow_event(
&attacker,
&actor.public_key().to_hex(),
KIND_STREAM_MESSAGE,
vec![],
);

assert_eq!(
trusted_workflow_actor(&event, Some(&relay.public_key().to_hex())),
None
);
assert_eq!(
effective_inbound_author(&event, Some(&relay.public_key().to_hex())),
attacker.public_key().to_hex()
);
}

#[test]
fn workflow_attribution_fails_closed_without_relay_identity_or_valid_signature() {
let relay = Keys::generate();
let actor = Keys::generate();
let event = workflow_event(
&relay,
&actor.public_key().to_hex(),
KIND_STREAM_MESSAGE,
vec![],
);
assert_eq!(trusted_workflow_actor(&event, None), None);

let mut tampered = event;
tampered.content = "tampered".into();
assert_eq!(
trusted_workflow_actor(&tampered, Some(&relay.public_key().to_hex())),
None
);
}

#[test]
fn workflow_attribution_rejects_wrong_kind_and_ambiguous_tags() {
let relay = Keys::generate();
let actor = Keys::generate();
let actor_hex = actor.public_key().to_hex();

let wrong_kind = workflow_event(&relay, &actor_hex, 1, vec![]);
assert_eq!(
trusted_workflow_actor(&wrong_kind, Some(&relay.public_key().to_hex())),
None
);

let duplicate_actor = workflow_event(
&relay,
&actor_hex,
KIND_STREAM_MESSAGE,
vec![Tag::parse(["actor", &actor_hex]).expect("duplicate actor tag")],
);
assert_eq!(
trusted_workflow_actor(&duplicate_actor, Some(&relay.public_key().to_hex())),
None
);

let duplicate_workflow = workflow_event(
&relay,
&actor_hex,
KIND_STREAM_MESSAGE,
vec![Tag::parse(["buzz:workflow", "true"]).expect("duplicate workflow tag")],
);
assert_eq!(
trusted_workflow_actor(&duplicate_workflow, Some(&relay.public_key().to_hex())),
None
);

let malformed_actor = workflow_event(&relay, "not-a-pubkey", KIND_STREAM_MESSAGE, vec![]);
assert_eq!(
trusted_workflow_actor(&malformed_actor, Some(&relay.public_key().to_hex())),
None
);

let malformed_duplicate_actor = workflow_event(
&relay,
&actor_hex,
KIND_STREAM_MESSAGE,
vec![Tag::parse(["actor", &actor_hex, "unexpected"]).expect("malformed actor tag")],
);
assert_eq!(
trusted_workflow_actor(
&malformed_duplicate_actor,
Some(&relay.public_key().to_hex())
),
None
);

let malformed_duplicate_workflow = workflow_event(
&relay,
&actor_hex,
KIND_STREAM_MESSAGE,
vec![Tag::parse(["buzz:workflow", "false"]).expect("malformed workflow tag")],
);
assert_eq!(
trusted_workflow_actor(
&malformed_duplicate_workflow,
Some(&relay.public_key().to_hex())
),
None
);
}
}

#[cfg(test)]
mod observer_snapshot_race_tests {
use super::*;
Expand Down
59 changes: 58 additions & 1 deletion crates/buzz-acp/src/relay.rs
Original file line number Diff line number Diff line change
Expand Up @@ -118,7 +118,7 @@ use buzz_core::kind::{
KIND_TYPING_INDICATOR,
};
use futures_util::{SinkExt, StreamExt};
use nostr::{Event, EventBuilder, Keys, Kind, RelayUrl, Tag};
use nostr::{Event, EventBuilder, Keys, Kind, PublicKey, RelayUrl, Tag};
use serde_json::{json, Value};
use tokio::sync::mpsc;
use tokio::time::timeout;
Expand Down Expand Up @@ -259,6 +259,28 @@ fn unix_now_secs() -> u64 {
}

impl RestClient {
/// Fetch the relay's NIP-11 `self` pubkey from its public root document.
///
/// This is intentionally unauthenticated metadata. Callers use it as the
/// expected signer for relay-authored events and must still verify each
/// event's id and signature locally.
pub async fn relay_self_pubkey(&self) -> Result<String, RelayError> {
let url = self.base_url.clone();
let response = self
.request_with_retry("GET", "/", || {
self.http
.get(&url)
.header("Accept", "application/nostr+json")
.send()
})
.await?;
let document: Value = response
.json()
.await
.map_err(|e| RelayError::Http(format!("invalid NIP-11 document: {e}")))?;
normalize_relay_self_pubkey(&document)
}

/// Sign a NIP-98 HTTP Auth event (kind:27235) for the given method/URL/body.
///
/// Returns the `Authorization: Nostr <base64>` header value (without the
Expand Down Expand Up @@ -436,6 +458,22 @@ impl RestClient {
}
}

/// Validate and canonicalize the NIP-11 relay-info `self` pubkey.
fn normalize_relay_self_pubkey(document: &Value) -> Result<String, RelayError> {
let self_hex = document
.get("self")
.and_then(Value::as_str)
.ok_or_else(|| RelayError::Http("NIP-11 document missing 'self' field".into()))?;
if self_hex.len() != 64 || !self_hex.chars().all(|c| c.is_ascii_hexdigit()) {
return Err(RelayError::Http(format!(
"NIP-11 'self' field is not a valid 64-hex pubkey: {self_hex}"
)));
}
PublicKey::from_hex(self_hex)
.map(|pubkey| pubkey.to_hex())
.map_err(|e| RelayError::Http(format!("invalid NIP-11 'self' pubkey: {e}")))
}

/// Events the harness cares about.
#[derive(Debug, Clone)]
pub struct BuzzEvent {
Expand Down Expand Up @@ -4008,6 +4046,25 @@ async fn wait_for_any_ok(
mod tests {
use super::*;

#[test]
fn normalize_relay_self_pubkey_accepts_and_canonicalizes_hex() {
let keys = Keys::generate();
let upper = keys.public_key().to_hex().to_ascii_uppercase();
let document = serde_json::json!({"self": upper});

assert_eq!(
normalize_relay_self_pubkey(&document).expect("valid relay self pubkey"),
keys.public_key().to_hex()
);
}

#[test]
fn normalize_relay_self_pubkey_rejects_missing_or_invalid_values() {
assert!(normalize_relay_self_pubkey(&serde_json::json!({})).is_err());
assert!(normalize_relay_self_pubkey(&serde_json::json!({"self": "not-a-pubkey"})).is_err());
assert!(normalize_relay_self_pubkey(&serde_json::json!({"self": "z".repeat(64)})).is_err());
}

#[test]
fn relay_ws_to_http_plain() {
assert_eq!(
Expand Down
17 changes: 15 additions & 2 deletions crates/buzz-acp/src/setup_mode.rs
Original file line number Diff line number Diff line change
Expand Up @@ -73,7 +73,7 @@ pub(crate) enum AcpAvailabilityStatus {
use crate::{
author_allowed,
config::Config,
event_mentions_agent, filter,
effective_inbound_author, event_mentions_agent, filter,
relay::{HarnessRelay, RelayEventPublisher},
};

Expand Down Expand Up @@ -382,6 +382,18 @@ pub(crate) async fn run_setup_listener(config: Config, payload: SetupPayload) ->

let publisher = relay.event_publisher();
let rest_client = relay.rest_client();
let trusted_relay_pubkey = match rest_client.relay_self_pubkey().await {
Ok(pubkey) => {
tracing::info!("setup-mode: trusted relay workflow signer: {pubkey}");
Some(pubkey)
}
Err(error) => {
tracing::warn!(
"setup-mode: failed to resolve trusted relay workflow signer; workflow attribution disabled: {error}"
);
None
}
};

let channel_info = crate::pool::ChannelInfoResolver::new(channel_info_map, rest_client.clone());

Expand Down Expand Up @@ -428,7 +440,8 @@ pub(crate) async fn run_setup_listener(config: Config, payload: SetupPayload) ->
// Apply the same author gate as normal mode so the nudge only goes
// to authors the real agent would have answered. Same DM hardening:
// in DMs only owner/siblings get a nudge (fail-closed on unknown type).
let author_hex = buzz_event.event.pubkey.to_hex();
let author_hex =
effective_inbound_author(&buzz_event.event, trusted_relay_pubkey.as_deref());
let is_dm = crate::is_dm_channel(buzz_event.channel_id, &channel_info).await;
let allowed = author_allowed(
&config.respond_to,
Expand Down
Loading