From 1ddbc0ef7781d5499674934a8df7c38e1185bbb1 Mon Sep 17 00:00:00 2001 From: coder 0 Date: Wed, 16 Sep 2026 18:58:50 -0400 Subject: [PATCH 1/4] feat(relay): add configurable relay banners Signed-off-by: coder 0 --- crates/buzz-core/src/kind.rs | 14 + crates/buzz-db/src/lib.rs | 8 +- crates/buzz-db/src/runtime/migration.rs | 8 + crates/buzz-db/src/store/mod.rs | 2 + crates/buzz-db/src/store/relay_banners.rs | 868 ++++++++++++++++++++++ crates/buzz-relay/src/api/admin/mod.rs | 225 +++++- crates/buzz-relay/src/api/banners.rs | 373 ++++++++++ crates/buzz-relay/src/api/bridge.rs | 33 + crates/buzz-relay/src/api/mod.rs | 1 + crates/buzz-relay/src/handlers/req.rs | 38 +- crates/buzz-relay/src/router.rs | 12 + migrations/0046_relay_banners.sql | 62 ++ schema/schema.sql | 66 +- 13 files changed, 1704 insertions(+), 6 deletions(-) create mode 100644 crates/buzz-db/src/store/relay_banners.rs create mode 100644 crates/buzz-relay/src/api/banners.rs create mode 100644 migrations/0046_relay_banners.sql diff --git a/crates/buzz-core/src/kind.rs b/crates/buzz-core/src/kind.rs index 4e1ab1c7f5e..7f89c3e20ab 100644 --- a/crates/buzz-core/src/kind.rs +++ b/crates/buzz-core/src/kind.rs @@ -437,6 +437,9 @@ pub const KIND_THREAD_SUMMARY: u32 = 39005; /// content = `{has_more, next_cursor}`. The only authority on exhaustion — /// clients must not infer `has_more` from row counts. pub const KIND_WINDOW_BOUNDS: u32 = 39006; +/// Deployment banner overlay: relay-signed active operator banner for this viewer. +/// Content is `{id, severity, text, maxDisplays}` and `d` tag is the banner public id. +pub const KIND_RELAY_BANNER: u32 = 13536; /// Workflow definition (parameterized replaceable, d=workflow_uuid). pub const KIND_WORKFLOW_DEF: u32 = 30620; @@ -694,6 +697,7 @@ pub const ALL_KINDS: &[u32] = &[ KIND_NIP29_GROUP_ROLES, KIND_THREAD_SUMMARY, KIND_WINDOW_BOUNDS, + KIND_RELAY_BANNER, KIND_PRESENCE_UPDATE, KIND_TYPING_INDICATOR, KIND_HUDDLE_REACTION, @@ -839,6 +843,7 @@ pub const fn is_relay_only_kind(kind: u32) -> bool { | KIND_DM_VISIBILITY | KIND_THREAD_SUMMARY | KIND_WINDOW_BOUNDS + | KIND_RELAY_BANNER ) } @@ -867,6 +872,7 @@ const _: () = assert!(is_parameterized_replaceable(KIND_DM_VISIBILITY)); // 3062 const _: () = assert!(is_parameterized_replaceable(KIND_PROJECT)); // 30621 ∈ 30000–39999 const _: () = assert!(is_parameterized_replaceable(KIND_THREAD_SUMMARY)); // 39005 ∈ 30000–39999 const _: () = assert!(is_parameterized_replaceable(KIND_WINDOW_BOUNDS)); // 39006 ∈ 30000–39999 +const _: () = assert!(is_replaceable(KIND_RELAY_BANNER)); // 13536 ∈ 10000–19999 // Compile-time: NIP-34 parameterized replaceable kinds are in the correct range. const _: () = assert!( @@ -936,6 +942,14 @@ mod tests { } } + #[test] + fn relay_banner_is_replaceable_relay_only_not_parameterized() { + assert_eq!(KIND_RELAY_BANNER, 13536); + assert!(is_replaceable(KIND_RELAY_BANNER)); + assert!(!is_parameterized_replaceable(KIND_RELAY_BANNER)); + assert!(is_relay_only_kind(KIND_RELAY_BANNER)); + } + // ── event_is_shared / is_unshared_gated_event ──────────────────────── fn make_event_of_kind(kind: u32, tags: &[&[&str]]) -> nostr::Event { diff --git a/crates/buzz-db/src/lib.rs b/crates/buzz-db/src/lib.rs index d444d9ce52b..4a11131e2e3 100644 --- a/crates/buzz-db/src/lib.rs +++ b/crates/buzz-db/src/lib.rs @@ -64,8 +64,8 @@ pub(crate) use runtime::{ pub use store::{ admin_moderation, allowlist, api_token, archived_identities, channel, channel_members, community, deletion, dm, event, feed, git_repo, moderation, partition, product_feedback, push, - reaction, relay_admin_actions, relay_invite, relay_members, relay_operators, reminder, - replaceable, storage_accounting, thread, usage, user, workflow, + reaction, relay_admin_actions, relay_banners, relay_invite, relay_members, relay_operators, + reminder, replaceable, storage_accounting, thread, usage, user, workflow, }; pub use allowlist::AllowlistEntry; @@ -78,6 +78,10 @@ pub use community::{ pub use error::{DbError, Result}; pub use event::{EventQuery, DEFAULT_MAX_PAGE_LIMIT}; pub use reaction::ReactionEventInsertOutcome; +pub use relay_banners::{ + RelayBannerDismissOutcome, RelayBannerRecord, RelayBannerScope, RelayBannerSeverity, + RelayBannerTargetScope, RelayBannerUpsert, RelayBannerViewOutcome, MAX_BANNER_MESSAGE_CHARS, +}; pub use reminder::DueReminder; pub use usage::UsageMetricsLeader; diff --git a/crates/buzz-db/src/runtime/migration.rs b/crates/buzz-db/src/runtime/migration.rs index aebf8b9092f..e984b2a836f 100644 --- a/crates/buzz-db/src/runtime/migration.rs +++ b/crates/buzz-db/src/runtime/migration.rs @@ -490,6 +490,9 @@ mod postgres_tests { "relay_admin_outbox", "relay_operator_audit", "storage_accounting_snapshots", + "relay_banners", + "relay_banner_communities", + "relay_banner_user_state", ] { if normalized[insert_pos..].contains(&format!("'{value}'")) { globals.insert(value.to_owned()); @@ -1828,6 +1831,8 @@ mod postgres_tests { let mut expected_fences = migration.fence_attachments.clone(); expected_fences.remove("product_feedback"); expected_fences.remove("rate_limit_violations"); + expected_fences.insert("relay_banner_communities".to_owned()); + expected_fences.insert("relay_banner_user_state".to_owned()); assert_eq!( expected_fences, schema.fence_attachments, "write-fence attachment targets differ after recovery policy" @@ -2267,6 +2272,9 @@ mod postgres_tests { "relay_admin_actions", "relay_admin_outbox", "relay_operator_audit", + "relay_banners", + "relay_banner_communities", + "relay_banner_user_state", ] { assert_eq!( columns(&desired, table).await, diff --git a/crates/buzz-db/src/store/mod.rs b/crates/buzz-db/src/store/mod.rs index 79739373d0f..3df8b257414 100644 --- a/crates/buzz-db/src/store/mod.rs +++ b/crates/buzz-db/src/store/mod.rs @@ -36,6 +36,8 @@ pub mod push; pub mod reaction; /// HTTP report-resolution enforcement state machine persistence. pub mod relay_admin_actions; +/// Deployment-global relay banner persistence. +pub mod relay_banners; /// Use-limited relay invite persistence (v2 opaque tokens). pub mod relay_invite; /// Relay-level membership persistence (NIP-43). diff --git a/crates/buzz-db/src/store/relay_banners.rs b/crates/buzz-db/src/store/relay_banners.rs new file mode 100644 index 00000000000..df623a17cf9 --- /dev/null +++ b/crates/buzz-db/src/store/relay_banners.rs @@ -0,0 +1,868 @@ +//! Deployment-global relay banner persistence and per-user state. +//! +//! A deployment has at most one active banner. Delivery is community-scoped at +//! read time: an active banner either targets every community, including future +//! communities, or an explicit non-empty set of community ids. + +use buzz_core::CommunityId; +use buzz_datastore_tracing::datastore_span; +use chrono::{DateTime, Utc}; +use serde::{Deserialize, Serialize}; +use sqlx::{Postgres, Row as _, Transaction}; +use uuid::Uuid; + +use crate::{Db, Result}; + +/// Maximum relay banner message length in Unicode scalar count. +pub const MAX_BANNER_MESSAGE_CHARS: usize = 2_000; + +/// Relay banner severity values accepted by the backend and emitted to clients. +#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)] +#[serde(rename_all = "lowercase")] +pub enum RelayBannerSeverity { + /// Informational banner. + Info, + /// Warning banner. + Warning, + /// Urgent banner. + Urgent, +} + +impl RelayBannerSeverity { + /// Returns the wire/database representation. + pub const fn as_str(self) -> &'static str { + match self { + Self::Info => "info", + Self::Warning => "warning", + Self::Urgent => "urgent", + } + } + + fn parse(value: String) -> Result { + match value.as_str() { + "info" => Ok(Self::Info), + "warning" => Ok(Self::Warning), + "urgent" => Ok(Self::Urgent), + other => Err(crate::DbError::InvalidData(format!( + "invalid relay banner severity: {other}" + ))), + } + } +} + +/// Public targeting scope value used by the admin/client banner contracts. +#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)] +#[serde(rename_all = "lowercase")] +pub enum RelayBannerTargetScope { + /// Banner targets every community. + All, + /// Banner targets selected communities only. + Communities, +} + +impl RelayBannerTargetScope { + /// Returns the wire representation. + pub const fn as_str(self) -> &'static str { + match self { + Self::All => "all", + Self::Communities => "communities", + } + } +} + +/// Targeting scope for an operator-configured relay banner. +#[derive(Debug, Clone, PartialEq, Eq)] +pub enum RelayBannerScope { + /// Banner applies to every community, including communities created later. + AllCommunities, + /// Banner applies only to this non-empty community set. + Communities(Vec), +} + +/// Input for replacing the deployment's active banner. +#[derive(Debug, Clone)] +pub struct RelayBannerUpsert { + /// Severity/type of the banner. + pub severity: RelayBannerSeverity, + /// Plain-text message, 1..=2000 characters. + pub message: String, + /// Number of successful views each user may receive before suppression. + pub max_displays: i32, + /// Targeting scope. + pub scope: RelayBannerScope, + /// Authenticated operator pubkey performing the write. + pub actor_pubkey: Vec, +} + +/// Stored relay banner plus targeting metadata. +#[derive(Debug, Clone, PartialEq, Eq)] +pub struct RelayBannerRecord { + /// Monotonic internal banner id. + pub id: i64, + /// Stable public banner id exposed on the wire. + pub public_id: Uuid, + /// Severity/type of the banner. + pub severity: RelayBannerSeverity, + /// Plain-text message. + pub message: String, + /// Number of successful views each user may receive before suppression. + pub max_displays: i32, + /// Whether the banner targets every community. + pub target_all_communities: bool, + /// Explicit community targets when not targeting all communities. + pub community_ids: Vec, + /// Row creation timestamp. + pub created_at: DateTime, + /// Row update timestamp. + pub updated_at: DateTime, + /// Operator pubkey that created the row. + pub created_by: Vec, + /// Disable timestamp, if disabled. + pub disabled_at: Option>, +} + +impl RelayBannerRecord { + /// Returns the public target-scope enum for this banner. + pub const fn target_scope(&self) -> RelayBannerTargetScope { + if self.target_all_communities { + RelayBannerTargetScope::All + } else { + RelayBannerTargetScope::Communities + } + } +} + +/// Result of recording a banner view acknowledgement. +#[derive(Debug, Clone, PartialEq, Eq)] +pub enum RelayBannerViewOutcome { + /// The display was accepted and consumed one remaining view. + Accepted { + /// Display count after accepting this view. + display_count: i32, + }, + /// The banner was already permanently dismissed by this user. + Dismissed, + /// The user had already exhausted `max_displays`. + Exhausted { + /// Display count already consumed by this user. + display_count: i32, + }, + /// The banner is not active or does not target the request community. + NotEligible, +} + +/// Result of recording a banner dismiss acknowledgement. +#[derive(Debug, Clone, PartialEq, Eq)] +pub enum RelayBannerDismissOutcome { + /// Dismiss state was present after the call. True means this call wrote it. + Dismissed { + /// Whether this call changed durable state from not-dismissed to dismissed. + changed: bool, + }, + /// The banner is not active or does not target the request community. + NotEligible, +} + +impl RelayBannerUpsert { + fn validate(&self) -> Result<()> { + let trimmed_empty = self.message.trim().is_empty(); + let char_count = self.message.chars().count(); + if trimmed_empty || char_count > MAX_BANNER_MESSAGE_CHARS { + return Err(crate::DbError::InvalidData(format!( + "relay banner message must be 1..={MAX_BANNER_MESSAGE_CHARS} characters" + ))); + } + if self.max_displays < 1 { + return Err(crate::DbError::InvalidData( + "relay banner max_displays must be >= 1".to_owned(), + )); + } + if self.actor_pubkey.len() != 32 { + return Err(crate::DbError::InvalidData( + "relay banner actor pubkey must be 32 bytes".to_owned(), + )); + } + if let RelayBannerScope::Communities(ids) = &self.scope { + if ids.is_empty() { + return Err(crate::DbError::InvalidData( + "relay banner community scope must be non-empty".to_owned(), + )); + } + } + Ok(()) + } +} + +async fn acquire_banner_lock(tx: &mut Transaction<'_, Postgres>) -> Result<()> { + sqlx::query("SELECT pg_advisory_xact_lock(hashtextextended('relay_banner_active', 0))") + .execute(&mut **tx) + .await?; + Ok(()) +} + +async fn banner_internal_id_for_public_id( + tx: &mut Transaction<'_, Postgres>, + public_id: Uuid, +) -> Result> { + sqlx::query_scalar::<_, i64>( + r#" + SELECT id + FROM relay_banners + WHERE public_id = $1 + AND disabled_at IS NULL + "#, + ) + .bind(public_id) + .fetch_optional(&mut **tx) + .await + .map_err(Into::into) +} + +async fn banner_targets_community( + tx: &mut Transaction<'_, Postgres>, + banner_id: i64, + community_id: CommunityId, +) -> Result { + let eligible = sqlx::query_scalar::<_, bool>( + r#" + SELECT EXISTS ( + SELECT 1 + FROM relay_banners b + WHERE b.id = $1 + AND b.disabled_at IS NULL + AND ( + b.target_all_communities + OR EXISTS ( + SELECT 1 + FROM relay_banner_communities bc + WHERE bc.banner_id = b.id + AND bc.community_id = $2 + ) + ) + ) + "#, + ) + .bind(banner_id) + .bind(community_id.as_uuid()) + .fetch_one(&mut **tx) + .await?; + Ok(eligible) +} + +impl Db { + /// Lists active communities for the admin banner community picker. + #[datastore_span(name = "admin_list_banner_communities", system = "postgresql")] + pub async fn admin_list_banner_communities(&self) -> Result> { + let mut connection = crate::observability::acquire_writer( + &self.pool, + crate::observability::WriterOperation::Authorization, + ) + .await?; + let rows = sqlx::query( + r#" + SELECT id, host + FROM communities + WHERE archived_at IS NULL + AND deleted_at IS NULL + AND deletion_state = 'active' + ORDER BY lower(host), id + "#, + ) + .fetch_all(&mut *connection) + .await?; + + rows.into_iter() + .map(|row| { + Ok(crate::CommunityRecord { + id: CommunityId::from_uuid(row.try_get("id")?), + host: row.try_get("host")?, + }) + }) + .collect() + } + + /// Lists relay banners, newest first, retaining disabled history. + #[datastore_span(name = "admin_list_relay_banners", system = "postgresql")] + pub async fn admin_list_relay_banners(&self, limit: i64) -> Result> { + let limit = limit.clamp(1, 200); + let mut connection = crate::observability::acquire_writer( + &self.pool, + crate::observability::WriterOperation::Authorization, + ) + .await?; + let rows = sqlx::query( + r#" + SELECT + b.id, + b.public_id, + b.severity, + b.message, + b.max_displays, + b.target_all_communities, + b.created_by, + b.created_at, + b.updated_at, + b.disabled_at, + COALESCE(array_agg(bc.community_id ORDER BY bc.community_id) + FILTER (WHERE bc.community_id IS NOT NULL), ARRAY[]::uuid[]) AS community_ids + FROM relay_banners b + LEFT JOIN relay_banner_communities bc ON bc.banner_id = b.id + GROUP BY b.id + ORDER BY b.updated_at DESC, b.id DESC + LIMIT $1 + "#, + ) + .bind(limit) + .fetch_all(&mut *connection) + .await?; + rows.into_iter().map(row_to_banner).collect() + } + + /// Returns the currently active deployment banner, if any. + #[datastore_span(name = "admin_get_active_relay_banner", system = "postgresql")] + pub async fn admin_get_active_relay_banner(&self) -> Result> { + let mut items = self.admin_list_active_relay_banners(1).await?; + Ok(items.pop()) + } + + async fn admin_list_active_relay_banners(&self, limit: i64) -> Result> { + let mut connection = crate::observability::acquire_writer( + &self.pool, + crate::observability::WriterOperation::Authorization, + ) + .await?; + let rows = sqlx::query( + r#" + SELECT + b.id, + b.public_id, + b.severity, + b.message, + b.max_displays, + b.target_all_communities, + b.created_by, + b.created_at, + b.updated_at, + b.disabled_at, + COALESCE(array_agg(bc.community_id ORDER BY bc.community_id) + FILTER (WHERE bc.community_id IS NOT NULL), ARRAY[]::uuid[]) AS community_ids + FROM relay_banners b + LEFT JOIN relay_banner_communities bc ON bc.banner_id = b.id + WHERE b.disabled_at IS NULL + GROUP BY b.id + ORDER BY b.id DESC + LIMIT $1 + "#, + ) + .bind(limit) + .fetch_all(&mut *connection) + .await?; + rows.into_iter().map(row_to_banner).collect() + } + + /// Replaces the single active deployment banner and returns the new row. + #[datastore_span(name = "admin_upsert_relay_banner", system = "postgresql")] + pub async fn admin_upsert_relay_banner( + &self, + input: RelayBannerUpsert, + ) -> Result { + input.validate()?; + let connection = crate::observability::acquire_writer( + &self.pool, + crate::observability::WriterOperation::Authorization, + ) + .await?; + let mut tx = sqlx::Transaction::begin(connection, None).await?; + acquire_banner_lock(&mut tx).await?; + + sqlx::query( + r#" + UPDATE relay_banners + SET disabled_at = now(), disabled_by = $1, updated_at = now(), updated_by = $1 + WHERE disabled_at IS NULL + "#, + ) + .bind(&input.actor_pubkey) + .execute(&mut *tx) + .await?; + + let target_all = matches!(input.scope, RelayBannerScope::AllCommunities); + let row = sqlx::query( + r#" + INSERT INTO relay_banners + (severity, message, max_displays, target_all_communities, created_by, updated_by) + VALUES ($1, $2, $3, $4, $5, $5) + RETURNING id, public_id, severity, message, max_displays, target_all_communities, + created_by, created_at, updated_at, disabled_at + "#, + ) + .bind(input.severity.as_str()) + .bind(&input.message) + .bind(input.max_displays) + .bind(target_all) + .bind(&input.actor_pubkey) + .fetch_one(&mut *tx) + .await?; + let banner_id: i64 = row.try_get("id")?; + + let community_ids = match input.scope { + RelayBannerScope::AllCommunities => Vec::new(), + RelayBannerScope::Communities(ids) => { + for community_id in &ids { + sqlx::query( + "INSERT INTO relay_banner_communities (banner_id, community_id) VALUES ($1, $2)", + ) + .bind(banner_id) + .bind(community_id.as_uuid()) + .execute(&mut *tx) + .await?; + } + ids + } + }; + + tx.commit().await?; + Ok(RelayBannerRecord { + id: banner_id, + public_id: row.try_get("public_id")?, + severity: RelayBannerSeverity::parse(row.try_get("severity")?)?, + message: row.try_get("message")?, + max_displays: row.try_get("max_displays")?, + target_all_communities: row.try_get("target_all_communities")?, + community_ids, + created_by: row.try_get("created_by")?, + created_at: row.try_get("created_at")?, + updated_at: row.try_get("updated_at")?, + disabled_at: row.try_get("disabled_at")?, + }) + } + + /// Disables the currently active deployment banner, preserving history and state. + #[datastore_span(name = "admin_disable_active_relay_banner", system = "postgresql")] + pub async fn admin_disable_active_relay_banner(&self, actor_pubkey: &[u8]) -> Result { + let connection = crate::observability::acquire_writer( + &self.pool, + crate::observability::WriterOperation::Authorization, + ) + .await?; + let mut tx = sqlx::Transaction::begin(connection, None).await?; + acquire_banner_lock(&mut tx).await?; + let result = sqlx::query( + r#" + UPDATE relay_banners + SET disabled_at = now(), disabled_by = $1, updated_at = now(), updated_by = $1 + WHERE disabled_at IS NULL + "#, + ) + .bind(actor_pubkey) + .execute(&mut *tx) + .await?; + tx.commit().await?; + Ok(result.rows_affected() > 0) + } + + /// Resolves the active banner eligible for this user in this community. + #[datastore_span(name = "active_relay_banner_for_user", system = "postgresql")] + pub async fn active_relay_banner_for_user( + &self, + community_id: CommunityId, + pubkey: &[u8], + ) -> Result> { + let mut connection = crate::observability::acquire_writer( + &self.pool, + crate::observability::WriterOperation::Authorization, + ) + .await?; + let row = sqlx::query( + r#" + SELECT + b.id, + b.public_id, + b.severity, + b.message, + b.max_displays, + b.target_all_communities, + b.created_by, + b.created_at, + b.updated_at, + b.disabled_at, + COALESCE(array_agg(bc_all.community_id ORDER BY bc_all.community_id) + FILTER (WHERE bc_all.community_id IS NOT NULL), ARRAY[]::uuid[]) AS community_ids + FROM relay_banners b + LEFT JOIN relay_banner_communities bc_all ON bc_all.banner_id = b.id + LEFT JOIN relay_banner_user_state us + ON us.banner_id = b.id + AND us.community_id = $1 + AND us.pubkey = $2 + WHERE b.disabled_at IS NULL + AND (b.target_all_communities OR EXISTS ( + SELECT 1 + FROM relay_banner_communities bc + WHERE bc.banner_id = b.id + AND bc.community_id = $1 + )) + AND us.dismissed_at IS NULL + AND COALESCE(us.display_count, 0) < b.max_displays + GROUP BY b.id, us.display_count, us.dismissed_at + ORDER BY b.id DESC + LIMIT 1 + "#, + ) + .bind(community_id.as_uuid()) + .bind(pubkey) + .fetch_optional(&mut *connection) + .await?; + row.map(row_to_banner).transpose() + } + + /// Records a successful client render/view acknowledgement. + #[datastore_span(name = "ack_relay_banner_view", system = "postgresql")] + pub async fn ack_relay_banner_view( + &self, + community_id: CommunityId, + public_id: Uuid, + pubkey: &[u8], + ) -> Result { + let connection = crate::observability::acquire_writer( + &self.pool, + crate::observability::WriterOperation::Authorization, + ) + .await?; + let mut tx = sqlx::Transaction::begin(connection, None).await?; + let Some(banner_id) = banner_internal_id_for_public_id(&mut tx, public_id).await? else { + tx.rollback().await?; + return Ok(RelayBannerViewOutcome::NotEligible); + }; + if !banner_targets_community(&mut tx, banner_id, community_id).await? { + tx.rollback().await?; + return Ok(RelayBannerViewOutcome::NotEligible); + } + + let row = sqlx::query( + r#" + INSERT INTO relay_banner_user_state + (banner_id, community_id, pubkey, display_count, first_viewed_at, last_viewed_at) + SELECT b.id, $2, $3, 1, now(), now() + FROM relay_banners b + WHERE b.id = $1 + AND b.disabled_at IS NULL + AND b.max_displays >= 1 + ON CONFLICT (banner_id, community_id, pubkey) DO UPDATE SET + display_count = relay_banner_user_state.display_count + 1, + first_viewed_at = COALESCE(relay_banner_user_state.first_viewed_at, now()), + last_viewed_at = now(), + updated_at = now() + WHERE relay_banner_user_state.dismissed_at IS NULL + AND relay_banner_user_state.display_count < ( + SELECT max_displays FROM relay_banners WHERE id = $1 AND disabled_at IS NULL + ) + RETURNING display_count + "#, + ) + .bind(banner_id) + .bind(community_id.as_uuid()) + .bind(pubkey) + .fetch_optional(&mut *tx) + .await?; + + let outcome = if let Some(row) = row { + RelayBannerViewOutcome::Accepted { + display_count: row.try_get("display_count")?, + } + } else { + let state = sqlx::query( + r#" + SELECT us.display_count, us.dismissed_at IS NOT NULL AS dismissed + FROM relay_banner_user_state us + WHERE us.banner_id = $1 AND us.community_id = $2 AND us.pubkey = $3 + "#, + ) + .bind(banner_id) + .bind(community_id.as_uuid()) + .bind(pubkey) + .fetch_optional(&mut *tx) + .await?; + match state { + Some(row) if row.try_get::("dismissed")? => { + RelayBannerViewOutcome::Dismissed + } + Some(row) => RelayBannerViewOutcome::Exhausted { + display_count: row.try_get("display_count")?, + }, + None => RelayBannerViewOutcome::NotEligible, + } + }; + tx.commit().await?; + Ok(outcome) + } + + /// Permanently dismisses a banner for a user in one community. + #[datastore_span(name = "ack_relay_banner_dismiss", system = "postgresql")] + pub async fn ack_relay_banner_dismiss( + &self, + community_id: CommunityId, + public_id: Uuid, + pubkey: &[u8], + ) -> Result { + let connection = crate::observability::acquire_writer( + &self.pool, + crate::observability::WriterOperation::Authorization, + ) + .await?; + let mut tx = sqlx::Transaction::begin(connection, None).await?; + let Some(banner_id) = banner_internal_id_for_public_id(&mut tx, public_id).await? else { + tx.rollback().await?; + return Ok(RelayBannerDismissOutcome::NotEligible); + }; + if !banner_targets_community(&mut tx, banner_id, community_id).await? { + tx.rollback().await?; + return Ok(RelayBannerDismissOutcome::NotEligible); + } + + let changed = sqlx::query_scalar::<_, bool>( + r#" + WITH inserted AS ( + INSERT INTO relay_banner_user_state + (banner_id, community_id, pubkey, dismissed_at) + VALUES ($1, $2, $3, now()) + ON CONFLICT (banner_id, community_id, pubkey) DO NOTHING + RETURNING TRUE AS changed + ), updated AS ( + UPDATE relay_banner_user_state + SET dismissed_at = now(), updated_at = now() + WHERE banner_id = $1 + AND community_id = $2 + AND pubkey = $3 + AND dismissed_at IS NULL + AND NOT EXISTS (SELECT 1 FROM inserted) + RETURNING TRUE AS changed + ) + SELECT COALESCE( + (SELECT changed FROM inserted), + (SELECT changed FROM updated), + FALSE + ) + "#, + ) + .bind(banner_id) + .bind(community_id.as_uuid()) + .bind(pubkey) + .fetch_one(&mut *tx) + .await?; + tx.commit().await?; + Ok(RelayBannerDismissOutcome::Dismissed { changed }) + } +} + +fn row_to_banner(row: sqlx::postgres::PgRow) -> Result { + let community_uuids: Vec = row.try_get("community_ids")?; + Ok(RelayBannerRecord { + id: row.try_get("id")?, + public_id: row.try_get("public_id")?, + severity: RelayBannerSeverity::parse(row.try_get("severity")?)?, + message: row.try_get("message")?, + max_displays: row.try_get("max_displays")?, + target_all_communities: row.try_get("target_all_communities")?, + community_ids: community_uuids + .into_iter() + .map(CommunityId::from_uuid) + .collect(), + created_by: row.try_get("created_by")?, + created_at: row.try_get("created_at")?, + updated_at: row.try_get("updated_at")?, + disabled_at: row.try_get("disabled_at")?, + }) +} + +#[cfg(test)] +mod tests { + use super::*; + use sqlx::PgPool; + + async fn setup_db() -> Db { + let pool = PgPool::connect(&crate::test_support::database_url()) + .await + .expect("connect to test DB"); + Db::from_pool(pool) + } + + async fn make_community(pool: &PgPool) -> CommunityId { + let id = Uuid::new_v4(); + sqlx::query("INSERT INTO communities (id, host) VALUES ($1, $2)") + .bind(id) + .bind(format!("banner-{}.example", id.simple())) + .execute(pool) + .await + .expect("insert community"); + CommunityId::from_uuid(id) + } + + fn actor(byte: u8) -> Vec { + vec![byte; 32] + } + + #[tokio::test] + #[ignore = "requires Postgres"] + async fn upsert_replaces_single_active_banner_and_retains_history() { + let db = setup_db().await; + let first = db + .admin_upsert_relay_banner(RelayBannerUpsert { + severity: RelayBannerSeverity::Info, + message: "first".to_owned(), + max_displays: 1, + scope: RelayBannerScope::AllCommunities, + actor_pubkey: actor(1), + }) + .await + .expect("insert first banner"); + let second = db + .admin_upsert_relay_banner(RelayBannerUpsert { + severity: RelayBannerSeverity::Warning, + message: "second".to_owned(), + max_displays: 2, + scope: RelayBannerScope::AllCommunities, + actor_pubkey: actor(2), + }) + .await + .expect("replace banner"); + + let active = db + .admin_get_active_relay_banner() + .await + .expect("get active") + .expect("active banner"); + assert_eq!(active.id, second.id); + let all = db.admin_list_relay_banners(10).await.expect("list banners"); + assert_eq!(all.len(), 2); + assert!(all + .iter() + .any(|b| b.id == first.id && b.disabled_at.is_some())); + } + + #[tokio::test] + #[ignore = "requires Postgres"] + async fn community_scope_filters_delivery() { + let db = setup_db().await; + let included = make_community(&db.pool).await; + let excluded = make_community(&db.pool).await; + let user = actor(3); + let banner = db + .admin_upsert_relay_banner(RelayBannerUpsert { + severity: RelayBannerSeverity::Urgent, + message: "scoped".to_owned(), + max_displays: 1, + scope: RelayBannerScope::Communities(vec![included]), + actor_pubkey: actor(4), + }) + .await + .expect("insert scoped banner"); + + assert_eq!(banner.community_ids, vec![included]); + assert!(db + .active_relay_banner_for_user(included, &user) + .await + .expect("included lookup") + .is_some()); + assert!(db + .active_relay_banner_for_user(excluded, &user) + .await + .expect("excluded lookup") + .is_none()); + } + + #[tokio::test] + #[ignore = "requires Postgres"] + async fn view_ack_consumes_display_and_then_exhausts() { + let db = setup_db().await; + let community = make_community(&db.pool).await; + let user = actor(5); + let banner = db + .admin_upsert_relay_banner(RelayBannerUpsert { + severity: RelayBannerSeverity::Info, + message: "limited".to_owned(), + max_displays: 1, + scope: RelayBannerScope::AllCommunities, + actor_pubkey: actor(6), + }) + .await + .expect("insert banner"); + + assert_eq!( + db.active_relay_banner_for_user(community, &user) + .await + .expect("lookup before view") + .as_ref() + .map(|active| active.id), + Some(banner.id), + "delivery lookup must not consume max_displays before client view ack" + ); + assert_eq!( + db.active_relay_banner_for_user(community, &user) + .await + .expect("second lookup before view") + .as_ref() + .map(|active| active.id), + Some(banner.id), + "repeated delivery lookups without view ack must not exhaust max_displays" + ); + assert_eq!( + db.ack_relay_banner_view(community, banner.public_id, &user) + .await + .expect("first view"), + RelayBannerViewOutcome::Accepted { display_count: 1 } + ); + assert_eq!( + db.ack_relay_banner_view(community, banner.public_id, &user) + .await + .expect("second view"), + RelayBannerViewOutcome::Exhausted { display_count: 1 } + ); + assert!(db + .active_relay_banner_for_user(community, &user) + .await + .expect("post-exhaust lookup") + .is_none()); + } + + #[tokio::test] + #[ignore = "requires Postgres"] + async fn dismiss_permanently_suppresses_banner() { + let db = setup_db().await; + let community = make_community(&db.pool).await; + let user = actor(7); + let banner = db + .admin_upsert_relay_banner(RelayBannerUpsert { + severity: RelayBannerSeverity::Info, + message: "dismiss me".to_owned(), + max_displays: 3, + scope: RelayBannerScope::AllCommunities, + actor_pubkey: actor(8), + }) + .await + .expect("insert banner"); + + assert_eq!( + db.ack_relay_banner_dismiss(community, banner.public_id, &user) + .await + .expect("dismiss"), + RelayBannerDismissOutcome::Dismissed { changed: true } + ); + assert_eq!( + db.ack_relay_banner_dismiss(community, banner.public_id, &user) + .await + .expect("second dismiss"), + RelayBannerDismissOutcome::Dismissed { changed: false } + ); + assert_eq!( + db.ack_relay_banner_view(community, banner.public_id, &user) + .await + .expect("view after dismiss"), + RelayBannerViewOutcome::Dismissed + ); + assert!(db + .active_relay_banner_for_user(community, &user) + .await + .expect("lookup after dismiss") + .is_none()); + } +} diff --git a/crates/buzz-relay/src/api/admin/mod.rs b/crates/buzz-relay/src/api/admin/mod.rs index 19f2153b95b..25ddf26155f 100644 --- a/crates/buzz-relay/src/api/admin/mod.rs +++ b/crates/buzz-relay/src/api/admin/mod.rs @@ -60,9 +60,13 @@ pub fn router(state: Arc) -> Router { .route("/operators", get(list_operators)) .route("/operators/{pubkey}", put(upsert_operator)) .route("/operators/{pubkey}", delete(delete_operator)) + .route("/communities", get(list_banner_communities)) + .route("/banners", get(banners)) + .route("/banners/current", put(upsert_banner)) + .route("/banners/current", delete(disable_banner)) .layer(middleware::from_fn(security_headers)) - // Mutation routes carry a JSON body (max ~4 KB); read-only routes have no body. - .layer(RequestBodyLimitLayer::new(4096)) + // Mutation routes carry a JSON body; banner text may be up to 2,000 chars. + .layer(RequestBodyLimitLayer::new(8192)) .with_state(state) } @@ -406,6 +410,223 @@ async fn feedback_attachment( Ok(response) } +#[derive(Serialize)] +#[serde(rename_all = "camelCase")] +struct BannerCommunityResponse { + id: Uuid, + host: String, +} + +#[derive(Serialize)] +#[serde(rename_all = "camelCase")] +struct AdminBannerResponse { + id: String, + severity: &'static str, + text: String, + max_displays: i32, + target_scope: &'static str, + community_ids: Vec, + created_by: String, + created_at: DateTime, + disabled_at: Option>, +} + +#[derive(Serialize)] +#[serde(rename_all = "camelCase")] +struct BannersResponse { + active: Option, + history: Vec, +} + +#[derive(Deserialize)] +#[serde(rename_all = "camelCase", deny_unknown_fields)] +struct BannerQuery { + limit: Option, +} + +#[derive(Deserialize)] +#[serde(rename_all = "camelCase", deny_unknown_fields)] +struct UpsertBannerBody { + severity: buzz_db::RelayBannerSeverity, + text: String, + max_displays: i32, + target_scope: buzz_db::RelayBannerTargetScope, + community_ids: Option>, +} + +impl From for AdminBannerResponse { + fn from(value: buzz_db::RelayBannerRecord) -> Self { + let target_scope = value.target_scope().as_str(); + Self { + id: value.public_id.to_string(), + severity: value.severity.as_str(), + text: value.message, + max_displays: value.max_displays, + target_scope, + community_ids: value + .community_ids + .into_iter() + .map(|id| *id.as_uuid()) + .collect(), + created_by: hex::encode(value.created_by), + created_at: value.created_at, + disabled_at: value.disabled_at, + } + } +} + +async fn list_banner_communities( + State(state): State>, + uri: Uri, + headers: HeaderMap, +) -> Result>, ApiError> { + let principal = authorize( + &state, + &headers, + uri.path_and_query() + .map_or_else(|| uri.path(), |pq| pq.as_str()), + "GET", + None, + ) + .await?; + let principal = require_mutation_principal(principal)?; + require_operator(&principal)?; + let communities = state.db.admin_list_banner_communities().await?; + Ok(Json( + communities + .into_iter() + .map(|community| BannerCommunityResponse { + id: *community.id.as_uuid(), + host: community.host, + }) + .collect(), + )) +} + +async fn banners( + State(state): State>, + uri: Uri, + headers: HeaderMap, + Query(query): Query, +) -> Result, ApiError> { + let principal = authorize( + &state, + &headers, + uri.path_and_query() + .map_or_else(|| uri.path(), |pq| pq.as_str()), + "GET", + None, + ) + .await?; + let principal = require_mutation_principal(principal)?; + require_operator(&principal)?; + let history = state + .db + .admin_list_relay_banners(limit(query.limit)?) + .await?; + let active = history + .iter() + .find(|banner| banner.disabled_at.is_none()) + .cloned() + .map(AdminBannerResponse::from); + Ok(Json(BannersResponse { + active, + history: history.into_iter().map(AdminBannerResponse::from).collect(), + })) +} + +async fn upsert_banner( + State(state): State>, + uri: Uri, + headers: HeaderMap, + body_bytes: Bytes, +) -> Result, ApiError> { + let principal = require_mutation_principal( + authorize( + &state, + &headers, + uri.path_and_query() + .map_or_else(|| uri.path(), |pq| pq.as_str()), + "PUT", + Some(&body_bytes), + ) + .await?, + )?; + require_operator(&principal)?; + let body: UpsertBannerBody = serde_json::from_slice(&body_bytes) + .map_err(|_| ApiError::bad_request("invalid_body", "invalid JSON body"))?; + let scope = if body.target_scope == buzz_db::RelayBannerTargetScope::All { + if body + .community_ids + .as_ref() + .is_some_and(|ids| !ids.is_empty()) + { + return Err(ApiError::bad_request( + "invalid_scope", + "communityIds must be empty when targetScope is all", + )); + } + buzz_db::RelayBannerScope::AllCommunities + } else { + let ids = body.community_ids.unwrap_or_default(); + if ids.is_empty() { + return Err(ApiError::bad_request( + "invalid_scope", + "communityIds must be non-empty when targetScope is communities", + )); + } + buzz_db::RelayBannerScope::Communities( + ids.into_iter() + .map(buzz_core::CommunityId::from_uuid) + .collect(), + ) + }; + + let banner = state + .db + .admin_upsert_relay_banner(buzz_db::RelayBannerUpsert { + severity: body.severity, + message: body.text, + max_displays: body.max_displays, + scope, + actor_pubkey: principal.pubkey.to_vec(), + }) + .await + .map_err(map_banner_db_error)?; + Ok(Json(AdminBannerResponse::from(banner))) +} + +async fn disable_banner( + State(state): State>, + uri: Uri, + headers: HeaderMap, +) -> Result, ApiError> { + let principal = require_mutation_principal( + authorize( + &state, + &headers, + uri.path_and_query() + .map_or_else(|| uri.path(), |pq| pq.as_str()), + "DELETE", + None, + ) + .await?, + )?; + require_operator(&principal)?; + let disabled = state + .db + .admin_disable_active_relay_banner(&principal.pubkey) + .await?; + Ok(Json(serde_json::json!({ "disabled": disabled }))) +} + +fn map_banner_db_error(error: buzz_db::DbError) -> ApiError { + match error { + buzz_db::DbError::InvalidData(message) => ApiError::bad_request("invalid_banner", &message), + _ => ApiError::internal(), + } +} + // ── Phase 2: Report resolution ──────────────────────────────────────────────── /// Request body for POST /reports/{id}/resolve. diff --git a/crates/buzz-relay/src/api/banners.rs b/crates/buzz-relay/src/api/banners.rs new file mode 100644 index 00000000000..8a680274201 --- /dev/null +++ b/crates/buzz-relay/src/api/banners.rs @@ -0,0 +1,373 @@ +//! Client-facing relay banner API. + +use std::sync::Arc; + +use axum::{ + body::Bytes, + extract::{Path, State}, + http::HeaderMap, + response::Json, +}; +use buzz_core::TenantContext; +use serde::Serialize; +use serde_json::Value; + +use crate::api::{api_error, bridge, internal_error, relay_members}; +use crate::state::AppState; + +/// Client route for retrieving the active banner for the authenticated user. +pub(crate) const BANNER_ACTIVE_PATH: &str = "/api/banners/current"; +/// Client route for acknowledging a successful banner render/view. +pub(crate) const BANNER_VIEW_ROUTE: &str = "/api/banners/{banner_id}/view"; +/// Client route for dismissing a banner. +pub(crate) const BANNER_DISMISS_ROUTE: &str = "/api/banners/{banner_id}/dismiss"; + +#[derive(Serialize)] +#[serde(rename_all = "camelCase")] +pub(crate) struct BannerResponse { + /// Public banner id. + pub id: String, + /// Banner event kind clients may also query over Nostr. + pub kind: u32, + /// Banner severity/type. + pub severity: &'static str, + /// Plain-text banner text. + pub text: String, + /// Number of successful views each user may receive before suppression. + pub max_displays: i32, + /// Every v1 severity is dismissible. + pub dismissible: bool, +} + +#[derive(Serialize)] +#[serde(rename_all = "camelCase")] +pub(crate) struct BannerAckResponse { + status: &'static str, +} + +impl From for BannerResponse { + fn from(value: buzz_db::RelayBannerRecord) -> Self { + Self { + id: value.public_id.to_string(), + kind: buzz_core::kind::KIND_RELAY_BANNER, + severity: value.severity.as_str(), + text: value.message, + max_displays: value.max_displays, + dismissible: true, + } + } +} + +pub(crate) fn banner_event( + relay_keypair: &nostr::Keys, + banner: &buzz_db::RelayBannerRecord, +) -> Result { + let id = banner.public_id.to_string(); + let scope = if banner.target_all_communities { + "all" + } else { + "communities" + }; + let content = serde_json::json!({ + "id": id, + "severity": banner.severity.as_str(), + "text": banner.message, + "maxDisplays": banner.max_displays, + "dismissible": true, + }) + .to_string(); + let tags = vec![ + nostr::Tag::parse(["d", id.as_str()]).map_err(|e| e.to_string())?, + nostr::Tag::parse(["scope", scope]).map_err(|e| e.to_string())?, + ]; + nostr::EventBuilder::new( + nostr::Kind::Custom(buzz_core::kind::KIND_RELAY_BANNER as u16), + content, + ) + .tags(tags) + .sign_with_keys(relay_keypair) + .map_err(|e| e.to_string()) +} + +pub(crate) async fn active_banner_event_for_user( + state: &AppState, + tenant: &TenantContext, + pubkey_bytes: &[u8], +) -> Result, String> { + let Some(banner) = state + .db + .active_relay_banner_for_user(tenant.community(), pubkey_bytes) + .await + .map_err(|e| e.to_string())? + else { + return Ok(None); + }; + banner_event(&state.relay_keypair, &banner).map(Some) +} + +pub(crate) async fn get_active_banner( + State(state): State>, + headers: HeaderMap, +) -> Result>, (axum::http::StatusCode, Json)> { + let (tenant, pubkey, event_id, signed_created_at) = + authenticate_client_request(&state, &headers, "GET", BANNER_ACTIVE_PATH, None, false) + .await?; + bridge::check_nip98_replay(&state, &tenant, event_id).await?; + enforce_member(&state, &tenant, &headers, &pubkey, signed_created_at).await?; + + let banner = state + .db + .active_relay_banner_for_user(tenant.community(), pubkey.as_bytes()) + .await + .map_err(|e| internal_error(&format!("banner lookup: {e}")))? + .map(BannerResponse::from); + Ok(Json(banner)) +} + +pub(crate) async fn ack_banner_view( + State(state): State>, + Path(banner_id): Path, + headers: HeaderMap, + body: Bytes, +) -> Result, (axum::http::StatusCode, Json)> { + let (tenant, pubkey, event_id, signed_created_at) = authenticate_client_request( + &state, + &headers, + "POST", + &format!("/api/banners/{banner_id}/view"), + Some(&body), + true, + ) + .await?; + bridge::check_nip98_replay(&state, &tenant, event_id).await?; + enforce_member(&state, &tenant, &headers, &pubkey, signed_created_at).await?; + if !body.is_empty() { + return Err(non_empty_ack_body_error()); + } + match state + .db + .ack_relay_banner_view(tenant.community(), banner_id, pubkey.as_bytes()) + .await + .map_err(|e| internal_error(&format!("banner view ack: {e}")))? + { + buzz_db::RelayBannerViewOutcome::Accepted { .. } => { + metrics::counter!( + "buzz_relay_banner_views_total", + "community" => tenant.host().to_owned() + ) + .increment(1); + Ok(Json(BannerAckResponse { status: "accepted" })) + } + buzz_db::RelayBannerViewOutcome::Dismissed => Ok(Json(BannerAckResponse { + status: "dismissed", + })), + buzz_db::RelayBannerViewOutcome::Exhausted { .. } => Ok(Json(BannerAckResponse { + status: "exhausted", + })), + buzz_db::RelayBannerViewOutcome::NotEligible => Err(api_error( + axum::http::StatusCode::NOT_FOUND, + "banner not eligible", + )), + } +} + +pub(crate) async fn ack_banner_dismiss( + State(state): State>, + Path(banner_id): Path, + headers: HeaderMap, + body: Bytes, +) -> Result, (axum::http::StatusCode, Json)> { + let (tenant, pubkey, event_id, signed_created_at) = authenticate_client_request( + &state, + &headers, + "POST", + &format!("/api/banners/{banner_id}/dismiss"), + Some(&body), + true, + ) + .await?; + bridge::check_nip98_replay(&state, &tenant, event_id).await?; + enforce_member(&state, &tenant, &headers, &pubkey, signed_created_at).await?; + if !body.is_empty() { + return Err(non_empty_ack_body_error()); + } + match state + .db + .ack_relay_banner_dismiss(tenant.community(), banner_id, pubkey.as_bytes()) + .await + .map_err(|e| internal_error(&format!("banner dismiss ack: {e}")))? + { + buzz_db::RelayBannerDismissOutcome::Dismissed { changed } => { + if changed { + metrics::counter!( + "buzz_relay_banner_dismissals_total", + "community" => tenant.host().to_owned() + ) + .increment(1); + } + Ok(Json(BannerAckResponse { + status: "dismissed", + })) + } + buzz_db::RelayBannerDismissOutcome::NotEligible => Err(api_error( + axum::http::StatusCode::NOT_FOUND, + "banner not eligible", + )), + } +} + +async fn authenticate_client_request( + state: &AppState, + headers: &HeaderMap, + method: &str, + path: &str, + body: Option<&[u8]>, + require_payload: bool, +) -> Result< + (TenantContext, nostr::PublicKey, [u8; 32], Option), + (axum::http::StatusCode, Json), +> { + let raw_host = headers + .get(axum::http::header::HOST) + .and_then(|v| v.to_str().ok()) + .unwrap_or(""); + let tenant = crate::tenant::bind_community(&state.db, raw_host) + .await + .map_err(|_| { + api_error( + axum::http::StatusCode::NOT_FOUND, + "relay: no community is configured for this host", + ) + })?; + let url = bridge::nip98_expected_url(&state.config.relay_url, &tenant, path); + let verified = bridge::verify_bridge_auth_with_options( + headers, + method, + &url, + body, + state.config.require_auth_token, + require_payload, + )?; + bridge::enforce_http_admission(state, &tenant, &verified.pubkey).await?; + Ok(( + tenant, + verified.pubkey, + verified.event_id_bytes, + verified.signed_created_at, + )) +} + +fn non_empty_ack_body_error() -> (axum::http::StatusCode, Json) { + api_error( + axum::http::StatusCode::BAD_REQUEST, + "banner acknowledgement body must be empty", + ) +} + +#[cfg(test)] +mod tests { + use super::*; + + fn banner_record(target_all_communities: bool) -> buzz_db::RelayBannerRecord { + buzz_db::RelayBannerRecord { + id: 42, + public_id: uuid::Uuid::from_u128(0x12345678123456781234567812345678), + severity: buzz_db::RelayBannerSeverity::Warning, + message: "scheduled maintenance".to_owned(), + max_displays: 3, + target_all_communities, + community_ids: Vec::new(), + created_at: chrono::DateTime::UNIX_EPOCH, + updated_at: chrono::DateTime::UNIX_EPOCH, + created_by: vec![7; 32], + disabled_at: None, + } + } + + #[test] + fn banner_event_uses_published_contract() { + let keys = nostr::Keys::generate(); + let banner = banner_record(false); + let event = banner_event(&keys, &banner).expect("banner event"); + + assert_eq!( + event.kind.as_u16() as u32, + buzz_core::kind::KIND_RELAY_BANNER + ); + assert_eq!(buzz_core::kind::KIND_RELAY_BANNER, 13536); + assert_eq!( + event + .tags + .iter() + .find(|tag| tag.kind().to_string() == "d") + .and_then(|tag| tag.content()), + Some(banner.public_id.to_string().as_str()) + ); + assert_eq!( + event + .tags + .iter() + .find(|tag| tag.kind().to_string() == "scope") + .and_then(|tag| tag.content()), + Some("communities") + ); + + let content: serde_json::Value = + serde_json::from_str(&event.content).expect("banner JSON content"); + assert_eq!(content["id"], banner.public_id.to_string()); + assert_eq!(content["severity"], "warning"); + assert_eq!(content["text"], "scheduled maintenance"); + assert_eq!(content["maxDisplays"], 3); + assert_eq!(content["dismissible"], true); + assert!(content.get("message").is_none()); + } + + #[test] + fn all_community_banner_event_has_all_scope() { + let keys = nostr::Keys::generate(); + let banner = banner_record(true); + let event = banner_event(&keys, &banner).expect("banner event"); + + assert_eq!( + event + .tags + .iter() + .find(|tag| tag.kind().to_string() == "scope") + .and_then(|tag| tag.content()), + Some("all") + ); + } + + #[test] + fn banner_ack_routes_are_uuid_path_routes() { + assert_eq!(BANNER_VIEW_ROUTE, "/api/banners/{banner_id}/view"); + assert_eq!(BANNER_DISMISS_ROUTE, "/api/banners/{banner_id}/dismiss"); + } + + #[test] + fn non_empty_ack_body_is_rejected() { + let (status, body) = non_empty_ack_body_error(); + + assert_eq!(status, axum::http::StatusCode::BAD_REQUEST); + assert_eq!(body["error"], "banner acknowledgement body must be empty"); + } +} + +async fn enforce_member( + state: &AppState, + tenant: &TenantContext, + headers: &HeaderMap, + pubkey: &nostr::PublicKey, + signed_created_at: Option, +) -> Result<(), (axum::http::StatusCode, Json)> { + let auth_tag = relay_members::extract_auth_tag_header(headers); + relay_members::enforce_relay_membership( + state, + tenant.community(), + pubkey.as_bytes(), + auth_tag, + signed_created_at, + ) + .await + .map(|_| ()) +} diff --git a/crates/buzz-relay/src/api/bridge.rs b/crates/buzz-relay/src/api/bridge.rs index 37c549610de..041a730bba3 100644 --- a/crates/buzz-relay/src/api/bridge.rs +++ b/crates/buzz-relay/src/api/bridge.rs @@ -1182,6 +1182,27 @@ async fn query_events_authed( return presence_result.map(|events| Json(Value::Array(events))); } + if filters_are_relay_banner_only(&filters) { + let mut events = Vec::new(); + if let Some(event) = + crate::api::banners::active_banner_event_for_user(state, tenant, &pubkey_bytes) + .await + .map_err(|e| internal_error(&format!("banner lookup: {e}")))? + { + let stored = + buzz_core::StoredEvent::with_received_at(event, chrono::Utc::now(), None, true); + if filters.iter().any(|filter| { + buzz_core::filter::filters_match(std::slice::from_ref(filter), &stored) + }) { + events.push( + serde_json::to_value(&stored.event) + .map_err(|e| internal_error(&format!("banner serialize: {e}")))?, + ); + } + } + return Ok(Json(Value::Array(events))); + } + let mut events: Vec = Vec::new(); let mut handled: std::collections::HashSet = std::collections::HashSet::new(); @@ -2329,6 +2350,18 @@ async fn synthesize_presence( Some(Ok(events)) } +fn filters_are_relay_banner_only(filters: &[nostr::Filter]) -> bool { + !filters.is_empty() + && filters.iter().all(|filter| { + filter.kinds.as_ref().is_some_and(|kinds| { + kinds.len() == 1 + && kinds + .iter() + .all(|kind| kind.as_u16() as u32 == buzz_core::kind::KIND_RELAY_BANNER) + }) + }) +} + // ── Moderation queue reads (L6 — Quinn) ─────────────────────────────────────── // // Mod-only structured rows (`moderation_reports`/`moderation_actions`/ diff --git a/crates/buzz-relay/src/api/mod.rs b/crates/buzz-relay/src/api/mod.rs index 5745b8d4e59..8b9bd280ad6 100644 --- a/crates/buzz-relay/src/api/mod.rs +++ b/crates/buzz-relay/src/api/mod.rs @@ -1,6 +1,7 @@ //! HTTP API — media, git, NIP-05, and the Nostr HTTP bridge. pub mod admin; +pub mod banners; pub mod bridge; pub mod events; pub mod gifs; diff --git a/crates/buzz-relay/src/handlers/req.rs b/crates/buzz-relay/src/handlers/req.rs index 97d252cb7e8..ded8348a583 100644 --- a/crates/buzz-relay/src/handlers/req.rs +++ b/crates/buzz-relay/src/handlers/req.rs @@ -8,7 +8,7 @@ use tracing::{debug, warn}; use buzz_core::filter::filters_match; use buzz_core::kind::{ is_unshared_gated_event, AUTHOR_ONLY_KINDS, KIND_AGENT_ENGRAM, KIND_AGENT_TURN_METRIC, - KIND_DM_VISIBILITY, KIND_HUDDLE_LIVENESS, P_GATED_KINDS, RESULT_GATED_KINDS, + KIND_DM_VISIBILITY, KIND_HUDDLE_LIVENESS, KIND_RELAY_BANNER, P_GATED_KINDS, RESULT_GATED_KINDS, SHARED_GATED_KINDS, }; use buzz_core::tenant::TenantContext; @@ -221,6 +221,30 @@ pub async fn handle_req( return; } + if filters_are_relay_banner_only(&filters) { + match crate::api::banners::active_banner_event_for_user(&state, &conn.tenant, &pubkey_bytes) + .await + { + Ok(Some(event)) => { + let stored = + buzz_core::StoredEvent::with_received_at(event, chrono::Utc::now(), None, true); + if filters_match(&filters, &stored) + && !conn.send(RelayMessage::event(&sub_id, &stored.event)) + { + return; + } + } + Ok(None) => {} + Err(error) => { + warn!(conn_id = %conn_id, sub_id = %sub_id, "Relay banner lookup failed: {error}"); + conn.send(RelayMessage::closed(&sub_id, "error: database error")); + return; + } + } + conn.send(RelayMessage::eose(&sub_id)); + return; + } + // Applied BEFORE the NIP-50 search branch so that an authenticated member // cannot use `{"search":"...","kinds":[30174]}` (or similar for p-gated // kinds) to harvest indexed-but-globally-stored sensitive events. Search @@ -1159,6 +1183,18 @@ fn filters_are_huddle_liveness_only(filters: &[Filter]) -> bool { }) } +fn filters_are_relay_banner_only(filters: &[Filter]) -> bool { + !filters.is_empty() + && filters.iter().all(|filter| { + filter.kinds.as_ref().is_some_and(|kinds| { + kinds.len() == 1 + && kinds + .iter() + .all(|kind| kind.as_u16() as u32 == KIND_RELAY_BANNER) + }) + }) +} + fn huddle_liveness_session_ids(filters: &[Filter]) -> Vec { let d_tag = nostr::SingleLetterTag::lowercase(nostr::Alphabet::D); let mut session_ids = Vec::new(); diff --git a/crates/buzz-relay/src/router.rs b/crates/buzz-relay/src/router.rs index 61aedf70be0..eabfbb80848 100644 --- a/crates/buzz-relay/src/router.rs +++ b/crates/buzz-relay/src/router.rs @@ -73,6 +73,18 @@ pub fn build_router(state: Arc) -> Router { .route("/events", post(api::bridge::submit_event)) .route("/query", post(api::bridge::query_events)) .route("/count", post(api::bridge::count_events)) + .route( + api::banners::BANNER_ACTIVE_PATH, + get(api::banners::get_active_banner), + ) + .route( + api::banners::BANNER_VIEW_ROUTE, + post(api::banners::ack_banner_view), + ) + .route( + api::banners::BANNER_DISMISS_ROUTE, + post(api::banners::ack_banner_dismiss), + ) // Relay-owned third-party GIF metadata proxy (NIP-98 auth). .route(api::gifs::SEARCH_PATH, post(api::gifs::search)) .route(api::gifs::SHARE_PATH, post(api::gifs::share)) diff --git a/migrations/0046_relay_banners.sql b/migrations/0046_relay_banners.sql new file mode 100644 index 00000000000..2be963e9ab8 --- /dev/null +++ b/migrations/0046_relay_banners.sql @@ -0,0 +1,62 @@ +-- Operator-configured deployment banners. +-- Exactly one active banner exists deployment-wide. Scope is either all +-- communities (target_all_communities=true) or an explicit non-empty set in +-- relay_banner_communities. + +CREATE TABLE relay_banners ( + id BIGSERIAL PRIMARY KEY, + public_id UUID NOT NULL DEFAULT gen_random_uuid(), + severity TEXT NOT NULL CHECK (severity IN ('info', 'warning', 'urgent')), + message TEXT NOT NULL CHECK (char_length(message) BETWEEN 1 AND 2000), + max_displays INTEGER NOT NULL CHECK (max_displays >= 1), + target_all_communities BOOLEAN NOT NULL, + created_at TIMESTAMPTZ NOT NULL DEFAULT now(), + updated_at TIMESTAMPTZ NOT NULL DEFAULT now(), + disabled_at TIMESTAMPTZ, + created_by BYTEA NOT NULL CHECK (length(created_by) = 32), + updated_by BYTEA NOT NULL CHECK (length(updated_by) = 32), + disabled_by BYTEA CHECK (disabled_by IS NULL OR length(disabled_by) = 32) +); + +CREATE UNIQUE INDEX idx_relay_banners_public_id + ON relay_banners (public_id); + +CREATE UNIQUE INDEX idx_relay_banners_one_active + ON relay_banners ((true)) + WHERE disabled_at IS NULL; + +CREATE INDEX idx_relay_banners_updated_at + ON relay_banners (updated_at DESC, id DESC); + +CREATE TABLE relay_banner_communities ( + banner_id BIGINT NOT NULL REFERENCES relay_banners(id) ON DELETE CASCADE, + community_id UUID NOT NULL REFERENCES communities(id) ON DELETE CASCADE, + PRIMARY KEY (banner_id, community_id) +); + +CREATE INDEX idx_relay_banner_communities_community + ON relay_banner_communities (community_id, banner_id); + +CREATE TABLE relay_banner_user_state ( + banner_id BIGINT NOT NULL REFERENCES relay_banners(id) ON DELETE CASCADE, + community_id UUID NOT NULL REFERENCES communities(id) ON DELETE CASCADE, + pubkey BYTEA NOT NULL CHECK (length(pubkey) = 32), + display_count INTEGER NOT NULL DEFAULT 0 CHECK (display_count >= 0), + dismissed_at TIMESTAMPTZ, + first_viewed_at TIMESTAMPTZ, + last_viewed_at TIMESTAMPTZ, + updated_at TIMESTAMPTZ NOT NULL DEFAULT now(), + PRIMARY KEY (banner_id, community_id, pubkey), + CHECK (first_viewed_at IS NULL OR last_viewed_at IS NOT NULL) +); + +CREATE INDEX idx_relay_banner_user_state_pubkey + ON relay_banner_user_state (community_id, pubkey, banner_id); + +INSERT INTO _operator_global_tables (table_name, reason) VALUES + ('relay_banners', 'deployment-global operator-configured banner; no community_id intentionally'), + ('relay_banner_communities', 'deployment-global banner targeting allowlist; community_id is target provenance only'), + ('relay_banner_user_state', 'deployment-global per-user banner display state; community_id scopes the targeted impression only'); + +SELECT attach_community_write_fence('relay_banner_communities'); +SELECT attach_community_write_fence('relay_banner_user_state'); diff --git a/schema/schema.sql b/schema/schema.sql index d7074f359d8..6e4e3750889 100644 --- a/schema/schema.sql +++ b/schema/schema.sql @@ -1700,6 +1700,69 @@ BEGIN END $$; + + +-- Operator-configured deployment banners. +-- Exactly one active banner exists deployment-wide. Scope is either all +-- communities (target_all_communities=true) or an explicit non-empty set in +-- relay_banner_communities. + +CREATE TABLE relay_banners ( + id BIGSERIAL PRIMARY KEY, + public_id UUID NOT NULL DEFAULT gen_random_uuid(), + severity TEXT NOT NULL CHECK (severity IN ('info', 'warning', 'urgent')), + message TEXT NOT NULL CHECK (char_length(message) BETWEEN 1 AND 2000), + max_displays INTEGER NOT NULL CHECK (max_displays >= 1), + target_all_communities BOOLEAN NOT NULL, + created_at TIMESTAMPTZ NOT NULL DEFAULT now(), + updated_at TIMESTAMPTZ NOT NULL DEFAULT now(), + disabled_at TIMESTAMPTZ, + created_by BYTEA NOT NULL CHECK (length(created_by) = 32), + updated_by BYTEA NOT NULL CHECK (length(updated_by) = 32), + disabled_by BYTEA CHECK (disabled_by IS NULL OR length(disabled_by) = 32) +); + +CREATE UNIQUE INDEX idx_relay_banners_public_id + ON relay_banners (public_id); + +CREATE UNIQUE INDEX idx_relay_banners_one_active + ON relay_banners ((true)) + WHERE disabled_at IS NULL; + +CREATE INDEX idx_relay_banners_updated_at + ON relay_banners (updated_at DESC, id DESC); + +CREATE TABLE relay_banner_communities ( + banner_id BIGINT NOT NULL REFERENCES relay_banners(id) ON DELETE CASCADE, + community_id UUID NOT NULL REFERENCES communities(id) ON DELETE CASCADE, + PRIMARY KEY (banner_id, community_id) +); + +CREATE INDEX idx_relay_banner_communities_community + ON relay_banner_communities (community_id, banner_id); + +CREATE TABLE relay_banner_user_state ( + banner_id BIGINT NOT NULL REFERENCES relay_banners(id) ON DELETE CASCADE, + community_id UUID NOT NULL REFERENCES communities(id) ON DELETE CASCADE, + pubkey BYTEA NOT NULL CHECK (length(pubkey) = 32), + display_count INTEGER NOT NULL DEFAULT 0 CHECK (display_count >= 0), + dismissed_at TIMESTAMPTZ, + first_viewed_at TIMESTAMPTZ, + last_viewed_at TIMESTAMPTZ, + updated_at TIMESTAMPTZ NOT NULL DEFAULT now(), + PRIMARY KEY (banner_id, community_id, pubkey), + CHECK (first_viewed_at IS NULL OR last_viewed_at IS NOT NULL) +); + +CREATE INDEX idx_relay_banner_user_state_pubkey + ON relay_banner_user_state (community_id, pubkey, banner_id); + +INSERT INTO _operator_global_tables (table_name, reason) VALUES + ('relay_banners', 'deployment-global operator-configured banner; no community_id intentionally'), + ('relay_banner_communities', 'deployment-global banner targeting allowlist; community_id is target provenance only'), + ('relay_banner_user_state', 'deployment-global per-user banner display state; community_id scopes the targeted impression only'); + + -- Attach the universal fence to every existing table carrying community_id, -- including deployment-private sidecars whose community_id is provenance. DO $$ @@ -1749,6 +1812,8 @@ SELECT attach_community_write_fence('push_wake_outbox'); SELECT attach_community_write_fence('reactions'); SELECT attach_community_write_fence('relay_invites'); SELECT attach_community_write_fence('relay_members'); +SELECT attach_community_write_fence('relay_banner_communities'); +SELECT attach_community_write_fence('relay_banner_user_state'); SELECT attach_community_write_fence('scheduled_workflow_fires'); SELECT attach_community_write_fence('subscriptions'); SELECT attach_community_write_fence('thread_metadata'); @@ -1761,7 +1826,6 @@ SELECT attach_community_write_fence('workflows'); -- Deployment-level principals staffed via the admin API. Config-backed operators -- (RELAY_OPERATOR_PUBKEYS, RELAY_OWNER_PUBKEY owner-fallback) are NOT seeded here; -- they are authoritative in config and outrank any DB row. - CREATE TABLE relay_operators ( pubkey BYTEA NOT NULL PRIMARY KEY CHECK (length(pubkey) = 32), role TEXT NOT NULL CHECK (role IN ('operator', 'moderator')), From bd4bb3bae6fc69032d4103aae5ba3ce84b1d3e27 Mon Sep 17 00:00:00 2001 From: coder 0 Date: Thu, 17 Sep 2026 11:00:05 -0400 Subject: [PATCH 2/4] Fix relay banner live delivery and view retries Signed-off-by: coder 0 --- crates/buzz-db/src/runtime/migration.rs | 3 + crates/buzz-db/src/store/relay_banners.rs | 437 +++++++++++++++++++++- crates/buzz-relay/src/api/admin/mod.rs | 112 +++++- crates/buzz-relay/src/api/banners.rs | 172 ++++++++- crates/buzz-relay/src/handlers/event.rs | 279 +++++++++++++- crates/buzz-relay/src/handlers/req.rs | 108 ++++-- migrations/0046_relay_banners.sql | 16 +- schema/schema.sql | 16 +- 8 files changed, 1068 insertions(+), 75 deletions(-) diff --git a/crates/buzz-db/src/runtime/migration.rs b/crates/buzz-db/src/runtime/migration.rs index e984b2a836f..2564f116c38 100644 --- a/crates/buzz-db/src/runtime/migration.rs +++ b/crates/buzz-db/src/runtime/migration.rs @@ -493,6 +493,7 @@ mod postgres_tests { "relay_banners", "relay_banner_communities", "relay_banner_user_state", + "relay_banner_view_acks", ] { if normalized[insert_pos..].contains(&format!("'{value}'")) { globals.insert(value.to_owned()); @@ -1833,6 +1834,7 @@ mod postgres_tests { expected_fences.remove("rate_limit_violations"); expected_fences.insert("relay_banner_communities".to_owned()); expected_fences.insert("relay_banner_user_state".to_owned()); + expected_fences.insert("relay_banner_view_acks".to_owned()); assert_eq!( expected_fences, schema.fence_attachments, "write-fence attachment targets differ after recovery policy" @@ -2275,6 +2277,7 @@ mod postgres_tests { "relay_banners", "relay_banner_communities", "relay_banner_user_state", + "relay_banner_view_acks", ] { assert_eq!( columns(&desired, table).await, diff --git a/crates/buzz-db/src/store/relay_banners.rs b/crates/buzz-db/src/store/relay_banners.rs index df623a17cf9..2d4e71660a8 100644 --- a/crates/buzz-db/src/store/relay_banners.rs +++ b/crates/buzz-db/src/store/relay_banners.rs @@ -9,6 +9,7 @@ use buzz_datastore_tracing::datastore_span; use chrono::{DateTime, Utc}; use serde::{Deserialize, Serialize}; use sqlx::{Postgres, Row as _, Transaction}; +use std::collections::HashSet; use uuid::Uuid; use crate::{Db, Result}; @@ -139,6 +140,8 @@ pub enum RelayBannerViewOutcome { Accepted { /// Display count after accepting this view. display_count: i32, + /// True only when this acknowledgement consumed a new display. + changed: bool, }, /// The banner was already permanently dismissed by this user. Dismissed, @@ -188,6 +191,12 @@ impl RelayBannerUpsert { "relay banner community scope must be non-empty".to_owned(), )); } + let mut seen = HashSet::with_capacity(ids.len()); + if ids.iter().any(|id| !seen.insert(*id)) { + return Err(crate::DbError::InvalidData( + "relay banner community scope must not contain duplicates".to_owned(), + )); + } } Ok(()) } @@ -203,16 +212,18 @@ async fn acquire_banner_lock(tx: &mut Transaction<'_, Postgres>) -> Result<()> { async fn banner_internal_id_for_public_id( tx: &mut Transaction<'_, Postgres>, public_id: Uuid, + require_active: bool, ) -> Result> { sqlx::query_scalar::<_, i64>( r#" SELECT id FROM relay_banners WHERE public_id = $1 - AND disabled_at IS NULL + AND (NOT $2 OR disabled_at IS NULL) "#, ) .bind(public_id) + .bind(require_active) .fetch_optional(&mut **tx) .await .map_err(Into::into) @@ -249,6 +260,47 @@ async fn banner_targets_community( Ok(eligible) } +async fn relay_banner_user_eligible_in_tx( + tx: &mut Transaction<'_, Postgres>, + banner_id: i64, + community_id: CommunityId, + pubkey: &[u8], + require_active: bool, +) -> Result { + let eligible = sqlx::query_scalar::<_, bool>( + r#" + SELECT EXISTS ( + SELECT 1 + FROM relay_banners b + LEFT JOIN relay_banner_user_state us + ON us.banner_id = b.id + AND us.community_id = $2 + AND us.pubkey = $3 + WHERE b.id = $1 + AND (NOT $4 OR b.disabled_at IS NULL) + AND ( + b.target_all_communities + OR EXISTS ( + SELECT 1 + FROM relay_banner_communities bc + WHERE bc.banner_id = b.id + AND bc.community_id = $2 + ) + ) + AND us.dismissed_at IS NULL + AND COALESCE(us.display_count, 0) < b.max_displays + ) + "#, + ) + .bind(banner_id) + .bind(community_id.as_uuid()) + .bind(pubkey) + .bind(require_active) + .fetch_one(&mut **tx) + .await?; + Ok(eligible) +} + impl Db { /// Lists active communities for the admin banner community picker. #[datastore_span(name = "admin_list_banner_communities", system = "postgresql")] @@ -318,6 +370,44 @@ impl Db { rows.into_iter().map(row_to_banner).collect() } + /// Returns a banner by its public id, including disabled history. + #[datastore_span(name = "relay_banner_by_public_id", system = "postgresql")] + pub async fn relay_banner_by_public_id( + &self, + public_id: Uuid, + ) -> Result> { + let mut connection = crate::observability::acquire_writer( + &self.pool, + crate::observability::WriterOperation::Authorization, + ) + .await?; + let row = sqlx::query( + r#" + SELECT + b.id, + b.public_id, + b.severity, + b.message, + b.max_displays, + b.target_all_communities, + b.created_by, + b.created_at, + b.updated_at, + b.disabled_at, + COALESCE(array_agg(bc.community_id ORDER BY bc.community_id) + FILTER (WHERE bc.community_id IS NOT NULL), ARRAY[]::uuid[]) AS community_ids + FROM relay_banners b + LEFT JOIN relay_banner_communities bc ON bc.banner_id = b.id + WHERE b.public_id = $1 + GROUP BY b.id + "#, + ) + .bind(public_id) + .fetch_optional(&mut *connection) + .await?; + row.map(row_to_banner).transpose() + } + /// Returns the currently active deployment banner, if any. #[datastore_span(name = "admin_get_active_relay_banner", system = "postgresql")] pub async fn admin_get_active_relay_banner(&self) -> Result> { @@ -439,7 +529,10 @@ impl Db { /// Disables the currently active deployment banner, preserving history and state. #[datastore_span(name = "admin_disable_active_relay_banner", system = "postgresql")] - pub async fn admin_disable_active_relay_banner(&self, actor_pubkey: &[u8]) -> Result { + pub async fn admin_disable_active_relay_banner( + &self, + actor_pubkey: &[u8], + ) -> Result> { let connection = crate::observability::acquire_writer( &self.pool, crate::observability::WriterOperation::Authorization, @@ -447,18 +540,79 @@ impl Db { .await?; let mut tx = sqlx::Transaction::begin(connection, None).await?; acquire_banner_lock(&mut tx).await?; - let result = sqlx::query( + let Some(active) = sqlx::query( + r#" + SELECT + b.id, + b.public_id, + b.severity, + b.message, + b.max_displays, + b.target_all_communities, + b.created_by, + b.created_at, + b.updated_at, + b.disabled_at, + COALESCE(array_agg(bc.community_id ORDER BY bc.community_id) + FILTER (WHERE bc.community_id IS NOT NULL), ARRAY[]::uuid[]) AS community_ids + FROM relay_banners b + LEFT JOIN relay_banner_communities bc ON bc.banner_id = b.id + WHERE b.disabled_at IS NULL + GROUP BY b.id + ORDER BY b.id DESC + LIMIT 1 + "#, + ) + .fetch_optional(&mut *tx) + .await? + else { + tx.commit().await?; + return Ok(None); + }; + let mut active = row_to_banner(active)?; + let updated = sqlx::query( r#" UPDATE relay_banners SET disabled_at = now(), disabled_by = $1, updated_at = now(), updated_by = $1 - WHERE disabled_at IS NULL + WHERE id = $2 AND disabled_at IS NULL + RETURNING updated_at, disabled_at "#, ) .bind(actor_pubkey) - .execute(&mut *tx) + .bind(active.id) + .fetch_one(&mut *tx) .await?; + active.updated_at = updated.try_get("updated_at")?; + active.disabled_at = updated.try_get("disabled_at")?; tx.commit().await?; - Ok(result.rows_affected() > 0) + Ok(Some(active)) + } + + /// Returns whether this user would receive the banner in this community. + #[datastore_span(name = "relay_banner_user_eligible", system = "postgresql")] + pub async fn relay_banner_user_eligible( + &self, + banner_id: i64, + community_id: CommunityId, + pubkey: &[u8], + require_active: bool, + ) -> Result { + let connection = crate::observability::acquire_writer( + &self.pool, + crate::observability::WriterOperation::Authorization, + ) + .await?; + let mut tx = sqlx::Transaction::begin(connection, None).await?; + let eligible = relay_banner_user_eligible_in_tx( + &mut tx, + banner_id, + community_id, + pubkey, + require_active, + ) + .await?; + tx.commit().await?; + Ok(eligible) } /// Resolves the active banner eligible for this user in this community. @@ -522,6 +676,7 @@ impl Db { community_id: CommunityId, public_id: Uuid, pubkey: &[u8], + view_id: Uuid, ) -> Result { let connection = crate::observability::acquire_writer( &self.pool, @@ -529,7 +684,8 @@ impl Db { ) .await?; let mut tx = sqlx::Transaction::begin(connection, None).await?; - let Some(banner_id) = banner_internal_id_for_public_id(&mut tx, public_id).await? else { + let Some(banner_id) = banner_internal_id_for_public_id(&mut tx, public_id, true).await? + else { tx.rollback().await?; return Ok(RelayBannerViewOutcome::NotEligible); }; @@ -538,6 +694,60 @@ impl Db { return Ok(RelayBannerViewOutcome::NotEligible); } + let inserted_view = sqlx::query_scalar::<_, bool>( + r#" + INSERT INTO relay_banner_view_acks + (banner_id, community_id, pubkey, view_id) + SELECT $1, $2, $3, $4 + WHERE EXISTS ( + SELECT 1 + FROM relay_banners b + WHERE b.id = $1 + AND b.disabled_at IS NULL + AND (b.target_all_communities OR EXISTS ( + SELECT 1 + FROM relay_banner_communities bc + WHERE bc.banner_id = b.id AND bc.community_id = $2 + )) + ) + ON CONFLICT (banner_id, community_id, pubkey, view_id) DO NOTHING + RETURNING TRUE + "#, + ) + .bind(banner_id) + .bind(community_id.as_uuid()) + .bind(pubkey) + .bind(view_id) + .fetch_optional(&mut *tx) + .await? + .unwrap_or(false); + + if !inserted_view { + let state = sqlx::query( + r#" + SELECT us.display_count, us.dismissed_at IS NOT NULL AS dismissed + FROM relay_banner_user_state us + WHERE us.banner_id = $1 AND us.community_id = $2 AND us.pubkey = $3 + "#, + ) + .bind(banner_id) + .bind(community_id.as_uuid()) + .bind(pubkey) + .fetch_optional(&mut *tx) + .await?; + tx.commit().await?; + return Ok(match state { + Some(row) if row.try_get::("dismissed")? => { + RelayBannerViewOutcome::Dismissed + } + Some(row) => RelayBannerViewOutcome::Accepted { + display_count: row.try_get("display_count")?, + changed: false, + }, + None => RelayBannerViewOutcome::NotEligible, + }); + } + let row = sqlx::query( r#" INSERT INTO relay_banner_user_state @@ -568,6 +778,7 @@ impl Db { let outcome = if let Some(row) = row { RelayBannerViewOutcome::Accepted { display_count: row.try_get("display_count")?, + changed: true, } } else { let state = sqlx::query( @@ -592,6 +803,13 @@ impl Db { None => RelayBannerViewOutcome::NotEligible, } }; + if !matches!( + outcome, + RelayBannerViewOutcome::Accepted { changed: true, .. } + ) { + tx.rollback().await?; + return Ok(outcome); + } tx.commit().await?; Ok(outcome) } @@ -610,7 +828,8 @@ impl Db { ) .await?; let mut tx = sqlx::Transaction::begin(connection, None).await?; - let Some(banner_id) = banner_internal_id_for_public_id(&mut tx, public_id).await? else { + let Some(banner_id) = banner_internal_id_for_public_id(&mut tx, public_id, true).await? + else { tx.rollback().await?; return Ok(RelayBannerDismissOutcome::NotEligible); }; @@ -686,17 +905,21 @@ mod tests { Db::from_pool(pool) } - async fn make_community(pool: &PgPool) -> CommunityId { + pub(crate) async fn insert_test_community(pool: &PgPool, host_prefix: &str) -> CommunityId { let id = Uuid::new_v4(); sqlx::query("INSERT INTO communities (id, host) VALUES ($1, $2)") .bind(id) - .bind(format!("banner-{}.example", id.simple())) + .bind(format!("{host_prefix}-{}.example", id.simple())) .execute(pool) .await .expect("insert community"); CommunityId::from_uuid(id) } + async fn make_community(pool: &PgPool) -> CommunityId { + insert_test_community(pool, "banner").await + } + fn actor(byte: u8) -> Vec { vec![byte; 32] } @@ -770,6 +993,60 @@ mod tests { .is_none()); } + #[tokio::test] + #[ignore = "requires Postgres"] + async fn duplicate_community_scope_is_rejected() { + let db = setup_db().await; + let community = make_community(&db.pool).await; + let result = db + .admin_upsert_relay_banner(RelayBannerUpsert { + severity: RelayBannerSeverity::Info, + message: "duplicates".to_owned(), + max_displays: 1, + scope: RelayBannerScope::Communities(vec![community, community]), + actor_pubkey: actor(15), + }) + .await; + + assert!( + matches!(result, Err(crate::DbError::InvalidData(message)) if message.contains("duplicates")) + ); + } + + #[tokio::test] + #[ignore = "requires Postgres"] + async fn disabled_banner_remains_eligible_when_active_requirement_is_lifted() { + let db = setup_db().await; + let community = make_community(&db.pool).await; + let user = actor(16); + let banner = db + .admin_upsert_relay_banner(RelayBannerUpsert { + severity: RelayBannerSeverity::Warning, + message: "disable eligibility".to_owned(), + max_displays: 2, + scope: RelayBannerScope::AllCommunities, + actor_pubkey: actor(17), + }) + .await + .expect("insert banner"); + let disabled = db + .admin_disable_active_relay_banner(&actor(17)) + .await + .expect("disable banner") + .expect("active banner"); + + assert_eq!(disabled.id, banner.id); + assert!(disabled.disabled_at.is_some()); + assert!(db + .relay_banner_user_eligible(banner.id, community, &user, false) + .await + .expect("pre-disable eligibility lookup")); + assert!(!db + .relay_banner_user_eligible(banner.id, community, &user, true) + .await + .expect("active eligibility lookup")); + } + #[tokio::test] #[ignore = "requires Postgres"] async fn view_ack_consumes_display_and_then_exhausts() { @@ -806,13 +1083,16 @@ mod tests { "repeated delivery lookups without view ack must not exhaust max_displays" ); assert_eq!( - db.ack_relay_banner_view(community, banner.public_id, &user) + db.ack_relay_banner_view(community, banner.public_id, &user, Uuid::new_v4()) .await .expect("first view"), - RelayBannerViewOutcome::Accepted { display_count: 1 } + RelayBannerViewOutcome::Accepted { + display_count: 1, + changed: true + } ); assert_eq!( - db.ack_relay_banner_view(community, banner.public_id, &user) + db.ack_relay_banner_view(community, banner.public_id, &user, Uuid::new_v4()) .await .expect("second view"), RelayBannerViewOutcome::Exhausted { display_count: 1 } @@ -824,6 +1104,135 @@ mod tests { .is_none()); } + #[tokio::test] + #[ignore = "requires Postgres"] + async fn view_ack_same_key_retry_does_not_consume_again() { + let db = setup_db().await; + let community = make_community(&db.pool).await; + let user = actor(9); + let view_id = Uuid::new_v4(); + let banner = db + .admin_upsert_relay_banner(RelayBannerUpsert { + severity: RelayBannerSeverity::Info, + message: "retry".to_owned(), + max_displays: 2, + scope: RelayBannerScope::AllCommunities, + actor_pubkey: actor(10), + }) + .await + .expect("insert banner"); + + assert_eq!( + db.ack_relay_banner_view(community, banner.public_id, &user, view_id) + .await + .expect("first view"), + RelayBannerViewOutcome::Accepted { + display_count: 1, + changed: true + } + ); + assert_eq!( + db.ack_relay_banner_view(community, banner.public_id, &user, view_id) + .await + .expect("retry view"), + RelayBannerViewOutcome::Accepted { + display_count: 1, + changed: false + } + ); + assert!(db + .active_relay_banner_for_user(community, &user) + .await + .expect("still eligible after duplicate") + .is_some()); + } + + #[tokio::test] + #[ignore = "requires Postgres"] + async fn distinct_view_keys_increment_until_cap() { + let db = setup_db().await; + let community = make_community(&db.pool).await; + let user = actor(11); + let banner = db + .admin_upsert_relay_banner(RelayBannerUpsert { + severity: RelayBannerSeverity::Warning, + message: "cap".to_owned(), + max_displays: 2, + scope: RelayBannerScope::AllCommunities, + actor_pubkey: actor(12), + }) + .await + .expect("insert banner"); + + assert_eq!( + db.ack_relay_banner_view(community, banner.public_id, &user, Uuid::new_v4()) + .await + .expect("first view"), + RelayBannerViewOutcome::Accepted { + display_count: 1, + changed: true + } + ); + assert_eq!( + db.ack_relay_banner_view(community, banner.public_id, &user, Uuid::new_v4()) + .await + .expect("second view"), + RelayBannerViewOutcome::Accepted { + display_count: 2, + changed: true + } + ); + assert_eq!( + db.ack_relay_banner_view(community, banner.public_id, &user, Uuid::new_v4()) + .await + .expect("third view"), + RelayBannerViewOutcome::Exhausted { display_count: 2 } + ); + } + + #[tokio::test] + #[ignore = "requires Postgres"] + async fn concurrent_same_view_key_consumes_once() { + let db = setup_db().await; + let community = make_community(&db.pool).await; + let user = actor(13); + let view_id = Uuid::new_v4(); + let banner = db + .admin_upsert_relay_banner(RelayBannerUpsert { + severity: RelayBannerSeverity::Urgent, + message: "concurrent".to_owned(), + max_displays: 2, + scope: RelayBannerScope::AllCommunities, + actor_pubkey: actor(14), + }) + .await + .expect("insert banner"); + + let (first, second) = tokio::join!( + db.ack_relay_banner_view(community, banner.public_id, &user, view_id), + db.ack_relay_banner_view(community, banner.public_id, &user, view_id), + ); + let changed = [first.expect("first"), second.expect("second")] + .into_iter() + .filter(|outcome| { + matches!( + outcome, + RelayBannerViewOutcome::Accepted { changed: true, .. } + ) + }) + .count(); + assert_eq!(changed, 1); + assert_eq!( + db.ack_relay_banner_view(community, banner.public_id, &user, Uuid::new_v4()) + .await + .expect("second distinct view"), + RelayBannerViewOutcome::Accepted { + display_count: 2, + changed: true + } + ); + } + #[tokio::test] #[ignore = "requires Postgres"] async fn dismiss_permanently_suppresses_banner() { @@ -854,7 +1263,7 @@ mod tests { RelayBannerDismissOutcome::Dismissed { changed: false } ); assert_eq!( - db.ack_relay_banner_view(community, banner.public_id, &user) + db.ack_relay_banner_view(community, banner.public_id, &user, Uuid::new_v4()) .await .expect("view after dismiss"), RelayBannerViewOutcome::Dismissed diff --git a/crates/buzz-relay/src/api/admin/mod.rs b/crates/buzz-relay/src/api/admin/mod.rs index 25ddf26155f..408b69a8dfa 100644 --- a/crates/buzz-relay/src/api/admin/mod.rs +++ b/crates/buzz-relay/src/api/admin/mod.rs @@ -7,6 +7,7 @@ mod auth; mod error; +use std::collections::HashSet; use std::sync::Arc; use auth::{ @@ -568,13 +569,16 @@ async fn upsert_banner( } buzz_db::RelayBannerScope::AllCommunities } else { - let ids = body.community_ids.unwrap_or_default(); + let mut ids = body.community_ids.unwrap_or_default(); + ids.sort_unstable(); + ids.dedup(); if ids.is_empty() { return Err(ApiError::bad_request( "invalid_scope", "communityIds must be non-empty when targetScope is communities", )); } + validate_active_banner_communities(&state, &ids).await?; buzz_db::RelayBannerScope::Communities( ids.into_iter() .map(buzz_core::CommunityId::from_uuid) @@ -593,6 +597,7 @@ async fn upsert_banner( }) .await .map_err(map_banner_db_error)?; + fan_out_relay_banner_update(&state, &banner, false).await; Ok(Json(AdminBannerResponse::from(banner))) } @@ -613,11 +618,112 @@ async fn disable_banner( .await?, )?; require_operator(&principal)?; - let disabled = state + let disabled_banner = state .db .admin_disable_active_relay_banner(&principal.pubkey) .await?; - Ok(Json(serde_json::json!({ "disabled": disabled }))) + if let Some(banner) = disabled_banner.as_ref() { + fan_out_relay_banner_update(&state, banner, true).await; + } + Ok(Json( + serde_json::json!({ "disabled": disabled_banner.is_some() }), + )) +} + +async fn validate_active_banner_communities( + state: &crate::state::AppState, + ids: &[Uuid], +) -> Result<(), ApiError> { + let active = state.db.admin_list_banner_communities().await?; + let active: HashSet = active + .into_iter() + .map(|community| *community.id.as_uuid()) + .collect(); + if ids.iter().all(|id| active.contains(id)) { + Ok(()) + } else { + Err(ApiError::bad_request( + "invalid_scope", + "communityIds must refer to active communities", + )) + } +} + +async fn fan_out_relay_banner_update( + state: &crate::state::AppState, + banner: &buzz_db::RelayBannerRecord, + disabled: bool, +) { + let created_at = nostr::Timestamp::from(banner.updated_at.timestamp().max(0) as u64); + let event = if disabled { + match crate::api::banners::banner_disabled_event( + &state.relay_keypair, + banner, + Some(created_at), + ) { + Ok(event) => event, + Err(error) => { + tracing::warn!("Relay banner disable event signing failed: {error}"); + return; + } + } + } else { + match crate::api::banners::banner_event(&state.relay_keypair, banner, Some(created_at)) { + Ok(event) => event, + Err(error) => { + tracing::warn!("Relay banner event signing failed: {error}"); + return; + } + } + }; + publish_relay_banner_update(state, banner, &event).await; + crate::handlers::event::fan_out_relay_banner_event(state, banner, &event, !disabled).await; +} + +async fn publish_relay_banner_update( + state: &crate::state::AppState, + banner: &buzz_db::RelayBannerRecord, + event: &nostr::Event, +) { + let communities = if banner.target_all_communities { + match state.db.admin_list_banner_communities().await { + Ok(communities) => communities + .into_iter() + .map(|community| community.id) + .collect(), + Err(error) => { + tracing::warn!("Relay banner community list failed during fan-out: {error}"); + return; + } + } + } else { + banner.community_ids.clone() + }; + for community_id in communities { + let Some(host) = state + .db + .lookup_community_host(community_id) + .await + .unwrap_or_else(|error| { + tracing::warn!(%community_id, "Relay banner community host lookup failed: {error}"); + None + }) + else { + continue; + }; + let tenant = buzz_core::TenantContext::resolved(community_id, host); + state.mark_local_event(community_id, &event.id); + if let Err(error) = state + .pubsub + .publish_event(&tenant, buzz_pubsub::EventTopic::Global, event) + .await + { + state + .local_event_ids + .invalidate(&(community_id, event.id.to_bytes())); + tracing::warn!(%community_id, "Relay banner Redis publish failed: {error}"); + } + } } fn map_banner_db_error(error: buzz_db::DbError) -> ApiError { diff --git a/crates/buzz-relay/src/api/banners.rs b/crates/buzz-relay/src/api/banners.rs index 8a680274201..56d8d2d1a74 100644 --- a/crates/buzz-relay/src/api/banners.rs +++ b/crates/buzz-relay/src/api/banners.rs @@ -11,6 +11,7 @@ use axum::{ use buzz_core::TenantContext; use serde::Serialize; use serde_json::Value; +use uuid::Uuid; use crate::api::{api_error, bridge, internal_error, relay_members}; use crate::state::AppState; @@ -19,6 +20,8 @@ use crate::state::AppState; pub(crate) const BANNER_ACTIVE_PATH: &str = "/api/banners/current"; /// Client route for acknowledging a successful banner render/view. pub(crate) const BANNER_VIEW_ROUTE: &str = "/api/banners/{banner_id}/view"; +/// Header carrying a client-stable UUID for idempotent view retries. +pub(crate) const BANNER_VIEW_ID_HEADER: &str = "x-buzz-banner-view-id"; /// Client route for dismissing a banner. pub(crate) const BANNER_DISMISS_ROUTE: &str = "/api/banners/{banner_id}/dismiss"; @@ -43,6 +46,8 @@ pub(crate) struct BannerResponse { #[serde(rename_all = "camelCase")] pub(crate) struct BannerAckResponse { status: &'static str, + #[serde(skip_serializing_if = "Option::is_none")] + changed: Option, } impl From for BannerResponse { @@ -61,6 +66,7 @@ impl From for BannerResponse { pub(crate) fn banner_event( relay_keypair: &nostr::Keys, banner: &buzz_db::RelayBannerRecord, + created_at: Option, ) -> Result { let id = banner.public_id.to_string(); let scope = if banner.target_all_communities { @@ -80,13 +86,50 @@ pub(crate) fn banner_event( nostr::Tag::parse(["d", id.as_str()]).map_err(|e| e.to_string())?, nostr::Tag::parse(["scope", scope]).map_err(|e| e.to_string())?, ]; - nostr::EventBuilder::new( + let builder = nostr::EventBuilder::new( nostr::Kind::Custom(buzz_core::kind::KIND_RELAY_BANNER as u16), content, ) - .tags(tags) - .sign_with_keys(relay_keypair) - .map_err(|e| e.to_string()) + .tags(tags); + let builder = if let Some(created_at) = created_at { + builder.custom_created_at(created_at) + } else { + builder + }; + builder + .sign_with_keys(relay_keypair) + .map_err(|e| e.to_string()) +} + +pub(crate) fn banner_disabled_event( + relay_keypair: &nostr::Keys, + banner: &buzz_db::RelayBannerRecord, + created_at: Option, +) -> Result { + let id = banner.public_id.to_string(); + let scope = if banner.target_all_communities { + "all" + } else { + "communities" + }; + let tags = vec![ + nostr::Tag::parse(["d", id.as_str()]).map_err(|e| e.to_string())?, + nostr::Tag::parse(["scope", scope]).map_err(|e| e.to_string())?, + nostr::Tag::parse(["status", "disabled"]).map_err(|e| e.to_string())?, + ]; + let builder = nostr::EventBuilder::new( + nostr::Kind::Custom(buzz_core::kind::KIND_RELAY_BANNER as u16), + "", + ) + .tags(tags); + let builder = if let Some(created_at) = created_at { + builder.custom_created_at(created_at) + } else { + builder + }; + builder + .sign_with_keys(relay_keypair) + .map_err(|e| e.to_string()) } pub(crate) async fn active_banner_event_for_user( @@ -102,7 +145,7 @@ pub(crate) async fn active_banner_event_for_user( else { return Ok(None); }; - banner_event(&state.relay_keypair, &banner).map(Some) + banner_event(&state.relay_keypair, &banner, None).map(Some) } pub(crate) async fn get_active_banner( @@ -144,25 +187,33 @@ pub(crate) async fn ack_banner_view( if !body.is_empty() { return Err(non_empty_ack_body_error()); } + let view_id = parse_view_id(&headers)?; match state .db - .ack_relay_banner_view(tenant.community(), banner_id, pubkey.as_bytes()) + .ack_relay_banner_view(tenant.community(), banner_id, pubkey.as_bytes(), view_id) .await .map_err(|e| internal_error(&format!("banner view ack: {e}")))? { - buzz_db::RelayBannerViewOutcome::Accepted { .. } => { - metrics::counter!( - "buzz_relay_banner_views_total", - "community" => tenant.host().to_owned() - ) - .increment(1); - Ok(Json(BannerAckResponse { status: "accepted" })) + buzz_db::RelayBannerViewOutcome::Accepted { changed, .. } => { + if changed { + metrics::counter!( + "buzz_relay_banner_views_total", + "community" => tenant.host().to_owned() + ) + .increment(1); + } + Ok(Json(BannerAckResponse { + status: "accepted", + changed: Some(changed), + })) } buzz_db::RelayBannerViewOutcome::Dismissed => Ok(Json(BannerAckResponse { status: "dismissed", + changed: None, })), buzz_db::RelayBannerViewOutcome::Exhausted { .. } => Ok(Json(BannerAckResponse { status: "exhausted", + changed: None, })), buzz_db::RelayBannerViewOutcome::NotEligible => Err(api_error( axum::http::StatusCode::NOT_FOUND, @@ -207,6 +258,7 @@ pub(crate) async fn ack_banner_dismiss( } Ok(Json(BannerAckResponse { status: "dismissed", + changed: Some(changed), })) } buzz_db::RelayBannerDismissOutcome::NotEligible => Err(api_error( @@ -257,6 +309,29 @@ async fn authenticate_client_request( )) } +fn parse_view_id(headers: &HeaderMap) -> Result)> { + let Some(raw) = headers + .get(BANNER_VIEW_ID_HEADER) + .and_then(|value| value.to_str().ok()) + else { + return Err(invalid_view_id_error()); + }; + let Ok(view_id) = Uuid::parse_str(raw) else { + return Err(invalid_view_id_error()); + }; + if raw != view_id.hyphenated().to_string() || view_id.is_nil() { + return Err(invalid_view_id_error()); + } + Ok(view_id) +} + +fn invalid_view_id_error() -> (axum::http::StatusCode, Json) { + api_error( + axum::http::StatusCode::BAD_REQUEST, + "missing or invalid banner view id", + ) +} + fn non_empty_ack_body_error() -> (axum::http::StatusCode, Json) { api_error( axum::http::StatusCode::BAD_REQUEST, @@ -288,7 +363,7 @@ mod tests { fn banner_event_uses_published_contract() { let keys = nostr::Keys::generate(); let banner = banner_record(false); - let event = banner_event(&keys, &banner).expect("banner event"); + let event = banner_event(&keys, &banner, None).expect("banner event"); assert_eq!( event.kind.as_u16() as u32, @@ -322,11 +397,44 @@ mod tests { assert!(content.get("message").is_none()); } + #[test] + fn disabled_banner_event_uses_clear_contract() { + let keys = nostr::Keys::generate(); + let banner = banner_record(false); + let event = banner_disabled_event(&keys, &banner, None).expect("disabled event"); + + assert_eq!( + event.kind.as_u16() as u32, + buzz_core::kind::KIND_RELAY_BANNER + ); + assert!(event.content.is_empty()); + assert_eq!( + event + .tags + .iter() + .find(|tag| tag.kind().to_string() == "scope") + .and_then(|tag| tag.content()), + Some("communities") + ); + assert_eq!( + event + .tags + .iter() + .find(|tag| tag.kind().to_string() == "status") + .and_then(|tag| tag.content()), + Some("disabled") + ); + assert!(event + .tags + .iter() + .all(|tag| tag.kind().to_string() != "text")); + } + #[test] fn all_community_banner_event_has_all_scope() { let keys = nostr::Keys::generate(); let banner = banner_record(true); - let event = banner_event(&keys, &banner).expect("banner event"); + let event = banner_event(&keys, &banner, None).expect("banner event"); assert_eq!( event @@ -344,6 +452,40 @@ mod tests { assert_eq!(BANNER_DISMISS_ROUTE, "/api/banners/{banner_id}/dismiss"); } + #[test] + fn parse_view_id_requires_canonical_non_nil_uuid() { + let mut headers = HeaderMap::new(); + assert!(parse_view_id(&headers).is_err()); + + headers.insert(BANNER_VIEW_ID_HEADER, "not-a-uuid".parse().expect("header")); + assert!(parse_view_id(&headers).is_err()); + + headers.insert( + BANNER_VIEW_ID_HEADER, + "00000000-0000-0000-0000-000000000000" + .parse() + .expect("header"), + ); + assert!(parse_view_id(&headers).is_err()); + + headers.insert( + BANNER_VIEW_ID_HEADER, + "12345678123456781234567812345678".parse().expect("header"), + ); + assert!(parse_view_id(&headers).is_err()); + + headers.insert( + BANNER_VIEW_ID_HEADER, + "12345678-1234-5678-1234-567812345678" + .parse() + .expect("header"), + ); + assert_eq!( + parse_view_id(&headers).expect("valid canonical uuid"), + Uuid::parse_str("12345678-1234-5678-1234-567812345678").expect("uuid") + ); + } + #[test] fn non_empty_ack_body_is_rejected() { let (status, body) = non_empty_ack_body_error(); diff --git a/crates/buzz-relay/src/handlers/event.rs b/crates/buzz-relay/src/handlers/event.rs index 3f767c18741..01d6451fe90 100644 --- a/crates/buzz-relay/src/handlers/event.rs +++ b/crates/buzz-relay/src/handlers/event.rs @@ -36,7 +36,7 @@ pub(crate) fn bounded_kind_label(kind: u32) -> String { match kind { 0..=9 | 1059 | 1063 => kind.to_string(), 8000..=8003 | 9000..=9022 | 9030..=9036 => kind.to_string(), - 13534..=13535 => kind.to_string(), + 13534..=13536 => kind.to_string(), 20000..=29999 => kind.to_string(), 30023 | 30315 | 39000..=39003 => kind.to_string(), 40002..=40100 => kind.to_string(), @@ -221,6 +221,88 @@ pub async fn filter_fanout_by_access( allowed } +/// Fan out a relay banner control event to live banner subscribers whose authenticated +/// user is eligible for that banner in their bound community. +pub(crate) async fn fan_out_relay_banner_event( + state: &AppState, + banner: &buzz_db::RelayBannerRecord, + event: &Event, + require_active: bool, +) { + let communities: Vec = if banner.target_all_communities { + state + .conn_manager + .per_community_ws_connections() + .keys() + .copied() + .collect() + } else { + banner.community_ids.clone() + }; + fan_out_relay_banner_event_to_communities(state, banner, event, require_active, communities) + .await; +} + +async fn fan_out_relay_banner_event_to_communities( + state: &AppState, + banner: &buzz_db::RelayBannerRecord, + event: &Event, + require_active: bool, + communities: Vec, +) { + let stored = StoredEvent::with_received_at(event.clone(), chrono::Utc::now(), None, true); + let mut matches = Vec::new(); + let mut seen = std::collections::HashSet::new(); + for community_id in communities { + for (conn_id, sub_id) in state.sub_registry.fan_out_scoped(community_id, &stored) { + if !seen.insert((conn_id, sub_id.clone())) { + continue; + } + let Some(pubkey) = state.conn_manager.pubkey_for_conn(conn_id) else { + continue; + }; + match state + .db + .relay_banner_user_eligible(banner.id, community_id, &pubkey, require_active) + .await + { + Ok(true) => matches.push((conn_id, sub_id)), + Ok(false) => {} + Err(error) => { + tracing::warn!(%community_id, "Relay banner eligibility lookup failed: {error}"); + } + } + } + } + if matches.is_empty() { + return; + } + let event_json = match serde_json::to_string(event) { + Ok(json) => json, + Err(error) => { + error!("Failed to serialize relay banner event for fan-out: {error}"); + return; + } + }; + let frames = fanout_frame_cache( + matches.iter().map(|(_, sub_id)| sub_id.as_str()), + &event_json, + ); + let drop_count = send_fanout_frames( + state, + matches + .iter() + .map(|(conn_id, sub_id)| (*conn_id, sub_id.as_str())), + &frames, + ); + if drop_count > 0 { + tracing::warn!( + drop_count, + "relay banner fan-out: {drop_count} connection(s) cancelled due to full/closed buffers" + ); + } +} + /// Deliver one event to this relay's local subscribers through the access gate. /// /// This is the single guarded send path for relay-local EVENT delivery. It runs @@ -303,6 +385,11 @@ pub async fn fan_out_pubsub_event(state: &Arc, channel_event: buzz_pub return; } + if event_kind_u32(&stored.event) == buzz_core::kind::KIND_RELAY_BANNER { + fan_out_pubsub_relay_banner_event(state, community_id, &stored.event).await; + return; + } + let matches = state.sub_registry.fan_out_scoped(community_id, &stored); let matches = filter_fanout_by_access(state, community_id, &stored, matches, None).await; metrics::counter!("buzz_multinode_fanout_total").increment(1); @@ -337,6 +424,38 @@ pub async fn fan_out_pubsub_event(state: &Arc, channel_event: buzz_pub } } +async fn fan_out_pubsub_relay_banner_event( + state: &AppState, + community_id: CommunityId, + event: &Event, +) { + let disabled = event + .tags + .iter() + .any(|tag| tag.kind().to_string() == "status" && tag.content() == Some("disabled")); + let banner_id = event + .tags + .iter() + .find(|tag| tag.kind().to_string() == "d") + .and_then(|tag| tag.content()) + .and_then(|value| uuid::Uuid::parse_str(value).ok()); + let Some(banner_id) = banner_id else { + tracing::warn!("Relay banner pubsub event missing valid d tag"); + return; + }; + let Some(banner) = (match state.db.relay_banner_by_public_id(banner_id).await { + Ok(banner) => banner, + Err(error) => { + tracing::warn!(%banner_id, "Relay banner pubsub lookup failed: {error}"); + return; + } + }) else { + return; + }; + fan_out_relay_banner_event_to_communities(state, &banner, event, !disabled, vec![community_id]) + .await; +} + /// Schedule post-commit delivery/side effects for a stored event. /// /// This intentionally returns after only the bounded audit enqueue has completed: @@ -2859,3 +2978,161 @@ mod tests { } } } + +#[cfg(test)] +mod relay_banner_fanout_tests { + use super::*; + use axum::extract::ws::Message as WsMessage; + use std::collections::HashMap; + use tokio::sync::{mpsc, Mutex}; + use tokio_util::sync::CancellationToken; + use uuid::Uuid; + + async fn setup_state() -> (Arc, sqlx::PgPool) { + let pool = sqlx::PgPool::connect(&crate::test_support::database_url()) + .await + .expect("connect test DB"); + let state = crate::state::tests::test_state_with_database_pool(pool.clone()).await; + (state, pool) + } + + async fn insert_test_community(pool: &sqlx::PgPool, host_prefix: &str) -> CommunityId { + let id = Uuid::new_v4(); + sqlx::query("INSERT INTO communities (id, host) VALUES ($1, $2)") + .bind(id) + .bind(format!("{host_prefix}-{}.example", id.simple())) + .execute(pool) + .await + .expect("insert community"); + CommunityId::from_uuid(id) + } + + fn register_banner_sub( + state: &AppState, + community: CommunityId, + pubkey: Vec, + ) -> (Uuid, mpsc::Receiver) { + let conn_id = Uuid::new_v4(); + let (tx, rx) = mpsc::channel(4); + let (ctrl_tx, _ctrl_rx) = mpsc::channel(4); + state.conn_manager.register( + conn_id, + tx, + ctrl_tx, + None, + CancellationToken::new(), + community, + Arc::new(std::sync::atomic::AtomicU8::new(0)), + Arc::new(Mutex::new(HashMap::new())), + 3, + ); + state.conn_manager.set_authenticated_pubkey(conn_id, pubkey); + state.sub_registry.register_scoped( + community, + conn_id, + "banner".to_owned(), + vec![nostr::Filter::new().kind(nostr::Kind::Custom( + buzz_core::kind::KIND_RELAY_BANNER as u16, + ))], + None, + ); + (conn_id, rx) + } + + #[tokio::test] + #[ignore = "requires Postgres"] + async fn relay_banner_upsert_live_routes_only_to_eligible_users() { + let (state, pool) = setup_state().await; + let included = insert_test_community(&pool, "banner-fanout-included").await; + let excluded = insert_test_community(&pool, "banner-fanout-excluded").await; + let (_eligible_conn, mut eligible_rx) = register_banner_sub(&state, included, vec![7; 32]); + let (_excluded_conn, mut excluded_rx) = register_banner_sub(&state, excluded, vec![8; 32]); + let banner = state + .db + .admin_upsert_relay_banner(buzz_db::RelayBannerUpsert { + severity: buzz_db::RelayBannerSeverity::Info, + message: "live".to_owned(), + max_displays: 1, + scope: buzz_db::RelayBannerScope::Communities(vec![included]), + actor_pubkey: vec![1; 32], + }) + .await + .expect("upsert banner"); + let event = crate::api::banners::banner_event(&state.relay_keypair, &banner, None) + .expect("banner event"); + + fan_out_relay_banner_event(&state, &banner, &event, true).await; + + let frame = eligible_rx.recv().await.expect("live banner frame"); + let WsMessage::Text(frame) = frame else { + panic!("expected text frame"); + }; + assert!(frame.contains("\"live\"")); + assert!(excluded_rx.try_recv().is_err()); + } + + #[tokio::test] + #[ignore = "requires Postgres"] + async fn relay_banner_disable_live_routes_using_pre_disable_eligibility() { + let (state, pool) = setup_state().await; + let eligible_community = insert_test_community(&pool, "banner-disable-eligible").await; + let excluded_community = insert_test_community(&pool, "banner-disable-excluded").await; + let exhausted_user = vec![8; 32]; + let eligible_user = vec![9; 32]; + let excluded_user = vec![10; 32]; + let (_exhausted_conn, mut exhausted_rx) = + register_banner_sub(&state, eligible_community, exhausted_user.clone()); + let (_eligible_conn, mut eligible_rx) = + register_banner_sub(&state, eligible_community, eligible_user); + let (_excluded_conn, mut excluded_rx) = + register_banner_sub(&state, excluded_community, excluded_user); + let banner = state + .db + .admin_upsert_relay_banner(buzz_db::RelayBannerUpsert { + severity: buzz_db::RelayBannerSeverity::Warning, + message: "clear".to_owned(), + max_displays: 1, + scope: buzz_db::RelayBannerScope::Communities(vec![eligible_community]), + actor_pubkey: vec![1; 32], + }) + .await + .expect("upsert banner"); + assert_eq!( + state + .db + .ack_relay_banner_view( + eligible_community, + banner.public_id, + &exhausted_user, + Uuid::new_v4(), + ) + .await + .expect("exhaust user"), + buzz_db::RelayBannerViewOutcome::Accepted { + display_count: 1, + changed: true, + } + ); + let disabled = state + .db + .admin_disable_active_relay_banner(&[1; 32]) + .await + .expect("disable banner") + .expect("disabled banner"); + assert_eq!(disabled.id, banner.id); + let event = + crate::api::banners::banner_disabled_event(&state.relay_keypair, &disabled, None) + .expect("disabled event"); + + fan_out_relay_banner_event(&state, &disabled, &event, false).await; + + let frame = eligible_rx.recv().await.expect("live disable frame"); + let WsMessage::Text(frame) = frame else { + panic!("expected text frame"); + }; + assert!(frame.contains("\"status\",\"disabled\"")); + assert!(frame.contains("\"scope\",\"communities\"")); + assert!(exhausted_rx.try_recv().is_err()); + assert!(excluded_rx.try_recv().is_err()); + } +} diff --git a/crates/buzz-relay/src/handlers/req.rs b/crates/buzz-relay/src/handlers/req.rs index ded8348a583..ea3958b4342 100644 --- a/crates/buzz-relay/src/handlers/req.rs +++ b/crates/buzz-relay/src/handlers/req.rs @@ -222,6 +222,15 @@ pub async fn handle_req( } if filters_are_relay_banner_only(&filters) { + register_subscription( + &state, + &conn, + conn_id, + &sub_id, + &filters, + authorized_requested_channels.as_ref(), + ) + .await; match crate::api::banners::active_banner_event_for_user(&state, &conn.tenant, &pubkey_bytes) .await { @@ -306,46 +315,15 @@ pub async fn handle_req( return; } - { - let mut subs = conn.subscriptions.lock().await; - subs.insert(sub_id.clone(), filters.clone()); - } - - let replaced = if let Some(channel_ids) = authorized_requested_channels.as_ref() { - state.sub_registry.register_channels_scoped( - conn.tenant.community(), - conn_id, - sub_id.clone(), - filters.clone(), - channel_ids.clone(), - ) - } else { - state.sub_registry.register_scoped( - conn.tenant.community(), - conn_id, - sub_id.clone(), - filters.clone(), - None, - ) - }; - if let Some(replaced) = replaced { - release_subscription_topics(&state, &conn.tenant, &replaced.scope).await; - } - if let Some(channel_ids) = authorized_requested_channels.as_ref() { - for &channel_id in channel_ids { - state - .pubsub - .retain_topic(&conn.tenant, EventTopic::Channel(channel_id)) - .await; - } - } else { - state - .pubsub - .retain_topic(&conn.tenant, EventTopic::Global) - .await; - } - - debug!(conn_id = %conn_id, sub_id = %sub_id, "Subscription registered"); + register_subscription( + &state, + &conn, + conn_id, + &sub_id, + &filters, + authorized_requested_channels.as_ref(), + ) + .await; // NIP-01 OR semantics: execute one DB query per filter and deduplicate results // by event ID. Collapsing all filters into a single query would merge their @@ -1171,6 +1149,56 @@ pub(crate) fn extract_channel_ids_from_filters(filters: &[Filter]) -> Option>, +) { + { + let mut subs = conn.subscriptions.lock().await; + subs.insert(sub_id.to_owned(), filters.to_vec()); + } + + let replaced = if let Some(channel_ids) = authorized_requested_channels { + state.sub_registry.register_channels_scoped( + conn.tenant.community(), + conn_id, + sub_id.to_owned(), + filters.to_vec(), + channel_ids.clone(), + ) + } else { + state.sub_registry.register_scoped( + conn.tenant.community(), + conn_id, + sub_id.to_owned(), + filters.to_vec(), + None, + ) + }; + if let Some(replaced) = replaced { + release_subscription_topics(state, &conn.tenant, &replaced.scope).await; + } + if let Some(channel_ids) = authorized_requested_channels { + for &channel_id in channel_ids { + state + .pubsub + .retain_topic(&conn.tenant, EventTopic::Channel(channel_id)) + .await; + } + } else { + state + .pubsub + .retain_topic(&conn.tenant, EventTopic::Global) + .await; + } + + debug!(conn_id = %conn_id, sub_id = %sub_id, "Subscription registered"); +} + fn filters_are_huddle_liveness_only(filters: &[Filter]) -> bool { !filters.is_empty() && filters.iter().all(|filter| { diff --git a/migrations/0046_relay_banners.sql b/migrations/0046_relay_banners.sql index 2be963e9ab8..d93f4f6ca0a 100644 --- a/migrations/0046_relay_banners.sql +++ b/migrations/0046_relay_banners.sql @@ -53,10 +53,24 @@ CREATE TABLE relay_banner_user_state ( CREATE INDEX idx_relay_banner_user_state_pubkey ON relay_banner_user_state (community_id, pubkey, banner_id); +CREATE TABLE relay_banner_view_acks ( + banner_id BIGINT NOT NULL REFERENCES relay_banners(id) ON DELETE CASCADE, + community_id UUID NOT NULL REFERENCES communities(id) ON DELETE CASCADE, + pubkey BYTEA NOT NULL CHECK (length(pubkey) = 32), + view_id UUID NOT NULL, + viewed_at TIMESTAMPTZ NOT NULL DEFAULT now(), + PRIMARY KEY (banner_id, community_id, pubkey, view_id) +); + +CREATE INDEX idx_relay_banner_view_acks_pubkey + ON relay_banner_view_acks (community_id, pubkey, banner_id); + INSERT INTO _operator_global_tables (table_name, reason) VALUES ('relay_banners', 'deployment-global operator-configured banner; no community_id intentionally'), ('relay_banner_communities', 'deployment-global banner targeting allowlist; community_id is target provenance only'), - ('relay_banner_user_state', 'deployment-global per-user banner display state; community_id scopes the targeted impression only'); + ('relay_banner_user_state', 'deployment-global per-user banner display state; community_id scopes the targeted impression only'), + ('relay_banner_view_acks', 'deployment-global per-render banner idempotency keys; community_id scopes the targeted impression only'); SELECT attach_community_write_fence('relay_banner_communities'); SELECT attach_community_write_fence('relay_banner_user_state'); +SELECT attach_community_write_fence('relay_banner_view_acks'); diff --git a/schema/schema.sql b/schema/schema.sql index 6e4e3750889..8aa3775b3bb 100644 --- a/schema/schema.sql +++ b/schema/schema.sql @@ -1757,10 +1757,23 @@ CREATE TABLE relay_banner_user_state ( CREATE INDEX idx_relay_banner_user_state_pubkey ON relay_banner_user_state (community_id, pubkey, banner_id); +CREATE TABLE relay_banner_view_acks ( + banner_id BIGINT NOT NULL REFERENCES relay_banners(id) ON DELETE CASCADE, + community_id UUID NOT NULL REFERENCES communities(id) ON DELETE CASCADE, + pubkey BYTEA NOT NULL CHECK (length(pubkey) = 32), + view_id UUID NOT NULL, + viewed_at TIMESTAMPTZ NOT NULL DEFAULT now(), + PRIMARY KEY (banner_id, community_id, pubkey, view_id) +); + +CREATE INDEX idx_relay_banner_view_acks_pubkey + ON relay_banner_view_acks (community_id, pubkey, banner_id); + INSERT INTO _operator_global_tables (table_name, reason) VALUES ('relay_banners', 'deployment-global operator-configured banner; no community_id intentionally'), ('relay_banner_communities', 'deployment-global banner targeting allowlist; community_id is target provenance only'), - ('relay_banner_user_state', 'deployment-global per-user banner display state; community_id scopes the targeted impression only'); + ('relay_banner_user_state', 'deployment-global per-user banner display state; community_id scopes the targeted impression only'), + ('relay_banner_view_acks', 'deployment-global per-render banner idempotency keys; community_id scopes the targeted impression only'); -- Attach the universal fence to every existing table carrying community_id, @@ -1814,6 +1827,7 @@ SELECT attach_community_write_fence('relay_invites'); SELECT attach_community_write_fence('relay_members'); SELECT attach_community_write_fence('relay_banner_communities'); SELECT attach_community_write_fence('relay_banner_user_state'); +SELECT attach_community_write_fence('relay_banner_view_acks'); SELECT attach_community_write_fence('scheduled_workflow_fires'); SELECT attach_community_write_fence('subscriptions'); SELECT attach_community_write_fence('thread_metadata'); From 69f26c9a86aff1b5d66b0254c245bbf93c4682b9 Mon Sep 17 00:00:00 2001 From: coder 0 Date: Thu, 17 Sep 2026 11:49:22 -0400 Subject: [PATCH 3/4] Fix relay banner replacement fanout Signed-off-by: coder 0 --- crates/buzz-db/src/lib.rs | 3 +- crates/buzz-db/src/store/relay_banners.rs | 308 ++++++++++++----- crates/buzz-relay/src/api/admin/mod.rs | 385 +++++++++++++++++++++- crates/buzz-relay/src/handlers/event.rs | 115 ++++++- crates/buzz-relay/src/test_support.rs | 4 + 5 files changed, 712 insertions(+), 103 deletions(-) diff --git a/crates/buzz-db/src/lib.rs b/crates/buzz-db/src/lib.rs index 4a11131e2e3..c360616e342 100644 --- a/crates/buzz-db/src/lib.rs +++ b/crates/buzz-db/src/lib.rs @@ -80,7 +80,8 @@ pub use event::{EventQuery, DEFAULT_MAX_PAGE_LIMIT}; pub use reaction::ReactionEventInsertOutcome; pub use relay_banners::{ RelayBannerDismissOutcome, RelayBannerRecord, RelayBannerScope, RelayBannerSeverity, - RelayBannerTargetScope, RelayBannerUpsert, RelayBannerViewOutcome, MAX_BANNER_MESSAGE_CHARS, + RelayBannerTargetScope, RelayBannerUpsert, RelayBannerUpsertOutcome, RelayBannerViewOutcome, + MAX_BANNER_MESSAGE_CHARS, }; pub use reminder::DueReminder; pub use usage::UsageMetricsLeader; diff --git a/crates/buzz-db/src/store/relay_banners.rs b/crates/buzz-db/src/store/relay_banners.rs index 2d4e71660a8..91c9755636f 100644 --- a/crates/buzz-db/src/store/relay_banners.rs +++ b/crates/buzz-db/src/store/relay_banners.rs @@ -95,6 +95,15 @@ pub struct RelayBannerUpsert { pub actor_pubkey: Vec, } +/// Result of replacing the deployment's active banner. +#[derive(Debug, Clone, PartialEq, Eq)] +pub struct RelayBannerUpsertOutcome { + /// Previously active banner after being disabled, retaining its original target audience. + pub previous: Option, + /// Newly active banner. + pub active: RelayBannerRecord, +} + /// Stored relay banner plus targeting metadata. #[derive(Debug, Clone, PartialEq, Eq)] pub struct RelayBannerRecord { @@ -209,6 +218,77 @@ async fn acquire_banner_lock(tx: &mut Transaction<'_, Postgres>) -> Result<()> { Ok(()) } +async fn next_banner_protocol_timestamp( + tx: &mut Transaction<'_, Postgres>, +) -> Result> { + let timestamp = sqlx::query_scalar::<_, DateTime>( + r#" + SELECT to_timestamp(GREATEST( + EXTRACT(EPOCH FROM date_trunc('second', now()))::bigint, + COALESCE((SELECT MAX(FLOOR(EXTRACT(EPOCH FROM updated_at))::bigint) + 1 FROM relay_banners), 0) + ))::timestamptz + "#, + ) + .fetch_one(&mut **tx) + .await?; + Ok(timestamp) +} + +async fn active_banner_in_tx( + tx: &mut Transaction<'_, Postgres>, +) -> Result> { + let row = sqlx::query( + r#" + SELECT + b.id, + b.public_id, + b.severity, + b.message, + b.max_displays, + b.target_all_communities, + b.created_by, + b.created_at, + b.updated_at, + b.disabled_at, + COALESCE(array_agg(bc.community_id ORDER BY bc.community_id) + FILTER (WHERE bc.community_id IS NOT NULL), ARRAY[]::uuid[]) AS community_ids + FROM relay_banners b + LEFT JOIN relay_banner_communities bc ON bc.banner_id = b.id + WHERE b.disabled_at IS NULL + GROUP BY b.id + ORDER BY b.id DESC + LIMIT 1 + "#, + ) + .fetch_optional(&mut **tx) + .await?; + row.map(row_to_banner).transpose() +} + +async fn disable_active_banner_in_tx( + tx: &mut Transaction<'_, Postgres>, + active: &mut RelayBannerRecord, + actor_pubkey: &[u8], +) -> Result<()> { + let protocol_timestamp = next_banner_protocol_timestamp(tx).await?; + let updated = sqlx::query( + r#" + UPDATE relay_banners + SET disabled_at = $1, disabled_by = $2, updated_at = $1, updated_by = $2 + WHERE id = $3 AND disabled_at IS NULL + RETURNING updated_at, disabled_at + "#, + ) + .bind(protocol_timestamp) + .bind(actor_pubkey) + .bind(active.id) + .fetch_one(&mut **tx) + .await?; + active.updated_at = updated.try_get("updated_at")?; + active.disabled_at = updated.try_get("disabled_at")?; + Ok(()) +} + async fn banner_internal_id_for_public_id( tx: &mut Transaction<'_, Postgres>, public_id: Uuid, @@ -288,7 +368,7 @@ async fn relay_banner_user_eligible_in_tx( ) ) AND us.dismissed_at IS NULL - AND COALESCE(us.display_count, 0) < b.max_displays + AND (NOT $4 OR COALESCE(us.display_count, 0) < b.max_displays) ) "#, ) @@ -455,7 +535,7 @@ impl Db { pub async fn admin_upsert_relay_banner( &self, input: RelayBannerUpsert, - ) -> Result { + ) -> Result { input.validate()?; let connection = crate::observability::acquire_writer( &self.pool, @@ -465,23 +545,18 @@ impl Db { let mut tx = sqlx::Transaction::begin(connection, None).await?; acquire_banner_lock(&mut tx).await?; - sqlx::query( - r#" - UPDATE relay_banners - SET disabled_at = now(), disabled_by = $1, updated_at = now(), updated_by = $1 - WHERE disabled_at IS NULL - "#, - ) - .bind(&input.actor_pubkey) - .execute(&mut *tx) - .await?; + let mut previous = active_banner_in_tx(&mut tx).await?; + if let Some(active) = previous.as_mut() { + disable_active_banner_in_tx(&mut tx, active, &input.actor_pubkey).await?; + } + let protocol_timestamp = next_banner_protocol_timestamp(&mut tx).await?; let target_all = matches!(input.scope, RelayBannerScope::AllCommunities); let row = sqlx::query( r#" INSERT INTO relay_banners - (severity, message, max_displays, target_all_communities, created_by, updated_by) - VALUES ($1, $2, $3, $4, $5, $5) + (severity, message, max_displays, target_all_communities, created_at, updated_at, created_by, updated_by) + VALUES ($1, $2, $3, $4, $5, $5, $6, $6) RETURNING id, public_id, severity, message, max_displays, target_all_communities, created_by, created_at, updated_at, disabled_at "#, @@ -490,6 +565,7 @@ impl Db { .bind(&input.message) .bind(input.max_displays) .bind(target_all) + .bind(protocol_timestamp) .bind(&input.actor_pubkey) .fetch_one(&mut *tx) .await?; @@ -512,18 +588,21 @@ impl Db { }; tx.commit().await?; - Ok(RelayBannerRecord { - id: banner_id, - public_id: row.try_get("public_id")?, - severity: RelayBannerSeverity::parse(row.try_get("severity")?)?, - message: row.try_get("message")?, - max_displays: row.try_get("max_displays")?, - target_all_communities: row.try_get("target_all_communities")?, - community_ids, - created_by: row.try_get("created_by")?, - created_at: row.try_get("created_at")?, - updated_at: row.try_get("updated_at")?, - disabled_at: row.try_get("disabled_at")?, + Ok(RelayBannerUpsertOutcome { + previous, + active: RelayBannerRecord { + id: banner_id, + public_id: row.try_get("public_id")?, + severity: RelayBannerSeverity::parse(row.try_get("severity")?)?, + message: row.try_get("message")?, + max_displays: row.try_get("max_displays")?, + target_all_communities: row.try_get("target_all_communities")?, + community_ids, + created_by: row.try_get("created_by")?, + created_at: row.try_get("created_at")?, + updated_at: row.try_get("updated_at")?, + disabled_at: row.try_get("disabled_at")?, + }, }) } @@ -540,50 +619,11 @@ impl Db { .await?; let mut tx = sqlx::Transaction::begin(connection, None).await?; acquire_banner_lock(&mut tx).await?; - let Some(active) = sqlx::query( - r#" - SELECT - b.id, - b.public_id, - b.severity, - b.message, - b.max_displays, - b.target_all_communities, - b.created_by, - b.created_at, - b.updated_at, - b.disabled_at, - COALESCE(array_agg(bc.community_id ORDER BY bc.community_id) - FILTER (WHERE bc.community_id IS NOT NULL), ARRAY[]::uuid[]) AS community_ids - FROM relay_banners b - LEFT JOIN relay_banner_communities bc ON bc.banner_id = b.id - WHERE b.disabled_at IS NULL - GROUP BY b.id - ORDER BY b.id DESC - LIMIT 1 - "#, - ) - .fetch_optional(&mut *tx) - .await? - else { + let Some(mut active) = active_banner_in_tx(&mut tx).await? else { tx.commit().await?; return Ok(None); }; - let mut active = row_to_banner(active)?; - let updated = sqlx::query( - r#" - UPDATE relay_banners - SET disabled_at = now(), disabled_by = $1, updated_at = now(), updated_by = $1 - WHERE id = $2 AND disabled_at IS NULL - RETURNING updated_at, disabled_at - "#, - ) - .bind(actor_pubkey) - .bind(active.id) - .fetch_one(&mut *tx) - .await?; - active.updated_at = updated.try_get("updated_at")?; - active.disabled_at = updated.try_get("disabled_at")?; + disable_active_banner_in_tx(&mut tx, &mut active, actor_pubkey).await?; tx.commit().await?; Ok(Some(active)) } @@ -898,11 +938,14 @@ mod tests { use super::*; use sqlx::PgPool; - async fn setup_db() -> Db { + static RELAY_BANNER_TEST_LOCK: tokio::sync::Mutex<()> = tokio::sync::Mutex::const_new(()); + + async fn setup_db() -> (tokio::sync::MutexGuard<'static, ()>, Db) { + let guard = RELAY_BANNER_TEST_LOCK.lock().await; let pool = PgPool::connect(&crate::test_support::database_url()) .await .expect("connect to test DB"); - Db::from_pool(pool) + (guard, Db::from_pool(pool)) } pub(crate) async fn insert_test_community(pool: &PgPool, host_prefix: &str) -> CommunityId { @@ -927,7 +970,7 @@ mod tests { #[tokio::test] #[ignore = "requires Postgres"] async fn upsert_replaces_single_active_banner_and_retains_history() { - let db = setup_db().await; + let (_guard, db) = setup_db().await; let first = db .admin_upsert_relay_banner(RelayBannerUpsert { severity: RelayBannerSeverity::Info, @@ -937,7 +980,8 @@ mod tests { actor_pubkey: actor(1), }) .await - .expect("insert first banner"); + .expect("insert first banner") + .active; let second = db .admin_upsert_relay_banner(RelayBannerUpsert { severity: RelayBannerSeverity::Warning, @@ -947,7 +991,8 @@ mod tests { actor_pubkey: actor(2), }) .await - .expect("replace banner"); + .expect("replace banner") + .active; let active = db .admin_get_active_relay_banner() @@ -955,17 +1000,76 @@ mod tests { .expect("get active") .expect("active banner"); assert_eq!(active.id, second.id); - let all = db.admin_list_relay_banners(10).await.expect("list banners"); - assert_eq!(all.len(), 2); - assert!(all + let all = db + .admin_list_relay_banners(200) + .await + .expect("list banners"); + let created: Vec<_> = all + .iter() + .filter(|banner| [first.public_id, second.public_id].contains(&banner.public_id)) + .collect(); + assert_eq!(created.len(), 2); + assert_eq!( + created + .iter() + .map(|banner| banner.public_id) + .collect::>(), + vec![second.public_id, first.public_id], + "history list should return only this test's created banners in newest-first order" + ); + assert!(created .iter() .any(|b| b.id == first.id && b.disabled_at.is_some())); } + #[tokio::test] + #[ignore = "requires Postgres"] + async fn rapid_replacement_and_disable_use_monotonic_protocol_seconds() { + let (_guard, db) = setup_db().await; + let first = db + .admin_upsert_relay_banner(RelayBannerUpsert { + severity: RelayBannerSeverity::Info, + message: "first monotonic".to_owned(), + max_displays: 1, + scope: RelayBannerScope::AllCommunities, + actor_pubkey: actor(18), + }) + .await + .expect("insert first banner") + .active; + let second = db + .admin_upsert_relay_banner(RelayBannerUpsert { + severity: RelayBannerSeverity::Warning, + message: "second monotonic".to_owned(), + max_displays: 1, + scope: RelayBannerScope::AllCommunities, + actor_pubkey: actor(18), + }) + .await + .expect("replace banner"); + let disabled_first = second.previous.as_ref().expect("disabled first"); + assert_eq!(disabled_first.id, first.id); + assert!( + disabled_first.updated_at.timestamp() < second.active.updated_at.timestamp(), + "disabled clear must sort before replacement active at protocol-second precision" + ); + + let disabled_second = db + .admin_disable_active_relay_banner(&actor(18)) + .await + .expect("disable active banner") + .expect("disabled second"); + assert_eq!(disabled_second.id, second.active.id); + assert!( + second.active.updated_at.timestamp() < disabled_second.updated_at.timestamp(), + "final disable must sort after the replacement active at protocol-second precision" + ); + } + #[tokio::test] #[ignore = "requires Postgres"] async fn community_scope_filters_delivery() { - let db = setup_db().await; + let (_guard, db) = setup_db().await; let included = make_community(&db.pool).await; let excluded = make_community(&db.pool).await; let user = actor(3); @@ -978,7 +1082,8 @@ mod tests { actor_pubkey: actor(4), }) .await - .expect("insert scoped banner"); + .expect("insert scoped banner") + .active; assert_eq!(banner.community_ids, vec![included]); assert!(db @@ -996,7 +1101,7 @@ mod tests { #[tokio::test] #[ignore = "requires Postgres"] async fn duplicate_community_scope_is_rejected() { - let db = setup_db().await; + let (_guard, db) = setup_db().await; let community = make_community(&db.pool).await; let result = db .admin_upsert_relay_banner(RelayBannerUpsert { @@ -1016,19 +1121,35 @@ mod tests { #[tokio::test] #[ignore = "requires Postgres"] async fn disabled_banner_remains_eligible_when_active_requirement_is_lifted() { - let db = setup_db().await; + let (_guard, db) = setup_db().await; let community = make_community(&db.pool).await; let user = actor(16); let banner = db .admin_upsert_relay_banner(RelayBannerUpsert { severity: RelayBannerSeverity::Warning, message: "disable eligibility".to_owned(), - max_displays: 2, + max_displays: 1, scope: RelayBannerScope::AllCommunities, actor_pubkey: actor(17), }) .await - .expect("insert banner"); + .expect("insert banner") + .active; + assert_eq!( + db.ack_relay_banner_view(community, banner.public_id, &user, Uuid::new_v4()) + .await + .expect("exhaust display"), + RelayBannerViewOutcome::Accepted { + display_count: 1, + changed: true, + } + ); + assert_eq!( + db.ack_relay_banner_view(community, banner.public_id, &user, Uuid::new_v4()) + .await + .expect("already exhausted"), + RelayBannerViewOutcome::Exhausted { display_count: 1 } + ); let disabled = db .admin_disable_active_relay_banner(&actor(17)) .await @@ -1050,7 +1171,7 @@ mod tests { #[tokio::test] #[ignore = "requires Postgres"] async fn view_ack_consumes_display_and_then_exhausts() { - let db = setup_db().await; + let (_guard, db) = setup_db().await; let community = make_community(&db.pool).await; let user = actor(5); let banner = db @@ -1062,7 +1183,8 @@ mod tests { actor_pubkey: actor(6), }) .await - .expect("insert banner"); + .expect("insert banner") + .active; assert_eq!( db.active_relay_banner_for_user(community, &user) @@ -1107,7 +1229,7 @@ mod tests { #[tokio::test] #[ignore = "requires Postgres"] async fn view_ack_same_key_retry_does_not_consume_again() { - let db = setup_db().await; + let (_guard, db) = setup_db().await; let community = make_community(&db.pool).await; let user = actor(9); let view_id = Uuid::new_v4(); @@ -1120,7 +1242,8 @@ mod tests { actor_pubkey: actor(10), }) .await - .expect("insert banner"); + .expect("insert banner") + .active; assert_eq!( db.ack_relay_banner_view(community, banner.public_id, &user, view_id) @@ -1150,7 +1273,7 @@ mod tests { #[tokio::test] #[ignore = "requires Postgres"] async fn distinct_view_keys_increment_until_cap() { - let db = setup_db().await; + let (_guard, db) = setup_db().await; let community = make_community(&db.pool).await; let user = actor(11); let banner = db @@ -1162,7 +1285,8 @@ mod tests { actor_pubkey: actor(12), }) .await - .expect("insert banner"); + .expect("insert banner") + .active; assert_eq!( db.ack_relay_banner_view(community, banner.public_id, &user, Uuid::new_v4()) @@ -1193,7 +1317,7 @@ mod tests { #[tokio::test] #[ignore = "requires Postgres"] async fn concurrent_same_view_key_consumes_once() { - let db = setup_db().await; + let (_guard, db) = setup_db().await; let community = make_community(&db.pool).await; let user = actor(13); let view_id = Uuid::new_v4(); @@ -1206,7 +1330,8 @@ mod tests { actor_pubkey: actor(14), }) .await - .expect("insert banner"); + .expect("insert banner") + .active; let (first, second) = tokio::join!( db.ack_relay_banner_view(community, banner.public_id, &user, view_id), @@ -1236,7 +1361,7 @@ mod tests { #[tokio::test] #[ignore = "requires Postgres"] async fn dismiss_permanently_suppresses_banner() { - let db = setup_db().await; + let (_guard, db) = setup_db().await; let community = make_community(&db.pool).await; let user = actor(7); let banner = db @@ -1248,7 +1373,8 @@ mod tests { actor_pubkey: actor(8), }) .await - .expect("insert banner"); + .expect("insert banner") + .active; assert_eq!( db.ack_relay_banner_dismiss(community, banner.public_id, &user) diff --git a/crates/buzz-relay/src/api/admin/mod.rs b/crates/buzz-relay/src/api/admin/mod.rs index 408b69a8dfa..1bb5bcc6faf 100644 --- a/crates/buzz-relay/src/api/admin/mod.rs +++ b/crates/buzz-relay/src/api/admin/mod.rs @@ -586,7 +586,7 @@ async fn upsert_banner( ) }; - let banner = state + let outcome = state .db .admin_upsert_relay_banner(buzz_db::RelayBannerUpsert { severity: body.severity, @@ -597,8 +597,11 @@ async fn upsert_banner( }) .await .map_err(map_banner_db_error)?; - fan_out_relay_banner_update(&state, &banner, false).await; - Ok(Json(AdminBannerResponse::from(banner))) + if let Some(previous) = outcome.previous.as_ref() { + fan_out_relay_banner_update(&state, previous, true).await; + } + fan_out_relay_banner_update(&state, &outcome.active, false).await; + Ok(Json(AdminBannerResponse::from(outcome.active))) } async fn disable_banner( @@ -1695,6 +1698,141 @@ mod postgres_tests { builder.body(Body::empty()).expect("request") } + fn banner_connection( + state: &Arc, + community: buzz_core::CommunityId, + host: &str, + pubkey: nostr::PublicKey, + ) -> ( + Arc, + tokio::sync::mpsc::Receiver, + ) { + let conn_id = Uuid::new_v4(); + let (send_tx, send_rx) = tokio::sync::mpsc::channel(16); + let (ctrl_tx, _ctrl_rx) = tokio::sync::mpsc::channel(4); + let cancel = tokio_util::sync::CancellationToken::new(); + let backpressure_count = Arc::new(std::sync::atomic::AtomicU8::new(0)); + let subscriptions = Arc::new(tokio::sync::Mutex::new(std::collections::HashMap::new())); + state.conn_manager.register( + conn_id, + send_tx.clone(), + ctrl_tx.clone(), + None, + cancel.clone(), + community, + backpressure_count.clone(), + subscriptions.clone(), + 3, + ); + state + .conn_manager + .set_authenticated_pubkey(conn_id, pubkey.to_bytes().to_vec()); + let conn = Arc::new(crate::connection::ConnectionState { + conn_id, + tenant: buzz_core::TenantContext::resolved(community, host.to_owned()), + remote_addr: "127.0.0.1:1234".parse().expect("socket addr"), + auth_state: std::sync::Mutex::new(crate::connection::AuthState::Authenticated( + buzz_auth::AuthContext { + pubkey, + scopes: vec![], + channel_ids: None, + auth_method: buzz_auth::AuthMethod::Nip42, + agent_owner_pubkey: None, + }, + )), + subscriptions, + send_tx, + ctrl_tx, + cancel, + backpressure_count, + grace_limit: 3, + }); + (conn, send_rx) + } + + async fn register_banner_req( + state: Arc, + community: buzz_core::CommunityId, + host: &str, + pubkey: nostr::PublicKey, + sub_id: &str, + ) -> tokio::sync::mpsc::Receiver { + let (conn, mut rx) = banner_connection(&state, community, host, pubkey); + crate::handlers::req::handle_req( + sub_id.to_owned(), + vec![nostr::Filter::new().kind(nostr::Kind::Custom( + buzz_core::kind::KIND_RELAY_BANNER as u16, + ))], + Vec::new(), + conn, + state, + ) + .await; + loop { + let axum::extract::ws::Message::Text(frame) = + rx.recv().await.expect("initial banner frame") + else { + panic!("expected text frame"); + }; + let frame: serde_json::Value = serde_json::from_str(&frame).expect("relay frame JSON"); + if frame[0] == "EOSE" { + break; + } + assert_eq!(frame[0], "EVENT"); + } + rx + } + + async fn recv_banner_frame( + rx: &mut tokio::sync::mpsc::Receiver, + ) -> nostr::Event { + let axum::extract::ws::Message::Text(text) = rx.recv().await.expect("banner EVENT") else { + panic!("expected text frame"); + }; + let frame: serde_json::Value = serde_json::from_str(&text).expect("EVENT frame JSON"); + assert_eq!(frame[0], "EVENT"); + serde_json::from_value(frame[2].clone()).expect("banner event") + } + + fn is_disabled_banner(event: &nostr::Event) -> bool { + event + .tags + .iter() + .any(|tag| tag.kind().to_string() == "status" && tag.content() == Some("disabled")) + } + + fn banner_text(event: &nostr::Event) -> String { + serde_json::from_str::(&event.content).expect("banner content JSON") + ["text"] + .as_str() + .expect("banner text") + .to_owned() + } + + fn spawn_banner_pubsub_loop(state: Arc) -> tokio::task::JoinHandle<()> { + let mut rx = state.pubsub.subscribe_local(); + tokio::spawn(async move { + while let Ok(channel_event) = rx.recv().await { + crate::handlers::event::fan_out_pubsub_event(&state, channel_event).await; + } + }) + } + + async fn banner_test_community( + pool: &sqlx::PgPool, + prefix: &str, + ) -> (buzz_core::CommunityId, String) { + let id = Uuid::new_v4(); + let host = format!("{prefix}-{}.example", id.simple()); + sqlx::query("INSERT INTO communities (id, host) VALUES ($1, $2)") + .bind(id) + .bind(&host) + .execute(pool) + .await + .expect("insert community"); + (buzz_core::CommunityId::from_uuid(id), host) + } + async fn status_for( state: Arc, request: Request, @@ -1702,6 +1840,196 @@ mod postgres_tests { router(state).oneshot(request).await.expect("response") } + #[tokio::test] + #[ignore = "requires Postgres and Redis"] + async fn banner_admin_req_local_and_redis_lifecycle_orders_scope_shrink() { + let _guard = crate::test_support::RELAY_BANNER_TEST_LOCK.lock().await; + let redis_url = + std::env::var("REDIS_URL").unwrap_or_else(|_| "redis://127.0.0.1:6379".to_string()); + let redis_pool = deadpool_redis::Config::from_url(&redis_url) + .create_pool(Some(deadpool_redis::Runtime::Tokio1)) + .expect("redis pool"); + let mut redis_conn = match redis_pool.get().await { + Ok(conn) => conn, + Err(_) => { + eprintln!("skipping banner lifecycle test: Redis unavailable"); + return; + } + }; + if redis::cmd("PING") + .query_async::(&mut redis_conn) + .await + .is_err() + { + eprintln!("skipping banner lifecycle test: Redis unavailable"); + return; + } + drop(redis_conn); + + let pool = sqlx::PgPool::connect(&database_url()) + .await + .expect("connect test DB"); + let operator = test_operator_keys(); + let origin = nip98_state_with_database_pool_and_redis( + pool.clone(), + &redis_url, + vec![operator.public_key().to_hex()], + ) + .await; + let receiver = nip98_state_with_database_pool_and_redis( + pool.clone(), + &redis_url, + vec![operator.public_key().to_hex()], + ) + .await; + let origin_subscriber = tokio::spawn(origin.pubsub.clone().run_subscriber()); + let receiver_subscriber = tokio::spawn(receiver.pubsub.clone().run_subscriber()); + let origin_fanout = spawn_banner_pubsub_loop(origin.clone()); + let receiver_fanout = spawn_banner_pubsub_loop(receiver.clone()); + + let (retained, retained_host) = banner_test_community(&pool, "banner-admin-retained").await; + let (removed, removed_host) = banner_test_community(&pool, "banner-admin-removed").await; + origin + .db + .admin_disable_active_relay_banner(&operator.public_key().to_bytes()) + .await + .expect("clear pre-existing active banner"); + let retained_user = nostr::Keys::generate().public_key(); + let removed_user = nostr::Keys::generate().public_key(); + let remote_retained_user = nostr::Keys::generate().public_key(); + let mut retained_rx = register_banner_req( + origin.clone(), + retained, + &retained_host, + retained_user, + "retained-banner", + ) + .await; + let mut removed_rx = register_banner_req( + origin.clone(), + removed, + &removed_host, + removed_user, + "removed-banner", + ) + .await; + let mut remote_retained_rx = register_banner_req( + receiver.clone(), + retained, + &retained_host, + remote_retained_user, + "remote-retained-banner", + ) + .await; + + tokio::time::sleep(std::time::Duration::from_millis(200)).await; + + let first_body = serde_json::json!({ + "severity": "info", + "text": "initial all", + "maxDisplays": 1, + "targetScope": "all" + }) + .to_string(); + let response = status_for( + origin.clone(), + Request::builder() + .method("PUT") + .uri("/banners/current") + .header(header::HOST, "admin.example") + .header( + header::AUTHORIZATION, + make_nostr_auth_put(&operator, "/banners/current", first_body.as_bytes()), + ) + .header(header::CONTENT_TYPE, "application/json") + .body(Body::from(first_body)) + .expect("request"), + ) + .await; + assert_eq!(response.status(), StatusCode::OK); + let retained_initial = recv_banner_frame(&mut retained_rx).await; + let removed_initial = recv_banner_frame(&mut removed_rx).await; + let remote_initial = recv_banner_frame(&mut remote_retained_rx).await; + assert_eq!(banner_text(&retained_initial), "initial all"); + assert_eq!(banner_text(&removed_initial), "initial all"); + assert_eq!(banner_text(&remote_initial), "initial all"); + assert!( + tokio::time::timeout(std::time::Duration::from_millis(150), retained_rx.recv()) + .await + .is_err() + ); + + let second_body = serde_json::json!({ + "severity": "warning", + "text": "retained only", + "maxDisplays": 1, + "targetScope": "communities", + "communityIds": [*retained.as_uuid()] + }) + .to_string(); + let response = status_for( + origin.clone(), + Request::builder() + .method("PUT") + .uri("/banners/current") + .header(header::HOST, "admin.example") + .header( + header::AUTHORIZATION, + make_nostr_auth_put(&operator, "/banners/current", second_body.as_bytes()), + ) + .header(header::CONTENT_TYPE, "application/json") + .body(Body::from(second_body)) + .expect("request"), + ) + .await; + assert_eq!(response.status(), StatusCode::OK); + + let retained_clear = recv_banner_frame(&mut retained_rx).await; + let retained_active = recv_banner_frame(&mut retained_rx).await; + assert!(is_disabled_banner(&retained_clear)); + assert_eq!(banner_text(&retained_active), "retained only"); + assert!(retained_clear.created_at < retained_active.created_at); + + let removed_clear = recv_banner_frame(&mut removed_rx).await; + assert!(is_disabled_banner(&removed_clear)); + assert!( + tokio::time::timeout(std::time::Duration::from_millis(150), removed_rx.recv()) + .await + .is_err() + ); + + let remote_clear = recv_banner_frame(&mut remote_retained_rx).await; + let remote_active = recv_banner_frame(&mut remote_retained_rx).await; + assert!(is_disabled_banner(&remote_clear)); + assert_eq!(banner_text(&remote_active), "retained only"); + assert_eq!(remote_clear.created_at, retained_clear.created_at); + assert_eq!(remote_active.created_at, retained_active.created_at); + + let response = status_for( + origin.clone(), + Request::builder() + .method("DELETE") + .uri("/banners/current") + .header(header::HOST, "admin.example") + .header( + header::AUTHORIZATION, + make_nostr_auth_delete(&operator, "/banners/current"), + ) + .body(Body::empty()) + .expect("request"), + ) + .await; + assert_eq!(response.status(), StatusCode::OK); + let retained_disable = recv_banner_frame(&mut retained_rx).await; + assert!(is_disabled_banner(&retained_disable)); + assert!(retained_active.created_at < retained_disable.created_at); + + origin_subscriber.abort(); + receiver_subscriber.abort(); + origin_fanout.abort(); + receiver_fanout.abort(); + } + #[tokio::test] async fn every_route_rejects_a_missing_credential_before_database_access() { let state = test_state().await; @@ -2111,6 +2439,57 @@ mod postgres_tests { nip98_state_with_replay(pubkeys, Arc::new(AlwaysFreshReplayGuard)).await } + async fn nip98_state_with_database_pool_and_redis( + pool: sqlx::PgPool, + redis_url: &str, + pubkeys: Vec, + ) -> Arc { + let mut config = crate::config::Config::from_env().expect("default config loads"); + config.require_relay_membership = false; + config.redis_url = redis_url.to_owned(); + config.read_database_url = None; + config.relay_operator_pubkeys = pubkeys; + if !config.relay_operator_pubkeys.is_empty() { + config.relay_operator_api_origin = Some("https://admin.example".to_string()); + } + config.admin = Some(crate::config::AdminConfig { + host: "admin.example".to_string(), + auth: crate::config::AdminAuth::Nip98, + web_dir: None, + }); + let db = buzz_db::Db::from_pool(pool.clone()); + let redis_pool = deadpool_redis::Config::from_url(&config.redis_url) + .create_pool(Some(deadpool_redis::Runtime::Tokio1)) + .expect("redis pool"); + let pubsub = Arc::new( + buzz_pubsub::PubSubManager::new(&config.redis_url, redis_pool.clone()) + .await + .expect("pubsub manager"), + ); + let audit = buzz_audit::AuditService::new(pool.clone()); + let auth = buzz_auth::AuthService::new(config.auth.clone()); + let search = buzz_search::SearchService::new(pool.clone()); + let workflow_engine = Arc::new(buzz_workflow::WorkflowEngine::new( + db.clone(), + buzz_workflow::WorkflowConfig::default(), + )); + let media_storage = buzz_media::MediaStorage::new(&config.media).expect("media storage"); + let (mut state, _audit_shutdown) = crate::state::AppState::new( + config, + db, + redis_pool, + audit, + pubsub, + auth, + search, + workflow_engine, + nostr::Keys::generate(), + media_storage, + ); + state.nip98_replay = Arc::new(AlwaysFreshReplayGuard); + Arc::new(state) + } + async fn nip98_state_with_replay( pubkeys: Vec, replay: Arc, diff --git a/crates/buzz-relay/src/handlers/event.rs b/crates/buzz-relay/src/handlers/event.rs index 01d6451fe90..03fafb3f9d6 100644 --- a/crates/buzz-relay/src/handlers/event.rs +++ b/crates/buzz-relay/src/handlers/event.rs @@ -2988,12 +2988,17 @@ mod relay_banner_fanout_tests { use tokio_util::sync::CancellationToken; use uuid::Uuid; - async fn setup_state() -> (Arc, sqlx::PgPool) { + async fn setup_state() -> ( + tokio::sync::MutexGuard<'static, ()>, + Arc, + sqlx::PgPool, + ) { + let guard = crate::test_support::RELAY_BANNER_TEST_LOCK.lock().await; let pool = sqlx::PgPool::connect(&crate::test_support::database_url()) .await .expect("connect test DB"); let state = crate::state::tests::test_state_with_database_pool(pool.clone()).await; - (state, pool) + (guard, state, pool) } async fn insert_test_community(pool: &sqlx::PgPool, host_prefix: &str) -> CommunityId { @@ -3039,10 +3044,85 @@ mod relay_banner_fanout_tests { (conn_id, rx) } + async fn recv_banner_event(rx: &mut mpsc::Receiver) -> nostr::Event { + let frame = rx.recv().await.expect("banner frame"); + let WsMessage::Text(frame) = frame else { + panic!("expected text frame"); + }; + let frame: serde_json::Value = serde_json::from_str(&frame).expect("EVENT frame JSON"); + assert_eq!(frame[0], "EVENT"); + serde_json::from_value(frame[2].clone()).expect("banner event") + } + + fn banner_content(event: &nostr::Event) -> serde_json::Value { + serde_json::from_str(&event.content).expect("banner content JSON") + } + + #[tokio::test] + #[ignore = "requires Postgres"] + async fn relay_banner_scope_shrink_clears_old_audience_before_new_active() { + let (_guard, state, pool) = setup_state().await; + let retained = insert_test_community(&pool, "banner-shrink-retained").await; + let removed = insert_test_community(&pool, "banner-shrink-removed").await; + let (_retained_conn, mut retained_rx) = register_banner_sub(&state, retained, vec![11; 32]); + let (_removed_conn, mut removed_rx) = register_banner_sub(&state, removed, vec![12; 32]); + + let initial = state + .db + .admin_upsert_relay_banner(buzz_db::RelayBannerUpsert { + severity: buzz_db::RelayBannerSeverity::Info, + message: "all communities".to_owned(), + max_displays: 1, + scope: buzz_db::RelayBannerScope::AllCommunities, + actor_pubkey: vec![1; 32], + }) + .await + .expect("insert initial banner"); + let replacement = state + .db + .admin_upsert_relay_banner(buzz_db::RelayBannerUpsert { + severity: buzz_db::RelayBannerSeverity::Warning, + message: "retained only".to_owned(), + max_displays: 1, + scope: buzz_db::RelayBannerScope::Communities(vec![retained]), + actor_pubkey: vec![1; 32], + }) + .await + .expect("replace banner"); + assert_eq!( + replacement.previous.as_ref().map(|banner| banner.id), + Some(initial.active.id) + ); + let disabled = replacement.previous.as_ref().expect("disabled previous"); + let disabled_event = + crate::api::banners::banner_disabled_event(&state.relay_keypair, disabled, None) + .expect("disabled event"); + let active_event = + crate::api::banners::banner_event(&state.relay_keypair, &replacement.active, None) + .expect("active event"); + + fan_out_relay_banner_event(&state, disabled, &disabled_event, false).await; + fan_out_relay_banner_event(&state, &replacement.active, &active_event, true).await; + + let retained_clear = recv_banner_event(&mut retained_rx).await; + assert!(retained_clear + .tags + .iter() + .any(|tag| tag.kind().to_string() == "status" && tag.content() == Some("disabled"))); + let retained_active = recv_banner_event(&mut retained_rx).await; + assert_eq!(banner_content(&retained_active)["text"], "retained only"); + let removed_clear = recv_banner_event(&mut removed_rx).await; + assert!(removed_clear + .tags + .iter() + .any(|tag| tag.kind().to_string() == "status" && tag.content() == Some("disabled"))); + assert!(removed_rx.try_recv().is_err()); + } + #[tokio::test] #[ignore = "requires Postgres"] async fn relay_banner_upsert_live_routes_only_to_eligible_users() { - let (state, pool) = setup_state().await; + let (_guard, state, pool) = setup_state().await; let included = insert_test_community(&pool, "banner-fanout-included").await; let excluded = insert_test_community(&pool, "banner-fanout-excluded").await; let (_eligible_conn, mut eligible_rx) = register_banner_sub(&state, included, vec![7; 32]); @@ -3057,7 +3137,8 @@ mod relay_banner_fanout_tests { actor_pubkey: vec![1; 32], }) .await - .expect("upsert banner"); + .expect("upsert banner") + .active; let event = crate::api::banners::banner_event(&state.relay_keypair, &banner, None) .expect("banner event"); @@ -3067,14 +3148,23 @@ mod relay_banner_fanout_tests { let WsMessage::Text(frame) = frame else { panic!("expected text frame"); }; - assert!(frame.contains("\"live\"")); + let frame: serde_json::Value = serde_json::from_str(&frame).expect("EVENT frame JSON"); + assert_eq!(frame[0], "EVENT"); + assert_eq!(frame[1], "banner"); + let content: serde_json::Value = serde_json::from_str( + frame[2]["content"] + .as_str() + .expect("banner event content string"), + ) + .expect("banner content JSON"); + assert_eq!(content["text"], "live"); assert!(excluded_rx.try_recv().is_err()); } #[tokio::test] #[ignore = "requires Postgres"] async fn relay_banner_disable_live_routes_using_pre_disable_eligibility() { - let (state, pool) = setup_state().await; + let (_guard, state, pool) = setup_state().await; let eligible_community = insert_test_community(&pool, "banner-disable-eligible").await; let excluded_community = insert_test_community(&pool, "banner-disable-excluded").await; let exhausted_user = vec![8; 32]; @@ -3096,7 +3186,8 @@ mod relay_banner_fanout_tests { actor_pubkey: vec![1; 32], }) .await - .expect("upsert banner"); + .expect("upsert banner") + .active; assert_eq!( state .db @@ -3126,13 +3217,21 @@ mod relay_banner_fanout_tests { fan_out_relay_banner_event(&state, &disabled, &event, false).await; + let frame = exhausted_rx + .recv() + .await + .expect("exhausted user clear frame"); + let WsMessage::Text(frame) = frame else { + panic!("expected text frame"); + }; + assert!(frame.contains("\"status\",\"disabled\"")); + assert!(frame.contains("\"scope\",\"communities\"")); let frame = eligible_rx.recv().await.expect("live disable frame"); let WsMessage::Text(frame) = frame else { panic!("expected text frame"); }; assert!(frame.contains("\"status\",\"disabled\"")); assert!(frame.contains("\"scope\",\"communities\"")); - assert!(exhausted_rx.try_recv().is_err()); assert!(excluded_rx.try_recv().is_err()); } } diff --git a/crates/buzz-relay/src/test_support.rs b/crates/buzz-relay/src/test_support.rs index a0b9f685374..011d2e82a49 100644 --- a/crates/buzz-relay/src/test_support.rs +++ b/crates/buzz-relay/src/test_support.rs @@ -8,6 +8,10 @@ pub(crate) fn database_url() -> String { .unwrap_or_else(|_| DEFAULT_DATABASE_URL.to_owned()) } +#[cfg(test)] +pub(crate) static RELAY_BANNER_TEST_LOCK: tokio::sync::Mutex<()> = + tokio::sync::Mutex::const_new(()); + #[cfg(test)] const CHILD_TEST_TIMEOUT: std::time::Duration = std::time::Duration::from_secs(10); #[cfg(test)] From 7d768e465ae0c4a579e78aafe94c077d0d7733ec Mon Sep 17 00:00:00 2001 From: coder 0 Date: Thu, 17 Sep 2026 15:15:01 -0400 Subject: [PATCH 4/4] Fix relay banner migration inventory Signed-off-by: coder 0 --- crates/buzz-db/src/runtime/migration.rs | 17 ++++++++++++++++- ...relay_banners.sql => 0047_relay_banners.sql} | 0 2 files changed, 16 insertions(+), 1 deletion(-) rename migrations/{0046_relay_banners.sql => 0047_relay_banners.sql} (100%) diff --git a/crates/buzz-db/src/runtime/migration.rs b/crates/buzz-db/src/runtime/migration.rs index 2564f116c38..5fa2d63a518 100644 --- a/crates/buzz-db/src/runtime/migration.rs +++ b/crates/buzz-db/src/runtime/migration.rs @@ -707,7 +707,7 @@ mod postgres_tests { let mut migrations: Vec<_> = MIGRATOR.iter().collect(); migrations.sort_by_key(|migration| migration.version); - assert_eq!(migrations.len(), 46); + assert_eq!(migrations.len(), 47); assert_eq!(migrations[0].version, 1); assert_eq!(&*migrations[0].description, "initial schema"); assert!(migrations[0] @@ -1292,6 +1292,21 @@ mod postgres_tests { .sql .as_str() .contains("CREATE TABLE storage_accounting_snapshots")); + + // Relay banners are deployment-global operator configuration and + // per-user display state. Keep them additive after storage accounting + // so the brownfield 0046 checksum from main remains immutable. + assert_eq!(migrations[46].version, 47); + let relay_banners = migrations[46].sql.as_str(); + assert!(relay_banners.contains("CREATE TABLE relay_banners")); + assert!(relay_banners.contains("CREATE TABLE relay_banner_communities")); + assert!(relay_banners.contains("CREATE TABLE relay_banner_user_state")); + assert!(relay_banners.contains("CREATE TABLE relay_banner_view_acks")); + assert!(relay_banners.contains("_operator_global_tables")); + assert!(relay_banners.contains("attach_community_write_fence('relay_banner_communities')")); + assert!(relay_banners.contains("attach_community_write_fence('relay_banner_user_state')")); + assert!(relay_banners.contains("attach_community_write_fence('relay_banner_view_acks')")); + assert!(!migrations[0].sql.as_str().contains("relay_banners")); // schema.sql exclusion list must match the restored (pre-0041) body. assert!( desired_schema.contains("'rate_limit_violations'\n ]::TEXT[])"), diff --git a/migrations/0046_relay_banners.sql b/migrations/0047_relay_banners.sql similarity index 100% rename from migrations/0046_relay_banners.sql rename to migrations/0047_relay_banners.sql