diff --git a/crates/buzz-cli/src/client.rs b/crates/buzz-cli/src/client.rs index 95dece2fc6..ad26d24b3a 100644 --- a/crates/buzz-cli/src/client.rs +++ b/crates/buzz-cli/src/client.rs @@ -335,6 +335,29 @@ mod media_download_tests { } } +const QUERY_PAGE_SIZE: u32 = 500; + +fn advance_query_cursor( + filter: &mut serde_json::Value, + page: &[serde_json::Value], +) -> Result<(), CliError> { + let last = page + .last() + .expect("a full query page always has a last event"); + let created_at = last + .get("created_at") + .and_then(serde_json::Value::as_u64) + .ok_or_else(|| CliError::Other("query event missing created_at".into()))?; + let id = last + .get("id") + .and_then(serde_json::Value::as_str) + .filter(|id| id.len() == 64 && id.chars().all(|c| c.is_ascii_hexdigit())) + .ok_or_else(|| CliError::Other("query event missing valid id".into()))?; + filter["until"] = serde_json::json!(created_at); + filter["before_id"] = serde_json::json!(id); + Ok(()) +} + pub struct BuzzClient { http: reqwest::Client, relay_url: String, // base URL, no trailing slash, e.g. "https://relay.buzz.place" @@ -428,6 +451,54 @@ impl BuzzClient { } } + async fn query_pages( + &self, + mut filter: serde_json::Value, + limit: Option, + ) -> Result, CliError> { + let mut events = Vec::new(); + + while limit.is_none_or(|limit| events.len() < limit as usize) { + let page_limit = limit + .map(|limit| (limit as usize - events.len()).min(QUERY_PAGE_SIZE as usize)) + .unwrap_or(QUERY_PAGE_SIZE as usize); + filter["limit"] = serde_json::json!(page_limit); + + let raw = self.query(&filter).await?; + let page: Vec = serde_json::from_str(&raw) + .map_err(|e| CliError::Other(format!("failed to parse query response: {e}")))?; + let done = page.len() < page_limit; + + if !done { + advance_query_cursor(&mut filter, &page)?; + } + events.extend(page); + if done { + break; + } + } + + Ok(events) + } + + /// Query up to `limit` historical events, following the relay bridge's + /// composite `(until, before_id)` cursor across bounded result pages. + pub async fn query_paginated( + &self, + filter: serde_json::Value, + limit: u32, + ) -> Result, CliError> { + self.query_pages(filter, Some(limit)).await + } + + /// Query every historical event matching a filter across bounded pages. + pub async fn query_all( + &self, + filter: serde_json::Value, + ) -> Result, CliError> { + self.query_pages(filter, None).await + } + /// Execute a one-shot query via the HTTP bridge. /// `filter` is a Nostr filter object (will be wrapped in an array). /// Returns the raw JSON response (array of events). @@ -849,7 +920,35 @@ pub fn normalize_write_response(raw: &str) -> String { #[cfg(test)] mod tests { - use super::{create_response_with_id, extract_relay_response_field}; + use super::{advance_query_cursor, create_response_with_id, extract_relay_response_field}; + + #[test] + fn query_cursor_uses_last_events_composite_sort_key() { + let mut filter = serde_json::json!({"kinds": [39000], "limit": 500}); + let page = vec![ + serde_json::json!({"id": "a".repeat(64), "created_at": 20}), + serde_json::json!({"id": "b".repeat(64), "created_at": 10}), + ]; + + advance_query_cursor(&mut filter, &page).unwrap(); + + assert_eq!(filter["until"], serde_json::json!(10)); + assert_eq!(filter["before_id"], serde_json::json!("b".repeat(64))); + } + + #[test] + fn query_cursor_rejects_missing_or_malformed_sort_key() { + let mut filter = serde_json::json!({}); + assert!( + advance_query_cursor(&mut filter, &[serde_json::json!({"id": "a".repeat(64)})]) + .is_err() + ); + assert!(advance_query_cursor( + &mut filter, + &[serde_json::json!({"id": "not-an-event-id", "created_at": 10})] + ) + .is_err()); + } #[test] fn extract_relay_response_field_reads_response_message_json() { diff --git a/crates/buzz-cli/src/commands/channels.rs b/crates/buzz-cli/src/commands/channels.rs index b2eb8ccb16..18c67f928f 100644 --- a/crates/buzz-cli/src/commands/channels.rs +++ b/crates/buzz-cli/src/commands/channels.rs @@ -29,17 +29,16 @@ pub async fn cmd_list_channels( format: &crate::OutputFormat, ) -> Result<(), CliError> { let effective_limit = limit.unwrap_or(500); - let raw = if member == Some(true) { + let events = if member == Some(true) { // Step 1: find channel IDs where we're a member (kind:39002) let my_pk = client.keys().public_key().to_hex(); let member_filter = serde_json::json!({ "kinds": [39002], "#p": [my_pk], - "limit": effective_limit }); - let member_resp = client.query(&member_filter).await?; - let member_events: Vec = - serde_json::from_str(&member_resp).unwrap_or_default(); + let member_events = client + .query_paginated(member_filter, effective_limit) + .await?; let channel_ids: Vec = member_events .iter() .map(extract_d_tag) @@ -49,22 +48,21 @@ pub async fn cmd_list_channels( println!("[]"); return Ok(()); } - // Step 2: fetch kind:39000 metadata for those channels + // Step 2: fetch kind:39000 metadata for those channels. let metadata_filter = serde_json::json!({ "kinds": [39000], "#d": channel_ids, - "limit": effective_limit }); - client.query(&metadata_filter).await? + client + .query_paginated(metadata_filter, effective_limit) + .await? } else { let filter = serde_json::json!({ "kinds": [39000], - "limit": effective_limit }); - client.query(&filter).await? + client.query_paginated(filter, effective_limit).await? }; - let events: Vec = serde_json::from_str(&raw).unwrap_or_default(); let channels: Vec = events .iter() .filter(|e| { @@ -130,16 +128,8 @@ pub async fn cmd_search_channels( let filter = serde_json::json!({ "kinds": [39000], - "limit": limit, }); - let raw = client.query(&filter).await?; - - let events: serde_json::Value = serde_json::from_str(&raw) - .map_err(|e| CliError::Other(format!("failed to parse response: {e}")))?; - let Some(arr) = events.as_array() else { - println!("[]"); - return Ok(()); - }; + let arr = client.query_paginated(filter, limit).await?; let needle = query.to_ascii_lowercase(); let mut matches: Vec = arr @@ -337,14 +327,6 @@ pub async fn cmd_create_channel( Ok(()) } -/// Relay-side cap on a single historical query (`buzz-relay/src/handlers/req.rs`). -/// Never request above this — a clamped response compared against a larger -/// requested size would falsely look like "no more pages." -const RELAY_MAX_HISTORICAL_LIMIT: u32 = 2000; - -/// One page of a keyset-paginated kind:30177 scan. -const MANAGED_AGENT_PAGE_SIZE: u32 = RELAY_MAX_HISTORICAL_LIMIT; - /// A resolved live managed-agent instance backing a template persona slug. #[derive(Debug, Clone, PartialEq, Eq)] struct ResolvedAgent { @@ -424,62 +406,29 @@ async fn scan_managed_agents_by_owner( owner: &str, slugs: &HashSet<&str>, ) -> Result, CliError> { + let filter = serde_json::json!({ + "kinds": [KIND_MANAGED_AGENT], + "authors": [owner], + }); + let events = client.query_all(filter).await?; let mut found: Vec = Vec::new(); - let mut seen_event_ids: HashSet = HashSet::new(); - let mut cursor: Option<(i64, String)> = None; - - loop { - let mut filter = serde_json::json!({ - "kinds": [KIND_MANAGED_AGENT], - "authors": [owner], - "limit": MANAGED_AGENT_PAGE_SIZE, - }); - if let Some((until, ref before_id)) = cursor { - filter["until"] = serde_json::json!(until); - filter["before_id"] = serde_json::json!(before_id); - } - - let raw = client.query(&filter).await?; - let events: Vec = serde_json::from_str(&raw).map_err(|e| { - CliError::Other(format!("failed to parse managed-agent query response: {e}")) - })?; - - let page_len = events.len(); - for event in &events { - let Some(event_id) = event.get("id").and_then(|v| v.as_str()) else { - continue; - }; - if !seen_event_ids.insert(event_id.to_string()) { - continue; - } - let pubkey = extract_d_tag(event); - if pubkey.is_empty() { - continue; - } - let content: ManagedAgentContent = event - .get("content") - .and_then(|c| c.as_str()) - .and_then(|s| serde_json::from_str(s).ok()) - .unwrap_or(ManagedAgentContent { persona_id: None }); - let Some(persona_id) = content.persona_id else { - continue; - }; - if slugs.contains(persona_id.as_str()) { - found.push(ResolvedAgent { persona_id, pubkey }); - } - } - if (page_len as u32) < MANAGED_AGENT_PAGE_SIZE { - break; + for event in &events { + let pubkey = extract_d_tag(event); + if pubkey.is_empty() { + continue; } - let Some(last) = events.last() else { break }; - let Some(created_at) = last.get("created_at").and_then(|v| v.as_i64()) else { - break; - }; - let Some(last_id) = last.get("id").and_then(|v| v.as_str()) else { - break; + let content: ManagedAgentContent = event + .get("content") + .and_then(|c| c.as_str()) + .and_then(|s| serde_json::from_str(s).ok()) + .unwrap_or(ManagedAgentContent { persona_id: None }); + let Some(persona_id) = content.persona_id else { + continue; }; - cursor = Some((created_at, last_id.to_string())); + if slugs.contains(persona_id.as_str()) { + found.push(ResolvedAgent { persona_id, pubkey }); + } } Ok(found) diff --git a/crates/buzz-db/src/event.rs b/crates/buzz-db/src/event.rs index 71f3119943..5adcb05bdc 100644 --- a/crates/buzz-db/src/event.rs +++ b/crates/buzz-db/src/event.rs @@ -63,8 +63,9 @@ pub struct EventQuery { /// Restrict results to events with an `e` tag referencing any of these event IDs (hex). /// Uses JSONB containment (`tags @> ...`) against the `tags` column. pub e_tags: Option>, - /// Restrict results to events in any of these channels (multi-channel `IN` pushdown). - /// Used by NIP-45 COUNT to enforce channel access without fetching all rows. + /// Restrict results to events in any of these channels, while retaining + /// channel-less global events. Applied before SQL `LIMIT` so access-filtered + /// historical pages have exact exhaustion semantics. pub channel_ids: Option>, /// Override the default limit clamp (1000). Used by COUNT fallback path /// which needs to fetch all matching events for post-filter counting. @@ -1760,6 +1761,60 @@ mod tests { .expect("sign") } + fn make_event_at(kind: u16, content: &str, created_at: u64) -> nostr::Event { + EventBuilder::new(Kind::Custom(kind), content) + .custom_created_at(nostr::Timestamp::from(created_at)) + .sign_with_keys(&Keys::generate()) + .expect("sign timestamped event") + } + + #[tokio::test] + #[ignore = "requires Postgres"] + async fn access_scope_is_applied_before_historical_page_limit() { + let pool = setup_pool().await; + let community_uuid = make_test_community(&pool).await; + let community = CommunityId::from_uuid(community_uuid); + let accessible = make_test_channel(&pool, community_uuid, None).await; + let inaccessible = make_test_channel(&pool, community_uuid, None).await; + let base = 1_800_000_000; + + // This is the bridge underfetch shape: newer inaccessible candidates + // outnumber the requested page, while the visible match is older. + for offset in 10..13 { + let event = make_event_at(39_000, "newer inaccessible", base + offset); + insert_event(&pool, community, &event, Some(inaccessible)) + .await + .expect("insert inaccessible candidate"); + } + let global = make_event_at(39_000, "newer global", base + 2); + insert_event(&pool, community, &global, None) + .await + .expect("insert global candidate"); + let older_accessible = make_event_at(39_000, "older accessible", base + 1); + insert_event(&pool, community, &older_accessible, Some(accessible)) + .await + .expect("insert accessible candidate"); + + let events = query_events( + &pool, + &EventQuery { + kinds: Some(vec![39_000]), + channel_ids: Some(vec![accessible]), + limit: Some(2), + ..EventQuery::for_community(community) + }, + ) + .await + .expect("query access-scoped page"); + + assert_eq!(events.len(), 2, "visible page must be filled before EOF"); + assert_eq!(events[0].event.id, global.id, "global rows remain visible"); + assert_eq!( + events[1].event.id, older_accessible.id, + "older accessible row must not be hidden behind newer inaccessible rows" + ); + } + fn make_text_event(content: &str) -> nostr::Event { let keys = Keys::generate(); EventBuilder::new(Kind::Custom(9), content) diff --git a/crates/buzz-relay/src/api/bridge.rs b/crates/buzz-relay/src/api/bridge.rs index f30e24b794..6f8628111d 100644 --- a/crates/buzz-relay/src/api/bridge.rs +++ b/crates/buzz-relay/src/api/bridge.rs @@ -955,6 +955,11 @@ pub async fn query_events( tenant.community(), ) .await; + crate::handlers::req::apply_access_scope_to_query( + &mut query, + extract_channel_from_filter(filter), + &accessible_channels, + ); match extract_before_id(raw) { BeforeId::Malformed => { diff --git a/crates/buzz-relay/src/handlers/req.rs b/crates/buzz-relay/src/handlers/req.rs index dde877e70f..ecc3ba6777 100644 --- a/crates/buzz-relay/src/handlers/req.rs +++ b/crates/buzz-relay/src/handlers/req.rs @@ -289,8 +289,9 @@ pub async fn handle_req( }) .or(channel_id) }; - let params = + let mut params = filter_to_query_params(filter, per_filter_channel, conn.tenant.community()); + apply_access_scope_to_query(&mut params, per_filter_channel, &accessible_channels); (idx, per_filter_channel, params) }) .collect(); @@ -985,6 +986,20 @@ fn filter_to_query_params( } } +/// Push the caller's authorized channel set into logically global historical +/// queries so SQL `LIMIT` counts visible rows. Channel-less events remain in +/// scope by `EventQuery::channel_ids` contract; an explicit single-channel +/// filter keeps its narrower `channel_id` predicate. +pub(crate) fn apply_access_scope_to_query( + query: &mut EventQuery, + channel_id: Option, + accessible_channels: &[uuid::Uuid], +) { + if channel_id.is_none() { + query.channel_ids = Some(accessible_channels.to_vec()); + } +} + /// Extract a single channel UUID from filter generic tags, or `None` if the /// subscription is logically global. /// @@ -1219,6 +1234,33 @@ mod tests { use super::*; use nostr::{Alphabet, Filter, SingleLetterTag}; + #[test] + fn global_queries_push_access_scope_before_limit() { + let accessible = vec![uuid::Uuid::new_v4(), uuid::Uuid::new_v4()]; + let mut query = EventQuery::for_community(buzz_core::tenant::CommunityId::from_uuid( + uuid::Uuid::new_v4(), + )); + + apply_access_scope_to_query(&mut query, None, &accessible); + + assert_eq!(query.channel_ids.as_deref(), Some(accessible.as_slice())); + } + + #[test] + fn channel_scoped_queries_keep_exact_channel_predicate() { + let channel = uuid::Uuid::new_v4(); + let accessible = vec![channel, uuid::Uuid::new_v4()]; + let mut query = EventQuery::for_community(buzz_core::tenant::CommunityId::from_uuid( + uuid::Uuid::new_v4(), + )); + query.channel_id = Some(channel); + + apply_access_scope_to_query(&mut query, Some(channel), &accessible); + + assert!(query.channel_ids.is_none()); + assert_eq!(query.channel_id, Some(channel)); + } + /// S2 invariant: the bounded-concurrency pipeline (phase 2) must yield /// per-filter results in original filter order even when an earlier /// filter's DB query completes *after* a later one. `buffered` guarantees