diff --git a/crates/buzz-db/src/store/event.rs b/crates/buzz-db/src/store/event.rs index cb45809eadb..ef0152543ec 100644 --- a/crates/buzz-db/src/store/event.rs +++ b/crates/buzz-db/src/store/event.rs @@ -10,8 +10,8 @@ use sqlx::{PgConnection, PgPool, Postgres, QueryBuilder, Row, Transaction}; use uuid::Uuid; use buzz_core::kind::{ - event_kind_i32, is_ephemeral, is_parameterized_replaceable, KIND_AUTH, KIND_EVENT_REMINDER, - KIND_HUDDLE_STARTED, SHARED_GATED_KINDS, + event_kind_i32, is_ephemeral, is_parameterized_replaceable, AUTHOR_ONLY_KINDS, KIND_AUTH, + KIND_EVENT_REMINDER, KIND_HUDDLE_STARTED, SHARED_GATED_KINDS, }; use buzz_core::{CommunityId, StoredEvent}; use buzz_datastore_tracing::datastore_span; @@ -113,6 +113,10 @@ pub struct EventQuery { /// SQL pushdown is sound. Keeping `event_visible_to_reader` as post-filter /// defense-in-depth catches any residual mismatch. pub shared_gated_reader: Option>, + /// Author-only visibility reader for [`query_events`]: exclude foreign + /// [`AUTHOR_ONLY_KINDS`] before ordering, offset and limit. This does not + /// affect `count_events`; callers must retain its existing fallback gate. + pub author_only_reader: Option>, } impl EventQuery { @@ -144,6 +148,7 @@ impl EventQuery { channel_ids_include_global: true, max_limit: None, shared_gated_reader: None, + author_only_reader: None, } } } @@ -682,6 +687,19 @@ pub(crate) async fn query_events_on( qb.push(")"); } + // Author-only visibility belongs before pagination, just like shared-gated + // visibility. Keep relay result checks as defense in depth. + if let Some(ref reader_bytes) = q.author_only_reader { + qb.push(format!(" AND ({col_prefix}kind NOT IN (")); + let mut sep = qb.separated(", "); + for kind in AUTHOR_ONLY_KINDS { + sep.push_bind(*kind as i32); + } + qb.push(format!(") OR {col_prefix}pubkey = ")); + qb.push_bind(reader_bytes.clone()); + qb.push(")"); + } + // Composite ordering for deterministic pagination across ALL callers of // query_events (WebSocket REQ, REST endpoints, canvas, notes, etc.). // The `id ASC` tiebreaker ensures stable results when events share the diff --git a/crates/buzz-relay/src/api/bridge.rs b/crates/buzz-relay/src/api/bridge.rs index 37c549610de..f10a0fc5de7 100644 --- a/crates/buzz-relay/src/api/bridge.rs +++ b/crates/buzz-relay/src/api/bridge.rs @@ -2529,6 +2529,10 @@ fn ban_json(b: &buzz_db::moderation::BanRecord) -> Value { }) } +#[cfg(test)] +#[path = "private_read_postgres_tests.rs"] +mod private_read_postgres_tests; + #[cfg(test)] mod postgres_tests { use super::*; @@ -3839,7 +3843,7 @@ mod postgres_tests { /// - Redis pool points at the local dev instance for the admission check. /// /// Returns `None` when local Postgres is not reachable. - async fn bridge_handler_test_state() -> Option> { + pub(super) async fn bridge_handler_test_state() -> Option> { let mut config = crate::config::Config::from_env().ok()?; config.database_url = crate::test_support::database_url(); // Use the real local Redis so enforce_http_admission can pass. diff --git a/crates/buzz-relay/src/api/private_read_postgres_tests.rs b/crates/buzz-relay/src/api/private_read_postgres_tests.rs new file mode 100644 index 00000000000..1d0f1c391f1 --- /dev/null +++ b/crates/buzz-relay/src/api/private_read_postgres_tests.rs @@ -0,0 +1,256 @@ +//! Existing author-only reads through HTTP and WebSocket storage seams. +use super::postgres_tests::bridge_handler_test_state; +use super::*; +use axum::{body::Body, http::Request}; +use buzz_core::kind::KIND_EVENT_REMINDER; +use nostr::{EventBuilder, Keys, Kind, Tag, Timestamp}; +use serde_json::json; +use tower::ServiceExt; + +async fn post( + state: &Arc, + host: &str, + path: &str, + keys: &Keys, + body: Value, + signed: bool, +) -> (StatusCode, Value) { + let mut request = Request::builder() + .method("POST") + .uri(path) + .header("host", host); + if signed { + let proof = EventBuilder::new(Kind::HttpAuth, "") + .tags([ + Tag::parse(["u", &format!("https://{host}{path}")]).unwrap(), + Tag::parse(["method", "POST"]).unwrap(), + ]) + .sign_with_keys(keys) + .unwrap(); + request = request.header( + "authorization", + format!( + "Nostr {}", + base64::engine::general_purpose::STANDARD + .encode(serde_json::to_vec(&proof).unwrap()) + ), + ); + } else { + request = request.header("x-pubkey", keys.public_key().to_hex()); + } + let response = crate::router::build_router(state.clone()) + .oneshot( + request + .body(Body::from(serde_json::to_vec(&body).unwrap())) + .unwrap(), + ) + .await + .unwrap(); + let status = response.status(); + let bytes = axum::body::to_bytes(response.into_body(), 1_048_576) + .await + .unwrap(); + (status, serde_json::from_slice(&bytes).unwrap()) +} + +fn drain(rx: &mut tokio::sync::mpsc::Receiver) -> Vec { + let mut frames = vec![]; + while let Ok(frame) = rx.try_recv() { + let axum::extract::ws::Message::Text(text) = frame else { + panic!("text frame") + }; + frames.push(serde_json::from_str(&text).unwrap()); + } + frames +} + +#[tokio::test] +#[ignore = "requires Postgres"] +async fn existing_private_reads_preserve_visible_pages_and_counts() { + let mut state = bridge_handler_test_state() + .await + .expect("test infrastructure"); + let host = format!("private-read-{}.example", uuid::Uuid::new_v4().simple()); + let community = state + .db + .ensure_configured_community(&host) + .await + .unwrap() + .id; + let tenant = TenantContext::resolved(community, &host); + let owner = Keys::generate(); + let outsider = Keys::generate(); + let now = Timestamp::now().as_secs(); + // Ingest the existing reminder envelope, not a new kind or test-only row. + // Terminal reminders may omit not_before; the relay never decrypts content. + let ciphertext = nostr::nips::nip44::encrypt( + owner.secret_key(), + &owner.public_key(), + r#"{"status":"done","note":"private reminder"}"#, + nostr::nips::nip44::Version::V2, + ) + .unwrap(); + let reminder = EventBuilder::new(Kind::Custom(KIND_EVENT_REMINDER as u16), &ciphertext) + .tags([Tag::parse(["d", &uuid::Uuid::new_v4().to_string()]).unwrap()]) + .custom_created_at(Timestamp::from(now)) + .sign_with_keys(&owner) + .unwrap(); + let mut public = vec![]; + for text in ["first public note", "second public note"] { + public.push( + EventBuilder::text_note(text) + .custom_created_at(Timestamp::from(now - 1)) + .sign_with_keys(&owner) + .unwrap(), + ); + } + public.sort_by_key(|event| event.id); + for event in std::iter::once(&reminder).chain(public.iter()) { + let (status, result) = post(&state, &host, "/events", &owner, json!(event), true).await; + assert_eq!(status, StatusCode::OK, "{result}"); + assert_eq!(result["accepted"], true, "{result}"); + } + let own = json!([{"kinds":[30300], "authors":[owner.public_key().to_hex()]}]); + let known = json!([{"ids":[reminder.id.to_hex()]}]); + let mixed = json!([{"kinds":[30300,1], "limit":1}]); + // Exercise the existing production NIP-98 contract and development X-Pubkey + // contract. Both must page/count using the reader selected by that mode. + for strict in [false, true] { + Arc::make_mut(&mut Arc::get_mut(&mut state).unwrap().config).require_auth_token = strict; + for route in ["/query", "/count"] { + if strict { + let (status, result) = post(&state, &host, route, &owner, own.clone(), false).await; + assert_eq!(status, StatusCode::UNAUTHORIZED, "{route}: {result}"); + } + let (status, result) = post(&state, &host, route, &outsider, own.clone(), strict).await; + assert_eq!(status, StatusCode::FORBIDDEN, "{result}"); + let (status, result) = post(&state, &host, route, &owner, own.clone(), strict).await; + assert_eq!(status, StatusCode::OK, "{result}"); + if route == "/query" { + assert_eq!(result[0]["id"], reminder.id.to_hex()); + } else { + assert_eq!(result["count"], 1); + } + let (status, result) = post( + &state, + &host, + route, + &outsider, + json!([{"kinds":[1]}]), + false, + ) + .await; + assert_eq!( + status, + if strict { + StatusCode::UNAUTHORIZED + } else { + StatusCode::OK + }, + "{result}" + ); + } + for who in [&owner, &outsider] { + let is_owner = who.public_key() == owner.public_key(); + let (status, rows) = post(&state, &host, "/query", who, mixed.clone(), strict).await; + assert_eq!(status, StatusCode::OK, "{rows}"); + assert_eq!(rows.as_array().unwrap().len(), 1); + assert_eq!( + rows[0]["id"], + if is_owner { + reminder.id.to_hex() + } else { + public[0].id.to_hex() + } + ); + // COUNT ignores the page limit, but must count only visible rows. + let (status, result) = post(&state, &host, "/count", who, mixed.clone(), strict).await; + assert_eq!(status, StatusCode::OK, "{result}"); + assert_eq!(result["count"], if is_owner { 3 } else { 2 }); + let (status, rows) = post(&state, &host, "/query", who, known.clone(), strict).await; + assert_eq!(status, StatusCode::OK, "{rows}"); + assert_eq!(rows.as_array().unwrap().len(), usize::from(is_owner)); + let (status, result) = post(&state, &host, "/count", who, known.clone(), strict).await; + assert_eq!(status, StatusCode::OK, "{result}"); + assert_eq!(result["count"], u64::from(is_owner)); + } + } + // Both offset and composite cursors must page over visible records, including ties. + for filters in [ + json!([{"kinds":[30300,1], "limit":1, "page":2}]), + json!([{"kinds":[30300,1], "limit":1, "until":now-1, "before_id":public[0].id.to_hex()}]), + ] { + let (status, rows) = post(&state, &host, "/query", &outsider, filters, true).await; + assert_eq!(status, StatusCode::OK, "{rows}"); + assert_eq!(rows.as_array().unwrap().len(), 1); + assert_eq!(rows[0]["id"], public[1].id.to_hex()); + } + // WS uses its existing authenticated principal, never an HTTP dev identity. + for who in [&owner, &outsider] { + let (mut conn, mut rx) = crate::connection::tests::test_conn_with_auth( + crate::connection::AuthState::Authenticated(buzz_auth::AuthContext { + pubkey: who.public_key(), + scopes: buzz_auth::Scope::all_known(), + channel_ids: None, + auth_method: buzz_auth::AuthMethod::Nip42, + agent_owner_pubkey: None, + }), + ); + Arc::get_mut(&mut conn).unwrap().tenant = tenant.clone(); + let filters = serde_json::from_value(mixed.clone()).unwrap(); + crate::handlers::req::handle_req( + "page".into(), + filters, + vec![], + conn.clone(), + state.clone(), + ) + .await; + let frames = drain(&mut rx); + let rows: Vec<_> = frames.iter().filter(|frame| frame[0] == "EVENT").collect(); + assert_eq!(rows.len(), 1, "{frames:?}"); + let is_owner = who.public_key() == owner.public_key(); + assert_eq!( + rows[0][2]["id"], + if is_owner { + reminder.id.to_hex() + } else { + public[0].id.to_hex() + } + ); + crate::handlers::count::handle_count( + "count".into(), + serde_json::from_value(mixed.clone()).unwrap(), + conn, + state.clone(), + ) + .await; + let frames = drain(&mut rx); + assert_eq!(frames[0][0], "COUNT", "{frames:?}"); + assert_eq!(frames[0][2]["count"], if is_owner { 3 } else { 2 }); + } + // More foreign reminders than COUNT's candidate budget must not turn an + // outsider's small visible count into a 'narrower constraints' error. Seed + // signed storage rows directly to avoid thousands of HTTP admission calls. + for _ in 0..crate::handlers::req::COUNT_FALLBACK_CANDIDATE_LIMIT { + let event = EventBuilder::new(reminder.kind, &ciphertext) + .tags([Tag::parse(["d", &uuid::Uuid::new_v4().to_string()]).unwrap()]) + .custom_created_at(Timestamp::from(now)) + .sign_with_keys(&owner) + .unwrap(); + assert!( + state + .db + .insert_event(community, &event, None) + .await + .unwrap() + .1 + ); + } + let (status, result) = post(&state, &host, "/count", &outsider, mixed.clone(), true).await; + assert_eq!(status, StatusCode::OK, "{result}"); + assert_eq!(result["count"], 2); + // The budget must still reject genuinely over-budget visible candidate sets. + let (status, _) = post(&state, &host, "/count", &owner, mixed.clone(), true).await; + assert_eq!(status, StatusCode::BAD_REQUEST); +} diff --git a/crates/buzz-relay/src/handlers/req.rs b/crates/buzz-relay/src/handlers/req.rs index 9a141e86ff1..257a9f4f1b4 100644 --- a/crates/buzz-relay/src/handlers/req.rs +++ b/crates/buzz-relay/src/handlers/req.rs @@ -357,6 +357,9 @@ pub async fn handle_req( let mut params = filter_to_query_params(filter, per_filter_channel, conn.tenant.community()); params.before_id = before_ids.get(idx).cloned().flatten(); + if filter_can_match_author_only_kinds(filter) { + params.author_only_reader = Some(pubkey_bytes.clone()); + } apply_channel_scope_to_query( &mut params, filter, @@ -826,12 +829,16 @@ async fn handle_search_req( /// Resolves accessible channels for the given pubkey and builds the query. pub async fn build_event_query_from_filter( filter: &Filter, - _pubkey_bytes: &[u8], + pubkey_bytes: &[u8], _state: &AppState, community: buzz_core::tenant::CommunityId, ) -> EventQuery { let channel_id = extract_channel_id_from_filter(filter); - filter_to_query_params(filter, channel_id, community) + let mut query = filter_to_query_params(filter, channel_id, community); + if filter_can_match_author_only_kinds(filter) { + query.author_only_reader = Some(pubkey_bytes.to_vec()); + } + query } /// Maximum SQL candidate rows a non-pushable COUNT filter may inspect before