diff --git a/Cargo.lock b/Cargo.lock index 421642e526d..2a689bf1c03 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -1123,6 +1123,7 @@ dependencies = [ "buzz-core", "buzz-datastore-tracing", "chrono", + "futures-util", "hex", "metrics", "metrics-util", diff --git a/crates/buzz-core/src/kind.rs b/crates/buzz-core/src/kind.rs index 4e1ab1c7f5e..45a7bf2e395 100644 --- a/crates/buzz-core/src/kind.rs +++ b/crates/buzz-core/src/kind.rs @@ -437,6 +437,8 @@ 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; +/// NIP-CW thread-mode query-time bounds, bound to a normalized newest-first thread request. +pub const KIND_THREAD_WINDOW_BOUNDS: u32 = 39007; /// Workflow definition (parameterized replaceable, d=workflow_uuid). pub const KIND_WORKFLOW_DEF: u32 = 30620; @@ -694,6 +696,7 @@ pub const ALL_KINDS: &[u32] = &[ KIND_NIP29_GROUP_ROLES, KIND_THREAD_SUMMARY, KIND_WINDOW_BOUNDS, + KIND_THREAD_WINDOW_BOUNDS, KIND_PRESENCE_UPDATE, KIND_TYPING_INDICATOR, KIND_HUDDLE_REACTION, @@ -839,6 +842,7 @@ pub const fn is_relay_only_kind(kind: u32) -> bool { | KIND_DM_VISIBILITY | KIND_THREAD_SUMMARY | KIND_WINDOW_BOUNDS + | KIND_THREAD_WINDOW_BOUNDS ) } @@ -916,6 +920,12 @@ mod tests { assert!(!is_relay_only_kind(KIND_NIP43_LEAVE_REQUEST)); } + #[test] + fn thread_window_bounds_is_relay_only() { + assert_eq!(KIND_THREAD_WINDOW_BOUNDS, 39007); + assert!(is_relay_only_kind(KIND_THREAD_WINDOW_BOUNDS)); + } + #[test] fn parameterized_replaceable_range() { assert!(!is_parameterized_replaceable(29999)); diff --git a/crates/buzz-core/src/lib.rs b/crates/buzz-core/src/lib.rs index 36dc772da3b..ec11dc0bf02 100644 --- a/crates/buzz-core/src/lib.rs +++ b/crates/buzz-core/src/lib.rs @@ -40,6 +40,8 @@ pub mod private_managed_agent; pub mod relay; /// Tenant identity — the server-resolved community key carried on scoped paths. pub mod tenant; +/// NIP-CW thread-mode normalized newest-first window contract. +pub mod thread_window; /// Schnorr signature and event ID verification. pub mod verification; diff --git a/crates/buzz-core/src/thread_window.rs b/crates/buzz-core/src/thread_window.rs new file mode 100644 index 00000000000..c7526d245ec --- /dev/null +++ b/crates/buzz-core/src/thread_window.rs @@ -0,0 +1,272 @@ +//! Strict, canonical NIP-CW thread-mode requests. Unknown constraints fail rather than +//! silently describing different rows from the ones the caller requested. + +use chrono::{DateTime, Utc}; +use serde::Serialize; +use serde_json::{json, Value}; +use sha2::{Digest, Sha256}; +use uuid::Uuid; + +/// Maximum reply scan candidates per window. +pub const MAX_LIMIT: u32 = 200; +/// Conversation row kinds. Edits, reactions, deletions and metadata are aux, +/// never reply-budget candidates, even if they have thread metadata. +pub const ROW_KINDS: [u32; 4] = [9, 40002, 45001, 45003]; + +/// A position in `(created_at DESC, id ASC)` order. +#[derive(Clone, Debug, PartialEq, Eq, Serialize)] +pub struct Cursor { + /// Nonnegative Unix seconds, representable by PostgreSQL/chrono. + pub created_at: i64, + /// Full lowercase hexadecimal event id. + pub id: String, +} + +impl Cursor { + /// Checked timestamp conversion (also used at the database boundary). + pub fn timestamp(&self) -> Result, String> { + DateTime::from_timestamp(self.created_at, 0) + .filter(|_| self.created_at >= 0) + .ok_or_else(|| "thread_window: cursor timestamp out of range".into()) + } +} + +/// Validated pagination and response-affecting arguments. Construct with +/// [`Request::parse`]; the database also validates public inputs defensively. +#[derive(Clone, Debug, PartialEq, Eq)] +pub struct Request { + /// Exactly one canonical channel UUID. + pub channel: Uuid, + /// Exactly one lowercase full root event id. + pub root: String, + /// Raw reply candidate budget, 1..=200. + pub limit: u32, + /// Maximum absolute depth below the root, 1..=100. + pub depth: u32, + /// Sorted, deduplicated conversation row kinds. + pub kinds: Vec, + /// Request upper bound; absent means newest page. + pub cursor: Option, + /// Include root/retained-row edits, reactions and deletion closure. + pub include_aux: bool, +} + +fn event_id(v: &Value) -> Result { + v.as_str() + .filter(|s| s.len() == 64 && s.bytes().all(|b| b.is_ascii_hexdigit())) + .map(str::to_ascii_lowercase) + .ok_or_else(|| "thread_window: expected a full 64-hex event id".into()) +} + +fn singleton<'a>(raw: &'a Value, key: &str) -> Result<&'a Value, String> { + raw.get(key) + .and_then(Value::as_array) + .filter(|a| a.len() == 1) + .and_then(|a| a.first()) + .ok_or_else(|| format!("thread_window requires exactly one {key}")) +} + +fn bounded(raw: &Value, key: &str, default: u32, max: u32) -> Result { + match raw.get(key) { + None => Ok(default), + Some(v) => v + .as_u64() + .filter(|n| (1..=u64::from(max)).contains(n)) + .map(|n| n as u32) + .ok_or_else(|| format!("thread_window: {key} must be an integer in 1..={max}")), + } +} + +impl Request { + /// Parse only an explicitly opted-in filter. Unsupported filter fields, + /// including legacy cursors, offsets and other bridge modes, are errors. + pub fn parse(raw: &Value) -> Result { + let object = raw.as_object().ok_or("thread_window: expected object")?; + const FIELDS: &[&str] = &[ + "thread_window", + "#h", + "#e", + "kinds", + "limit", + "depth_limit", + "until", + "before_id", + "include_aux", + ]; + if object.keys().any(|key| !FIELDS.contains(&key.as_str())) { + return Err("thread_window: unsupported or conflicting filter field".into()); + } + if raw.get("thread_window") != Some(&Value::Bool(true)) { + return Err("thread_window must be true".into()); + } + let channel = singleton(raw, "#h")? + .as_str() + .and_then(|s| Uuid::parse_str(s).ok()) + .ok_or("thread_window: #h must be a channel UUID")?; + let root = event_id(singleton(raw, "#e")?)?; + let limit = bounded(raw, "limit", 50, MAX_LIMIT)?; + let depth = bounded(raw, "depth_limit", 100, 100)?; + let mut kinds = raw + .get("kinds") + .and_then(Value::as_array) + .filter(|a| !a.is_empty() && a.len() <= ROW_KINDS.len()) + .ok_or("thread_window requires nonempty conversation kinds")? + .iter() + .map(|v| { + v.as_u64() + .and_then(|n| u32::try_from(n).ok()) + .filter(|n| ROW_KINDS.contains(n)) + .ok_or_else(|| { + "thread_window supports row kinds 9, 40002, 45001, 45003 only".to_string() + }) + }) + .collect::, _>>()?; + kinds.sort_unstable(); + kinds.dedup(); + let cursor = match (raw.get("until"), raw.get("before_id")) { + (None, None) => None, + (Some(ts), Some(id)) => { + let cursor = Cursor { + created_at: ts + .as_i64() + .filter(|n| *n >= 0) + .ok_or("thread_window: until must be nonnegative integer seconds")?, + id: event_id(id)?, + }; + cursor.timestamp()?; + Some(cursor) + } + _ => return Err("thread_window requires both until and before_id, or neither".into()), + }; + let include_aux = match raw.get("include_aux") { + None => false, + Some(v) => v + .as_bool() + .ok_or("thread_window: include_aux must be boolean")?, + }; + Ok(Self { + channel, + root, + limit, + depth, + kinds, + cursor, + include_aux, + }) + } + + /// Canonical NIP-CW thread-mode v1 request identity. SHA-256 over a compact JSON array + /// avoids ambiguous separators and binds all normalized response options. + /// The host is server-resolved; the reader is the authenticated lowercase + /// public key. Neither may be taken from filter-supplied fields. + pub fn binding(&self, resolved_host: &str, reader_hex: &str) -> String { + let cursor = self.cursor.as_ref().map(|c| json!([c.created_at, c.id])); + let canonical = json!([ + "tw", + 1, + "older", + resolved_host, + reader_hex, + self.channel.to_string(), + self.root, + self.limit, + self.depth, + self.kinds, + cursor, + self.include_aux + ]); + format!( + "tw:1:{}", + hex::encode(Sha256::digest(canonical.to_string().as_bytes())) + ) + } +} + +#[cfg(test)] +mod tests { + use super::*; + fn filter() -> Value { + json!({"thread_window":true,"#h":[Uuid::nil()],"#e":["ab".repeat(32)],"kinds":[9]}) + } + #[test] + fn rejects_malformed_and_conflicting_constraints() { + for (key, value) in [ + ("until", json!(1)), + ("before_id", json!("ab".repeat(32))), + ("limit", json!(0)), + ("limit", json!(201)), + ("depth_limit", json!(101)), + ("depth_limit", json!(-1)), + ("include_aux", json!("true")), + ("thread_cursor", json!(0)), + ("threadCursorId", json!(null)), + ("top_level", json!(false)), + ("page", json!(1)), + ("offset", json!(0)), + ("authors", json!([])), + ("since", json!(1)), + ("#p", json!([])), + ("kinds", json!([7])), + ("kinds", json!([])), + ("search", json!("hello")), + ("#h", json!([Uuid::nil(), Uuid::nil()])), + ("#e", json!(["bad"])), + ] { + let mut f = filter(); + f[key] = value; + assert!(Request::parse(&f).is_err(), "{f}"); + } + for ts in [ + json!(-1), + json!(1.5), + json!(null), + json!(u64::MAX), + json!(i64::MAX), + ] { + let mut f = filter(); + f["until"] = ts; + f["before_id"] = json!("a".repeat(64)); + assert!(Request::parse(&f).is_err(), "{f}"); + } + } + #[test] + fn binding_normalizes_defaults_and_binds_every_argument() { + let f = filter(); + let binding = |raw: &Value| { + Request::parse(raw) + .unwrap() + .binding("relay.example", &"ab".repeat(32)) + }; + let expected = binding(&f); + assert_eq!( + expected, + "tw:1:5252322dfd797ddb1d5f1150acd4cf914b9fc09e39bcf3048d25aed514b3546d" + ); + let mut explicit = f.clone(); + explicit["limit"] = json!(50); + explicit["depth_limit"] = json!(100); + explicit["include_aux"] = json!(false); + explicit["#e"] = json!(["AB".repeat(32)]); + assert_eq!(expected, binding(&explicit)); + for (key, val) in [ + ("limit", json!(49)), + ("depth_limit", json!(1)), + ("include_aux", json!(true)), + ("kinds", json!([40002])), + ("#e", json!(["cd".repeat(32)])), + ("#h", json!([Uuid::new_v4()])), + ] { + let mut changed = f.clone(); + changed[key] = val; + assert_ne!(expected, binding(&changed), "{key} must affect binding"); + } + let request = Request::parse(&f).unwrap(); + for (host, reader) in [("other.example", "ab"), ("relay.example", "cd")] { + assert_ne!(expected, request.binding(host, &reader.repeat(32))); + } + let mut changed = f; + changed["until"] = json!(0); + changed["before_id"] = json!("ab".repeat(32)); + assert_ne!(expected, binding(&changed)); + } +} diff --git a/crates/buzz-db/Cargo.toml b/crates/buzz-db/Cargo.toml index 380605728bb..47c8691bf02 100644 --- a/crates/buzz-db/Cargo.toml +++ b/crates/buzz-db/Cargo.toml @@ -17,6 +17,7 @@ serde_json = { workspace = true } uuid = { workspace = true } chrono = { workspace = true } hex = { workspace = true } +futures-util = { workspace = true } sha2 = { workspace = true } tracing = { workspace = true } thiserror = { workspace = true } diff --git a/crates/buzz-db/src/error.rs b/crates/buzz-db/src/error.rs index 79a371acb4f..e7d7f2d3246 100644 --- a/crates/buzz-db/src/error.rs +++ b/crates/buzz-db/src/error.rs @@ -68,6 +68,10 @@ pub enum DbError { #[error("read-state snapshot exceeds event or byte limit")] ReadStateSnapshotTooLarge, + /// A complete thread window exceeds its request-wide work allowance. + #[error("thread window exceeds {0} budget")] + ThreadWindowBudgetExceeded(&'static str), + /// A stored timestamp value could not be interpreted. #[error("invalid timestamp: {0}")] InvalidTimestamp(i64), diff --git a/crates/buzz-db/src/lib.rs b/crates/buzz-db/src/lib.rs index 36a70b758de..e70c5dd85b1 100644 --- a/crates/buzz-db/src/lib.rs +++ b/crates/buzz-db/src/lib.rs @@ -65,7 +65,7 @@ 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, read_state, relay_admin_actions, relay_invite, relay_members, relay_operators, - reminder, replaceable, storage_accounting, thread, usage, user, workflow, + reminder, replaceable, storage_accounting, thread, thread_window, usage, user, workflow, }; pub use allowlist::AllowlistEntry; diff --git a/crates/buzz-db/src/runtime/migration.rs b/crates/buzz-db/src/runtime/migration.rs index 6e5e14c2c62..20496ac1cfc 100644 --- a/crates/buzz-db/src/runtime/migration.rs +++ b/crates/buzz-db/src/runtime/migration.rs @@ -703,7 +703,12 @@ mod postgres_tests { let mut migrations: Vec<_> = MIGRATOR.iter().collect(); migrations.sort_by_key(|migration| migration.version); - assert_eq!(migrations.len(), 48); + assert_eq!(migrations.len(), 49); + assert_eq!(migrations[48].version, 49); + assert!(migrations[48] + .sql + .as_str() + .contains("idx_thread_metadata_window")); assert_eq!(migrations[0].version, 1); assert_eq!(&*migrations[0].description, "initial schema"); assert!(migrations[0] diff --git a/crates/buzz-db/src/runtime/tests.rs b/crates/buzz-db/src/runtime/tests.rs index cce1927be69..7d70d0a643e 100644 --- a/crates/buzz-db/src/runtime/tests.rs +++ b/crates/buzz-db/src/runtime/tests.rs @@ -3292,3 +3292,6 @@ async fn floor_guard_blocks_updates_that_move_rows_below_the_fence() { drop_scratch_db(&admin, pool, &name).await; } + +#[path = "tests/thread_window_postgres_tests.rs"] +mod thread_window_postgres_tests; diff --git a/crates/buzz-db/src/runtime/tests/thread_window_postgres_tests.rs b/crates/buzz-db/src/runtime/tests/thread_window_postgres_tests.rs new file mode 100644 index 00000000000..6062528d550 --- /dev/null +++ b/crates/buzz-db/src/runtime/tests/thread_window_postgres_tests.rs @@ -0,0 +1,317 @@ +use super::*; +use crate::thread_window::{AuxQuery, ScanBudget}; +use buzz_core::thread_window::Request; +use nostr::{EventBuilder, Kind, Tag, Timestamp}; + +fn request(channel: Uuid, root: &nostr::Event, upper: u64) -> Request { + Request::parse(&serde_json::json!({"thread_window":true,"#h":[channel], + "#e":[root.id.to_hex()],"kinds":[9],"limit":50, + "until":upper,"before_id":"00".repeat(32)})) + .unwrap() +} + +#[tokio::test] +#[ignore = "requires Postgres"] +async fn thread_window_upper_fence_terminal_snapshot_and_fallback() { + let admin = PgPool::connect(&admin_url().await).await.unwrap(); + let (writer, wname) = create_scratch_db(&admin, "tw_writer").await; + let (replica, rname) = create_scratch_db(&admin, "tw_replica").await; + let keys = nostr::Keys::generate(); + let cid = Uuid::new_v4(); + let channel = Uuid::new_v4(); + let base = 1_700_000_000; + let root = signed_event_at(&keys, "root", base); + let old = signed_event_at(&keys, "old", base + 10); + let middle = signed_event_at(&keys, "missing-middle", base + 20); + for pool in [&writer, &replica] { + seed_community_channel(pool, cid, channel, &keys).await; + insert_top_level(pool, cid, channel, &root).await; + insert_thread_reply(pool, cid, channel, &root, &old).await; + } + insert_thread_reply(&writer, cid, channel, &root, &middle).await; + let db = Db::from_pools(writer.clone(), replica.clone()); + let community = CommunityId::from_uuid(cid); + let mut req = request(channel, &root, base + 30); + // An old delivered tail does not justify an uncovered request upper bound. + db.fence() + .force_open_for_tests(chrono::DateTime::from_timestamp(base as i64 + 15, 0).unwrap()); + let (page, session) = db + .get_thread_window_with_session(community, &req, &mut ScanBudget::default()) + .await + .unwrap(); + assert!(!session.is_replica()); + assert_eq!( + page.rows.iter().map(|e| e.event.id).collect::>(), + [middle.id, old.id] + ); + assert!(!page.has_more); + drop(session); + // Counterfactual: over-claiming coverage demonstrably loses the middle row. + db.fence().force_open_for_tests(chrono::Utc::now()); + let (page, session) = db + .get_thread_window_with_session(community, &req, &mut ScanBudget::default()) + .await + .unwrap(); + assert!(session.is_replica()); + assert_eq!(page.rows.len(), 1); + assert_eq!(page.rows[0].event.id, old.id); + assert!(!page.has_more && page.next_cursor.is_none()); + drop(session); + // Default head route remains writer despite an open fence. + req.cursor = None; + let (head, session) = db + .get_thread_window_with_session(community, &req, &mut ScanBudget::default()) + .await + .unwrap(); + assert!(!session.is_replica()); + assert_eq!(head.rows[0].event.id, middle.id); + drop(session); + + // Complete the reply fixture, then lag *recent* edits and channel-less + // deletions. Coverage of old replies is not a freshness proof for aux. + insert_thread_reply(&replica, cid, channel, &root, &middle).await; + req = request(channel, &root, base + 30); + let (_, mut session) = db + .get_thread_window_with_session(community, &req, &mut ScanBudget::default()) + .await + .unwrap(); + assert!(session.is_replica()); + let edit = EventBuilder::new(Kind::Custom(40003), "new edit") + .tags([Tag::parse(["e", &old.id.to_hex()]).unwrap()]) + .custom_created_at(Timestamp::from(base + 1000)) + .sign_with_keys(&keys) + .unwrap(); + let deletion = EventBuilder::new(Kind::Custom(5), "") + .tags([Tag::parse(["e", &edit.id.to_hex()]).unwrap()]) + .custom_created_at(Timestamp::from(base + 1001)) + .sign_with_keys(&keys) + .unwrap(); + for pool in [&writer, &replica] { + event::insert_event(pool, community, &edit, Some(channel)) + .await + .unwrap(); + event::insert_event(pool, community, &deletion, None) + .await + .unwrap(); + } + let targets = [old.id.to_hex(), edit.id.to_hex()]; + let query = AuxQuery { + community, + targets: &targets, + kinds: &[5, 40003], + accessible: &[channel], + cursor: None, + }; + let mut budget = ScanBudget::default(); + let stale = session + .thread_window_aux(&query, &mut budget) + .await + .unwrap(); + assert!( + stale.events.is_empty(), + "held snapshot must not advance after page proof" + ); + let (_, mut fresh) = db + .get_thread_window_with_session(community, &req, &mut ScanBudget::default()) + .await + .unwrap(); + assert_eq!( + fresh + .thread_window_aux(&query, &mut budget) + .await + .unwrap() + .events + .len(), + 2, + "control: recent aux is visible on a fresh snapshot, with no reply timestamp bound" + ); + drop(fresh); + sqlx::query("SELECT pg_terminate_backend(pid) FROM pg_stat_activity WHERE datname=$1 AND pid<>pg_backend_pid()") + .bind(&rname).execute(&admin).await.unwrap(); + let degraded = session + .thread_window_aux(&query, &mut budget) + .await + .unwrap(); + assert_eq!(degraded.events.len(), 2); + assert!( + !session.is_replica(), + "fallback must permanently release the failed snapshot" + ); + drop(session); + // The shared router's reader acquisition failure must also reach the writer, + // not reinterpret an unavailable reader as an empty terminal window. + replica.close().await; + let (page, session) = db + .get_thread_window_with_session(community, &req, &mut ScanBudget::default()) + .await + .unwrap(); + assert!(!session.is_replica()); + assert_eq!(page.rows.len(), 2); + assert!(!page.has_more); + drop(session); + drop_scratch_db(&admin, replica, &rname).await; + drop_scratch_db(&admin, writer, &wname).await; +} + +#[tokio::test] +#[ignore = "requires Postgres"] +async fn migration_schema_thread_window_prebuild_validation_and_old_ledger() { + let admin = PgPool::connect(&admin_url().await).await.unwrap(); + let (pool, name) = create_scratch_db_through(&admin, "tw_prebuild", Some(48)).await; + for setup in [ + "CREATE INDEX idx_thread_metadata_window ON thread_metadata (community_id,root_event_id,event_created_at ASC,event_id ASC)", + // Simulate an invalid concurrent-build remnant only in this disposable superuser DB. + "CREATE INDEX idx_thread_metadata_window ON thread_metadata (community_id,root_event_id,event_created_at DESC,event_id ASC); \ + UPDATE pg_index SET indisvalid=false WHERE indexrelid='idx_thread_metadata_window'::regclass", + ] { + sqlx::raw_sql(setup).execute(&pool).await.unwrap(); + let error = migration::run_migrations(&pool).await.unwrap_err(); + assert!(error.to_string().contains("invalid or wrong definition"), + "must reject catalog shape, not merely time out: {error}"); + let version: i64 = sqlx::query_scalar("SELECT max(version) FROM _sqlx_migrations") + .fetch_one(&pool).await.unwrap(); + assert_eq!(version, 48); + sqlx::query("DROP INDEX CONCURRENTLY idx_thread_metadata_window").execute(&pool).await.unwrap(); + } + sqlx::query("CREATE INDEX CONCURRENTLY idx_thread_metadata_window ON thread_metadata (community_id,root_event_id,event_created_at DESC,event_id ASC)") + .execute(&pool).await.unwrap(); + let oid: i64 = sqlx::query_scalar("SELECT 'idx_thread_metadata_window'::regclass::oid::bigint") + .fetch_one(&pool) + .await + .unwrap(); + migration::run_migrations(&pool).await.unwrap(); + assert_eq!( + oid, + sqlx::query_scalar::<_, i64>("SELECT 'idx_thread_metadata_window'::regclass::oid::bigint") + .fetch_one(&pool) + .await + .unwrap() + ); + // Model the previous binary's exact embedded ledger, not run_to(48), + // which would still know version 49 and cannot test VersionMissing. + let current = sqlx::migrate::Migrator::new( + std::path::Path::new(env!("CARGO_MANIFEST_DIR")).join("../../migrations"), + ) + .await + .unwrap(); + let old = sqlx::migrate::Migrator::with_migrations( + current + .iter() + .filter(|m| m.version <= 48) + .cloned() + .collect(), + ); + assert!(matches!( + old.run(&pool).await, + Err(sqlx::migrate::MigrateError::VersionMissing(49)) + )); + drop_scratch_db(&admin, pool, &name).await; +} + +#[tokio::test] +#[ignore = "requires Postgres"] +async fn migration_schema_thread_window_prebuild_does_not_queue_behind_writer() { + let admin = PgPool::connect(&admin_url().await).await.unwrap(); + let (pool, name) = create_scratch_db_through(&admin, "tw_prebuild_writer", Some(48)).await; + let community = Uuid::new_v4(); + let channel = Uuid::new_v4(); + seed_community_channel(&pool, community, channel, &nostr::Keys::generate()).await; + sqlx::query("CREATE INDEX CONCURRENTLY idx_thread_metadata_window ON thread_metadata (community_id,root_event_id,event_created_at DESC,event_id ASC)") + .execute(&pool).await.unwrap(); + let oid: i64 = sqlx::query_scalar("SELECT 'idx_thread_metadata_window'::regclass::oid::bigint") + .fetch_one(&pool) + .await + .unwrap(); + + // Model an ingestion transaction already holding the table's writer lock. + // It stays held through every control: releasing it would hide the convoy. + let mut writer = pool.begin().await.unwrap(); + sqlx::query("LOCK TABLE thread_metadata IN ROW EXCLUSIVE MODE") + .execute(&mut *writer) + .await + .unwrap(); + let insert = + "INSERT INTO thread_metadata (community_id, channel_id, event_created_at, event_id) \ + VALUES ($1, $2, now(), $3)"; + let mut witness = pool.acquire().await.unwrap(); + sqlx::query("SET lock_timeout = '200ms'") + .execute(&mut *witness) + .await + .unwrap(); + sqlx::query(insert) + .bind(community) + .bind(channel) + .bind(vec![1_u8; 32]) + .execute(&mut *witness) + .await + .unwrap(); + + // Keep the exact migration transaction open after its SQL finishes so a + // second insert observes *all* locks it acquired. This is an explicit + // overlap barrier, not a race against a fast catalog-only migration. + let mut migration_tx = pool.begin().await.unwrap(); + let result = sqlx::raw_sql(include_str!( + "../../../../../migrations/0049_thread_window_index.sql" + )) + .execute(&mut *migration_tx) + .await; + let concurrent_insert = sqlx::query(insert) + .bind(community) + .bind(channel) + .bind(vec![2_u8; 32]) + .execute(&mut *witness) + .await; + migration_tx.rollback().await.unwrap(); + sqlx::query(insert) + .bind(community) + .bind(channel) + .bind(vec![3_u8; 32]) + .execute(&mut *witness) + .await + .unwrap(); + + // Also bind the public migrator: the advisory schema/destruction lock, + // ledger write and post-migration catalog checks must work with ingestion. + let production_result = migration::run_migrations(&pool).await; + sqlx::query(insert) + .bind(community) + .bind(channel) + .bind(vec![4_u8; 32]) + .execute(&mut *witness) + .await + .unwrap(); + let version: i64 = sqlx::query_scalar("SELECT max(version) FROM _sqlx_migrations") + .fetch_one(&mut *witness) + .await + .unwrap(); + let final_oid: i64 = + sqlx::query_scalar("SELECT 'idx_thread_metadata_window'::regclass::oid::bigint") + .fetch_one(&mut *witness) + .await + .unwrap(); + let count: i64 = + sqlx::query_scalar("SELECT count(*) FROM thread_metadata WHERE community_id=$1") + .bind(community) + .fetch_one(&mut *witness) + .await + .unwrap(); + drop(witness); + writer.rollback().await.unwrap(); + drop_scratch_db(&admin, pool, &name).await; + + assert!( + result.is_ok(), + "valid prebuild must bypass write-conflicting CREATE: {result:?}" + ); + assert!( + concurrent_insert.is_ok(), + "migration must allow metadata inserts: {concurrent_insert:?}" + ); + assert!( + production_result.is_ok(), + "production migrator must preserve ingestion progress: {production_result:?}" + ); + assert_eq!(version, 49); + assert_eq!(final_oid, oid, "prebuild must not be replaced"); + assert_eq!(count, 4, "all writer witnesses must persist"); +} diff --git a/crates/buzz-db/src/store/mod.rs b/crates/buzz-db/src/store/mod.rs index 3aeb2257aae..559460e15b9 100644 --- a/crates/buzz-db/src/store/mod.rs +++ b/crates/buzz-db/src/store/mod.rs @@ -52,6 +52,8 @@ pub mod replaceable; pub mod storage_accounting; /// Thread metadata persistence. pub mod thread; +/// Strict newest-first thread windows and auxiliary scans. +pub mod thread_window; /// Per-community usage rollup queries for Prometheus gauges. pub mod usage; /// User profile persistence. diff --git a/crates/buzz-db/src/store/thread_window.rs b/crates/buzz-db/src/store/thread_window.rs new file mode 100644 index 00000000000..659de340a69 --- /dev/null +++ b/crates/buzz-db/src/store/thread_window.rs @@ -0,0 +1,465 @@ +//! NIP-CW thread-mode newest-first reply windows. This path deliberately does not change +//! legacy forward threads or the permissive generic event reconstruction path. + +use buzz_core::{ + thread_window::{Cursor, Request, MAX_LIMIT, ROW_KINDS}, + CommunityId, StoredEvent, +}; +use chrono::{DateTime, Utc}; +use futures_util::TryStreamExt; +use sqlx::{PgConnection, PgPool, QueryBuilder, Row}; +use uuid::Uuid; + +use crate::{ + event::row_to_stored_event, Db, DbError, ReadSession, ReadSessionInner, Result, RouteDecision, + RoutePredicate, +}; +use buzz_datastore_tracing::datastore_span; + +/// One bounded raw scan; damaged reply rows can shorten `rows`, never the +/// authoritative scan bounds. +#[derive(Debug)] +pub struct ThreadWindow { + /// Reconstructed replies, newest first (id ascending for ties). + pub rows: Vec, + /// Another eligible raw reply exists after this page. + pub has_more: bool, + /// Last retained raw candidate, or None iff exhausted. + pub next_cursor: Option, + /// A supported conversation root exists in this channel, possibly as a + /// tombstone. Only then may the bridge serve rows, auxiliary events or bounds. + pub root_in_channel: bool, +} + +fn invalid(message: impl Into) -> DbError { + DbError::InvalidData(message.into()) +} + +fn cursor_key(cursor: &Cursor) -> Result<(DateTime, Vec)> { + let ts = cursor.timestamp().map_err(invalid)?; + let id = hex::decode(&cursor.id).map_err(|_| invalid("invalid thread cursor id"))?; + if id.len() != 32 { + return Err(invalid("invalid thread cursor length")); + } + Ok((ts, id)) +} + +fn scan_cursor(row: &sqlx::postgres::PgRow) -> Result { + let created_at: DateTime = row.try_get("created_at")?; + let id: Vec = row.try_get("id")?; + if id.len() != 32 || created_at.timestamp() < 0 { + return Err(invalid("unrepresentable thread scan position")); + } + Ok(Cursor { + created_at: created_at.timestamp(), + id: hex::encode(id), + }) +} + +/// Every new-mode SQL statement has a server-side deadline as well as the +/// bridge's whole-operation deadline. SET LOCAL cannot leak to pooled callers. +async fn set_deadline(conn: &mut PgConnection) -> Result<()> { + sqlx::query("SET LOCAL statement_timeout = '4000ms'") + .execute(&mut *conn) + .await?; + sqlx::query("SET LOCAL lock_timeout = '1000ms'") + .execute(&mut *conn) + .await?; + Ok(()) +} + +async fn select_window( + conn: &mut PgConnection, + community: CommunityId, + request: &Request, + budget: &mut ScanBudget, +) -> Result { + if !(1..=MAX_LIMIT).contains(&request.limit) + || !(1..=100).contains(&request.depth) + || request.kinds.is_empty() + || request.kinds.iter().any(|k| !ROW_KINDS.contains(k)) + { + return Err(invalid("invalid thread window arguments")); + } + let root = hex::decode(&request.root).map_err(|_| invalid("invalid thread root"))?; + if root.len() != 32 { + return Err(invalid("invalid thread root length")); + } + // Retain tombstones: deleting a root must not make its readable replies + // disappear. Never follow a root belonging to another channel/community. + let root_in_channel: bool = sqlx::query_scalar( + "SELECT EXISTS(SELECT 1 FROM events WHERE community_id = $1 AND channel_id = $2 AND id = $3 AND kind = ANY($4))") + .bind(community.as_uuid()).bind(request.channel).bind(&root) + .bind(ROW_KINDS.map(|k| k as i32).to_vec()) + .fetch_one(&mut *conn).await?; + if !root_in_channel { + return Ok(ThreadWindow { + rows: vec![], + has_more: false, + next_cursor: None, + root_in_channel, + }); + } + // Order metadata before event lookups. With stale statistics a root-wide + // sort can otherwise execute one lateral lookup for every reply first. + // OFFSET 0 keeps this boundary without limiting candidates prematurely. + let mut q = QueryBuilder::new( + "SELECT e.id, e.pubkey, e.created_at, e.kind, e.tags, e.content, e.sig, e.received_at, e.channel_id, \ + octet_length(e.content) + octet_length(e.tags::text) AS payload_bytes \ + FROM (SELECT tm.community_id, tm.event_id, tm.event_created_at \ + FROM thread_metadata tm WHERE tm.community_id = ", + ); + q.push_bind(community.as_uuid()) + .push(" AND tm.root_event_id = ") + .push_bind(root) + .push(" AND tm.channel_id = ") + .push_bind(request.channel) + .push(" AND tm.depth BETWEEN 1 AND ") + .push_bind(request.depth as i32); + if let Some(cursor) = &request.cursor { + let (ts, id) = cursor_key(cursor)?; + // The redundant upper timestamp bound gives PostgreSQL an index range + // start instead of filtering all newer keys at a deep cursor. + q.push(" AND tm.event_created_at <= ") + .push_bind(ts) + .push(" AND (tm.event_created_at < ") + .push_bind(ts) + .push(" OR (tm.event_created_at = ") + .push_bind(ts) + .push(" AND tm.event_id > ") + .push_bind(id) + .push("))"); + } + // Fetch by the event PK before testing visibility, so missing statistics + // cannot turn a selective live/kind index into a root-wide scan per reply. + // PK uniqueness makes the inner LIMIT 1 exact. All eligibility predicates + // still precede the *outer* limit+1 that establishes page bounds. + q.push(" ORDER BY tm.event_created_at DESC, tm.event_id ASC OFFSET 0) tm \ + JOIN LATERAL (SELECT e.id, e.pubkey, e.created_at, e.kind, e.tags, e.content, e.sig, e.received_at, e.channel_id, e.deleted_at \ + FROM events e WHERE e.community_id = tm.community_id \ + AND e.created_at = tm.event_created_at AND e.id = tm.event_id LIMIT 1) e ON true \ + WHERE e.channel_id = ") + .push_bind(request.channel) + .push(" AND e.deleted_at IS NULL AND e.kind = ANY(") + .push_bind(request.kinds.iter().map(|k| *k as i32).collect::>()) + .push(") ORDER BY tm.event_created_at DESC, tm.event_id ASC LIMIT ") + .push_bind(i64::from(request.limit) + 1); + // Selectivity varies drastically with root/depth/cursor. A cached generic + // plan can sort the entire root (and exceeded the SQL deadline in paging + // tests). An unnamed statement keeps parameter-aware planning for this + // bounded query without changing pooled session or legacy settings. + // Charge every consumed payload before reconstruction, including damaged + // rows and the limit+1 probe. The caller owns this ledger across all + // windows, auxiliary scans and replica-to-writer retries. + let mut raw = q.build().persistent(false).fetch(&mut *conn); + let mut rows = Vec::new(); + let mut retained = 0; + let mut last_cursor = None; + let mut has_more = false; + while let Some(row) = raw.try_next().await? { + budget.rows(1)?; + let payload_bytes: i32 = row.try_get("payload_bytes")?; + budget.bytes(payload_bytes as usize)?; + if retained == request.limit { + has_more = true; + break; + } + last_cursor = Some(scan_cursor(&row)?); + retained += 1; + if let Some(event) = row_to_stored_event(row)? { + rows.push(event); + } + } + let next_cursor = if has_more { last_cursor } else { None }; + Ok(ThreadWindow { + rows, + has_more, + next_cursor, + root_in_channel, + }) +} + +async fn writer_window( + pool: &PgPool, + community: CommunityId, + request: &Request, + budget: &mut ScanBudget, +) -> Result { + let conn = crate::observability::acquire_writer( + pool, + crate::observability::WriterOperation::SubscriptionHistory, + ) + .await?; + let mut tx = sqlx::Transaction::begin(conn, None).await?; + set_deadline(&mut tx).await?; + let window = select_window(&mut tx, community, request, budget).await?; + tx.rollback().await?; + Ok(window) +} + +impl Db { + /// Newest-first thread page and its proved replica transaction. Backward + /// cursors supply an upper bound, including terminal pages; no forward + /// thread "last delivered row is newest" inference is used. Head routing + /// inherits the existing default-off budget. Writer follow-ups are pooled, + /// not a snapshot spanning the response or subsequent history pages. The + /// caller must share `budget` with every window and auxiliary scan in the request. + #[datastore_span(name = "get_thread_window", system = "postgresql")] + pub async fn get_thread_window_with_session( + &self, + community: CommunityId, + request: &Request, + budget: &mut ScanBudget, + ) -> Result<(ThreadWindow, ReadSession)> { + let cursor = request.cursor.as_ref().map(cursor_key).transpose()?; + let path = if cursor.is_some() { + "thread_window_cursor" + } else { + "thread_window_head" + }; + if let RouteDecision::Replica(mut tx, _, reason) = self + .route_read( + path, + RoutePredicate::from_channel_cursor(request.channel, &cursor), + crate::observability::ReaderOperation::SubscriptionHistory, + ) + .await + { + let result = async { + set_deadline(&mut tx).await?; + select_window(&mut tx, community, request, budget).await + } + .await; + match result { + Ok(window) if !window.root_in_channel => { + // The cursor covers reply insertions, not a root authored + // later than its replies. Recheck missing anchors on writer. + Self::record_route(path, "writer", "missing_root"); + } + Ok(window) => { + Self::record_route(path, "replica", reason); + return Ok(( + window, + ReadSession { + inner: ReadSessionInner::Replica { + tx, + writer: self.pool.clone(), + }, + }, + )); + } + Err(error @ (DbError::InvalidData(_) | DbError::ThreadWindowBudgetExceeded(_))) => { + return Err(error); + } + Err(error) => { + tracing::warn!(%error, path, "replica thread window failed; re-running on writer"); + Self::record_route(path, "writer", "replica_error"); + } + } + } + let window = writer_window(&self.pool, community, request, budget).await?; + Ok(( + window, + ReadSession { + inner: ReadSessionInner::Writer(self.pool.clone()), + }, + )) + } +} + +/// Strict auxiliary scan. Targets are chunked by the bridge, limiting SQL +/// expression size. Access is applied before the probe, including channel-less +/// deletions. Deleted aux payloads are NOT returned, but their IDs are retained +/// to discover deletions-of-aux on the second hop. +pub struct AuxQuery<'a> { + /// Host-bound community. + pub community: CommunityId, + /// At most 200 retained target IDs (never the reply sentinel). + pub targets: &'a [String], + /// Auxiliary kinds for this hop. + pub kinds: &'a [u32], + /// Fresh writer-authorized channels for this page. + pub accessible: &'a [Uuid], + /// Last raw auxiliary scan position, not last delivered event. + pub cursor: Option, +} + +/// Fixed raw auxiliary page budget (the probe is one extra row). +pub const AUX_LIMIT: usize = 1000; + +/// Aggregate reply and auxiliary scan allowance for one query, not one window. +/// Bytes include content and tags; the bridge separately bounds serialized output. +/// Passed through replica retries so degradation cannot reset the allowance. +#[derive(Default)] +pub struct ScanBudget { + queries: usize, + rows: usize, + bytes: usize, +} + +impl ScanBudget { + fn query(&mut self) -> Result<()> { + self.queries += 1; + if self.queries > 64 { + return Err(DbError::ThreadWindowBudgetExceeded("query")); + } + Ok(()) + } + + fn bytes(&mut self, count: usize) -> Result<()> { + self.bytes = self.bytes.saturating_add(count); + if self.bytes > 8 * 1024 * 1024 { + return Err(DbError::ThreadWindowBudgetExceeded("payload byte")); + } + Ok(()) + } + + fn rows(&mut self, count: usize) -> Result<()> { + self.rows = self.rows.saturating_add(count); + if self.rows > 8192 { + return Err(DbError::ThreadWindowBudgetExceeded("raw row")); + } + Ok(()) + } +} + +/// A strict auxiliary page with raw traversal metadata. +pub struct AuxPage { + /// Live events. Unreconstructable live auxiliary events cause an error. + pub events: Vec, + /// Raw IDs including tombstones, for deletion-of-aux discovery. + pub target_ids: Vec, + /// Last retained scan position when another raw candidate exists. + pub next_cursor: Option, +} + +async fn select_aux( + conn: &mut PgConnection, + query: &AuxQuery<'_>, + budget: &mut ScanBudget, +) -> Result { + if query.targets.is_empty() || query.targets.len() > 200 { + return Err(invalid("invalid thread auxiliary target batch")); + } + budget.query()?; + let mut q = QueryBuilder::new( + "SELECT id, pubkey, created_at, kind, tags, content, sig, received_at, channel_id, deleted_at, \ + octet_length(content) + octet_length(tags::text) AS payload_bytes \ + FROM events WHERE community_id = "); + q.push_bind(query.community.as_uuid()) + .push(" AND ((channel_id IS NULL AND kind IN (5, 9005)) OR channel_id = ANY(") + .push_bind(query.accessible) + .push("))") + .push(" AND kind = ANY(") + .push_bind(query.kinds.iter().map(|k| *k as i32).collect::>()) + .push(") AND ("); + for (index, target) in query.targets.iter().enumerate() { + if index != 0 { + q.push(" OR "); + } + q.push("tags @> ") + .push_bind(serde_json::json!([["e", target]])); + } + // JSONB containment is an indexable prefilter, not a positional tag match. + q.push( + ") AND EXISTS (SELECT 1 FROM jsonb_array_elements(tags) tag \ + WHERE tag->>0 = 'e' AND tag->>1 = ANY(", + ) + .push_bind(query.targets) + .push("))"); + if let Some(cursor) = &query.cursor { + let (ts, id) = cursor_key(cursor)?; + q.push(" AND (created_at < ") + .push_bind(ts) + .push(" OR (created_at = ") + .push_bind(ts) + .push(" AND id > ") + .push_bind(id) + .push("))"); + } + q.push(" ORDER BY created_at DESC, id ASC LIMIT ") + .push_bind(AUX_LIMIT as i64 + 1); + // Never collect a full raw page: 1,001 ingest-valid edits can contain + // 250 MiB. Charge each row (including tombstones and the probe) before + // reconstruction, and preserve this request-wide ledger across retries. + let mut raw = q.build().fetch(&mut *conn); + let mut events = Vec::new(); + let mut target_ids = Vec::new(); + let mut last_cursor = None; + let mut next_cursor = None; + while let Some(row) = raw.try_next().await? { + budget.rows(1)?; + let payload_bytes: i32 = row.try_get("payload_bytes")?; + budget.bytes(payload_bytes as usize)?; + if target_ids.len() == AUX_LIMIT { + next_cursor = last_cursor; + break; + } + // Validate IDs even for deleted payloads: no ambiguous continuation. + let cursor = scan_cursor(&row)?; + target_ids.push(cursor.id.clone()); + last_cursor = Some(cursor); + if row + .try_get::>, _>("deleted_at")? + .is_none() + { + events + .push(row_to_stored_event(row)?.ok_or_else(|| { + invalid("cannot reconstruct required thread auxiliary event") + })?); + } + } + Ok(AuxPage { + events, + target_ids, + next_cursor, + }) +} + +async fn writer_aux( + pool: &PgPool, + query: &AuxQuery<'_>, + budget: &mut ScanBudget, +) -> Result { + let conn = crate::observability::acquire_writer( + pool, + crate::observability::WriterOperation::SubscriptionHistory, + ) + .await?; + let mut tx = sqlx::Transaction::begin(conn, None).await?; + set_deadline(&mut tx).await?; + let page = select_aux(&mut tx, query, budget).await?; + tx.rollback().await?; + Ok(page) +} + +impl ReadSession { + /// Strict NIP-CW thread-mode auxiliary scan on the same proved snapshot, with the + /// existing permanent writer degradation on mid-request replica failure. + #[datastore_span(name = "thread_window_aux", system = "postgresql")] + pub async fn thread_window_aux( + &mut self, + query: &AuxQuery<'_>, + budget: &mut ScanBudget, + ) -> Result { + let writer = match &mut self.inner { + ReadSessionInner::Replica { tx, writer } => match select_aux(tx, query, budget).await { + Ok(page) => return Ok(page), + Err(error @ (DbError::InvalidData(_) | DbError::ThreadWindowBudgetExceeded(_))) => { + return Err(error); + } + Err(error) => { + tracing::warn!(%error, "thread auxiliary read failed; degrading to writer"); + metrics::counter!("buzz_db_read_session_degraded").increment(1); + writer.clone() + } + }, + ReadSessionInner::Writer(pool) => return writer_aux(pool, query, budget).await, + }; + self.inner = ReadSessionInner::Writer(writer.clone()); + writer_aux(&writer, query, budget).await + } +} + +#[cfg(test)] +mod postgres_tests; diff --git a/crates/buzz-db/src/store/thread_window/postgres_tests.rs b/crates/buzz-db/src/store/thread_window/postgres_tests.rs new file mode 100644 index 00000000000..a2bda65123a --- /dev/null +++ b/crates/buzz-db/src/store/thread_window/postgres_tests.rs @@ -0,0 +1,571 @@ +use super::*; +use buzz_core::channel::{ChannelType, ChannelVisibility}; +use nostr::{EventBuilder, Keys, Kind, Tag, Timestamp}; + +async fn fixture() -> (Db, CommunityId, Uuid, Keys, nostr::Event) { + let pool = PgPool::connect(&crate::test_support::database_url()) + .await + .unwrap(); + if std::env::var("BUZZ_TEST_SCHEMA_MODE").as_deref() == Ok("migration") { + crate::migration::run_migrations(&pool).await.unwrap(); + } + let db = Db::from_pool(pool); + let community = db + .ensure_configured_community(&format!("tw-{}.local", Uuid::new_v4())) + .await + .unwrap() + .id; + let keys = Keys::generate(); + let channel = Uuid::new_v4(); + db.create_channel_with_id( + community, + channel, + "window", + ChannelType::Stream, + ChannelVisibility::Private, + None, + &keys.public_key().to_bytes(), + None, + ) + .await + .unwrap(); + let root = make_event(&keys, channel, 9, "root", None, Timestamp::now().as_secs()); + db.insert_event(community, &root, Some(channel)) + .await + .unwrap(); + (db, community, channel, keys, root) +} + +fn make_event( + keys: &Keys, + channel: Uuid, + kind: u16, + content: &str, + target: Option<&nostr::Event>, + ts: u64, +) -> nostr::Event { + let mut tags = vec![Tag::parse(["h", &channel.to_string()]).unwrap()]; + if let Some(target) = target { + tags.push(Tag::parse(["e", &target.id.to_hex(), "", "reply"]).unwrap()); + } + EventBuilder::new(Kind::Custom(kind), content) + .tags(tags) + .custom_created_at(Timestamp::from(ts)) + .sign_with_keys(keys) + .unwrap() +} + +fn request(channel: Uuid, root: &nostr::Event, limit: u32) -> Request { + Request::parse(&serde_json::json!({"thread_window":true,"#h":[channel], + "#e":[root.id.to_hex()], "kinds":[9], "limit":limit})) + .unwrap() +} + +async fn replies( + db: &Db, + cid: CommunityId, + channel: Uuid, + keys: &Keys, + root: &nostr::Event, + count: usize, +) -> Vec { + let mut events = Vec::new(); + for n in 0..count { + // Large same-second groups as well as timestamp boundaries. + events.push(make_event( + keys, + channel, + 9, + &format!("reply {n}"), + Some(root), + root.created_at.as_secs() + (n / 120) as u64, + )); + } + // Batch fixture writes, preserving real signed events and production table + // constraints. Production selection is exercised through the Db seam. + for batch in events.chunks(500) { + let mut q = QueryBuilder::new("INSERT INTO events (community_id,id,pubkey,created_at,kind,tags,content,sig,received_at,channel_id) "); + q.push_values(batch, |mut b, event| { + b.push_bind(cid.as_uuid()) + .push_bind(event.id.to_bytes().to_vec()) + .push_bind(event.pubkey.to_bytes().to_vec()) + .push_bind(DateTime::from_timestamp(event.created_at.as_secs() as i64, 0).unwrap()) + .push_bind(9_i32) + .push_bind(serde_json::to_value(&event.tags).unwrap()) + .push_bind(&event.content) + .push_bind(event.sig.serialize().to_vec()) + .push_bind(Utc::now()) + .push_bind(channel); + }); + q.build().execute(&db.pool).await.unwrap(); + let mut q = QueryBuilder::new("INSERT INTO thread_metadata (community_id,event_created_at,event_id,channel_id,parent_event_id,parent_event_created_at,root_event_id,root_event_created_at,depth,broadcast) "); + q.push_values(batch, |mut b, event| { + let root_ts = DateTime::from_timestamp(root.created_at.as_secs() as i64, 0).unwrap(); + b.push_bind(cid.as_uuid()) + .push_bind(DateTime::from_timestamp(event.created_at.as_secs() as i64, 0).unwrap()) + .push_bind(event.id.to_bytes().to_vec()) + .push_bind(channel) + .push_bind(root.id.to_bytes().to_vec()) + .push_bind(root_ts) + .push_bind(root.id.to_bytes().to_vec()) + .push_bind(root_ts) + .push_bind(1_i32) + .push_bind(false); + }); + q.build().execute(&db.pool).await.unwrap(); + } + events.sort_by(|a, b| { + b.created_at + .cmp(&a.created_at) + .then_with(|| a.id.cmp(&b.id)) + }); + events +} + +async fn assert_pages(db: &Db, cid: CommunityId, mut req: Request, expected: &[nostr::Event]) { + let mut delivered = Vec::new(); + loop { + let (window, _) = db + .get_thread_window_with_session(cid, &req, &mut ScanBudget::default()) + .await + .unwrap(); + assert_eq!(window.has_more, window.next_cursor.is_some()); + assert!(window.rows.len() <= req.limit as usize); + delivered.extend(window.rows.iter().map(|row| row.event.id)); + if let Some(next) = window.next_cursor { + assert_ne!(req.cursor.as_ref(), Some(&next)); + req.cursor = Some(next); + } else { + break; + } + assert!(delivered.len() <= expected.len()); + } + assert_eq!(delivered, expected.iter().map(|e| e.id).collect::>()); +} + +async fn cardinalities() { + let (db, cid, ch, keys, root) = fixture().await; + // Keep bulk-ingest statistics deliberately stale in both schema paths. + // Otherwise autoanalyze can hide a root-wide sort + repeated broad event + // scans that exceeded the production four-second deadline at 10k replies. + sqlx::raw_sql( + "ALTER TABLE thread_metadata SET (autovacuum_enabled=false); \ + DO $$ DECLARE part regclass; BEGIN \ + FOR part IN SELECT inhrelid::regclass FROM pg_inherits \ + WHERE inhparent='events'::regclass LOOP \ + EXECUTE format('ALTER TABLE %s SET (autovacuum_enabled=false)', part); \ + END LOOP; \ + END $$;", + ) + .execute(&db.pool) + .await + .unwrap(); + for count in [0, 1, 50, 51, 501, 10_000] { + let root = make_event( + &keys, + ch, + 9, + &format!("root {count}"), + None, + root.created_at.as_secs(), + ); + db.insert_event(cid, &root, Some(ch)).await.unwrap(); + let expected = replies(&db, cid, ch, &keys, &root, count).await; + assert_pages(&db, cid, request(ch, &root, 50), &expected).await; + let (head, _) = db + .get_thread_window_with_session( + cid, + &request(ch, &root, 50), + &mut ScanBudget::default(), + ) + .await + .unwrap(); + assert_eq!(head.has_more, count > 50); + assert_eq!(head.rows.len(), count.min(50)); + } + let index: (bool,bool,String) = sqlx::query_as("SELECT indisvalid,indisready,pg_get_indexdef(indexrelid) FROM pg_index WHERE indexrelid='idx_thread_metadata_window'::regclass") + .fetch_one(&db.pool).await.unwrap(); + assert!(index.0 && index.1); + assert!(index + .2 + .ends_with("(community_id, root_event_id, event_created_at DESC, event_id)")); +} + +#[tokio::test] +#[ignore = "requires Postgres"] +async fn thread_window_cardinalities_and_index() { + cardinalities().await; +} + +#[tokio::test] +#[ignore = "requires Postgres"] +async fn migration_schema_thread_window_cardinalities_and_index() { + cardinalities().await; +} + +#[tokio::test] +#[ignore = "requires Postgres"] +async fn thread_window_predicates_corruption_and_legacy() { + let (db, cid, ch, keys, root) = fixture().await; + let mut expected = replies(&db, cid, ch, &keys, &root, 51).await; + // The 50th raw row is damaged. It must STILL be the page's cursor, not + // the 49th delivered row; the sentinel cannot leak into the page. + let broken = expected[49].id; + sqlx::query("UPDATE events SET sig='\\x00' WHERE community_id=$1 AND id=$2") + .bind(cid.as_uuid()) + .bind(broken.to_bytes().to_vec()) + .execute(&db.pool) + .await + .unwrap(); + let req = request(ch, &root, 50); + let (page, _) = db + .get_thread_window_with_session(cid, &req, &mut ScanBudget::default()) + .await + .unwrap(); + assert_eq!(page.rows.len(), 49); + assert!(page.has_more); + assert_eq!(page.next_cursor.as_ref().unwrap().id, broken.to_hex()); + assert!(!page.rows.iter().any(|r| r.event.id == expected[50].id)); + expected.remove(49); + assert_pages(&db, cid, req, &expected).await; + // Put 600 ineligible rows before 51 eligible ones. Each rejection class + // exceeds a page independently; an early metadata LIMIT manufactures EOF. + let other_channel = Uuid::new_v4(); + db.create_channel_with_id( + cid, + other_channel, + "other", + ChannelType::Stream, + ChannelVisibility::Private, + None, + &keys.public_key().to_bytes(), + None, + ) + .await + .unwrap(); + let dense_root = make_event(&keys, ch, 9, "dense root", None, root.created_at.as_secs()); + db.insert_event(cid, &dense_root, Some(ch)).await.unwrap(); + let rows = replies(&db, cid, ch, &keys, &dense_root, 651).await; + for (class, sql) in [ + "UPDATE events SET deleted_at=now() WHERE community_id=$1 AND id=ANY($2)", + "UPDATE events SET kind=40002 WHERE community_id=$1 AND id=ANY($2)", + "UPDATE events SET channel_id=$3 WHERE community_id=$1 AND id=ANY($2)", + "UPDATE thread_metadata SET channel_id=$3 WHERE community_id=$1 AND event_id=ANY($2)", + "UPDATE thread_metadata SET depth=2 WHERE community_id=$1 AND event_id=ANY($2)", + "UPDATE thread_metadata SET root_event_id=$3 WHERE community_id=$1 AND event_id=ANY($2)", + ] + .into_iter() + .enumerate() + { + let ids: Vec<_> = rows[..600] + .iter() + .skip(class) + .step_by(6) + .map(|row| row.id.to_bytes().to_vec()) + .collect(); + let query = sqlx::query(sql).bind(cid.as_uuid()).bind(ids); + let query = match class { + 2 | 3 => query.bind(other_channel), + 5 => query.bind(root.id.to_bytes().to_vec()), + _ => query, + }; + assert_eq!(query.execute(&db.pool).await.unwrap().rows_affected(), 100); + } + let root = dense_root; + let expected = &rows[600..]; + let mut req = request(ch, &root, 50); + req.depth = 1; + let (page, _) = db + .get_thread_window_with_session(cid, &req, &mut ScanBudget::default()) + .await + .unwrap(); + assert_eq!(page.rows.len(), 50); + assert!(page.has_more); + assert_eq!(page.next_cursor.unwrap().id, expected[49].id.to_hex()); + assert_pages(&db, cid, req.clone(), expected).await; + sqlx::query("UPDATE events SET deleted_at=now() WHERE community_id=$1 AND id=$2") + .bind(cid.as_uuid()) + .bind(expected[50].id.to_bytes().to_vec()) + .execute(&db.pool) + .await + .unwrap(); + let expected = &expected[..50]; + let (page, _) = db + .get_thread_window_with_session(cid, &req, &mut ScanBudget::default()) + .await + .unwrap(); + assert_eq!(page.rows.len(), 50); + assert!(!page.has_more && page.next_cursor.is_none()); + // Deleting the root does not hide the descendants. + sqlx::query("UPDATE events SET deleted_at=now() WHERE community_id=$1 AND id=$2") + .bind(cid.as_uuid()) + .bind(root.id.to_bytes().to_vec()) + .execute(&db.pool) + .await + .unwrap(); + assert_pages(&db, cid, req.clone(), expected).await; + req.channel = Uuid::new_v4(); + let (page, _) = db + .get_thread_window_with_session(cid, &req, &mut ScanBudget::default()) + .await + .unwrap(); + assert!(page.rows.is_empty() && !page.has_more && !page.root_in_channel); + let legacy = db + .get_thread_replies(cid, &root.id.to_bytes(), None, 100, None) + .await + .unwrap(); + assert!(legacy + .windows(2) + .all(|p| (p[0].created_at, &p[0].event_id) < (p[1].created_at, &p[1].event_id))); +} + +#[tokio::test] +#[ignore = "requires Postgres"] +async fn thread_aux_tombstone_ids_and_raw_cursor_survive_damaged_payload() { + let (db, cid, ch, keys, root) = fixture().await; + let rows = replies(&db, cid, ch, &keys, &root, 1001).await; + // Change kinds only in the isolated corruption fixture. Reconstruction is + // structural in the existing store; signatures are not reverified on read. + sqlx::query("UPDATE events SET kind=40003 WHERE community_id=$1 AND id<>$2") + .bind(cid.as_uuid()) + .bind(root.id.to_bytes().to_vec()) + .execute(&db.pool) + .await + .unwrap(); + let targets = vec![root.id.to_hex()]; + let accessible = [ch]; + let mut query = AuxQuery { + community: cid, + targets: &targets, + kinds: &[40003], + accessible: &accessible, + cursor: None, + }; + let mut session = ReadSession { + inner: ReadSessionInner::Writer(db.pool.clone()), + }; + sqlx::query("UPDATE events SET sig='\\x00', deleted_at=now() WHERE community_id=$1 AND id=$2") + .bind(cid.as_uuid()) + .bind(rows[500].id.to_bytes().to_vec()) + .execute(&db.pool) + .await + .unwrap(); + // A deleted damaged payload still supplies its ID for delete-of-aux. + let first = session + .thread_window_aux(&query, &mut ScanBudget::default()) + .await + .unwrap(); + assert_eq!(first.events.len(), 999); + assert_eq!(first.target_ids.len(), 1000); + assert!(first.target_ids.contains(&rows[500].id.to_hex())); + assert!(first.next_cursor.is_some()); + query.cursor = first.next_cursor; + let tail = session + .thread_window_aux(&query, &mut ScanBudget::default()) + .await + .unwrap(); + assert_eq!(tail.events.len(), 1); + assert!(tail.next_cursor.is_none()); +} + +#[tokio::test] +#[ignore = "requires Postgres"] +async fn thread_aux_query_budget_fails_on_65th_real_scan() { + let (db, community, channel, _, root) = fixture().await; + let targets = [root.id.to_hex()]; + let query = AuxQuery { + community, + targets: &targets, + kinds: &[7], + accessible: &[channel], + cursor: None, + }; + let mut session = ReadSession { + inner: ReadSessionInner::Writer(db.pool.clone()), + }; + let mut budget = ScanBudget::default(); + for _ in 0..64 { + assert!(session + .thread_window_aux(&query, &mut budget) + .await + .unwrap() + .events + .is_empty()); + } + let result = session.thread_window_aux(&query, &mut budget).await; + assert!(matches!( + result, + Err(DbError::ThreadWindowBudgetExceeded("query")) + )); + assert_eq!(budget.queries, 65); + assert_eq!(budget.rows, 0); +} + +#[tokio::test] +#[ignore = "requires Postgres"] +async fn thread_aux_byte_budget_stops_consumption_and_survives_retry() { + let (db, community, channel, keys, root) = fixture().await; + // Each edit is signed and within ingest's 256 KiB content limit. The + // router regression exercises the same size through signed POST /events. + for n in 0..40 { + let edit = make_event( + &keys, + channel, + 40003, + &"x".repeat(256 * 1024), + Some(&root), + root.created_at.as_secs() + n, + ); + db.insert_event(community, &edit, Some(channel)) + .await + .unwrap(); + } + let targets = [root.id.to_hex()]; + let query = AuxQuery { + community, + targets: &targets, + kinds: &[40003], + accessible: &[channel], + cursor: None, + }; + let mut session = ReadSession { + inner: ReadSessionInner::Writer(db.pool.clone()), + }; + let mut budget = ScanBudget::default(); + let result = session.thread_window_aux(&query, &mut budget).await; + assert!(matches!( + result, + Err(DbError::ThreadWindowBudgetExceeded("payload byte")) + )); + assert_eq!( + budget.rows, 32, + "stop consuming at the first over-budget payload, not the full page" + ); + assert!(budget.bytes > 8 * 1024 * 1024); + let result = session.thread_window_aux(&query, &mut budget).await; + assert!(matches!( + result, + Err(DbError::ThreadWindowBudgetExceeded("payload byte")) + )); + assert_eq!( + budget.rows, 33, + "a retry cannot reset the request-wide allowance" + ); +} + +#[tokio::test] +#[ignore = "requires Postgres"] +async fn thread_reply_byte_budget_stops_before_collecting_page() { + let (db, community, channel, keys, root) = fixture().await; + // Real signed replies at the ingest content ceiling, not oversized DB rows. + for n in 0..201 { + let reply = make_event( + &keys, + channel, + 9, + &"x".repeat(256 * 1024), + Some(&root), + root.created_at.as_secs() + n, + ); + let ts = DateTime::from_timestamp(reply.created_at.as_secs() as i64, 0).unwrap(); + let root_ts = DateTime::from_timestamp(root.created_at.as_secs() as i64, 0).unwrap(); + db.insert_event_with_thread_metadata( + community, + &reply, + Some(channel), + Some(crate::event::ThreadMetadataParams { + event_id: &reply.id.to_bytes(), + event_created_at: ts, + channel_id: channel, + parent_event_id: Some(&root.id.to_bytes()), + parent_event_created_at: Some(root_ts), + root_event_id: Some(&root.id.to_bytes()), + root_event_created_at: Some(root_ts), + depth: 1, + broadcast: false, + }), + ) + .await + .unwrap(); + } + let mut req = request(channel, &root, 200); + req.include_aux = false; + let mut budget = ScanBudget::default(); + let result = db + .get_thread_window_with_session(community, &req, &mut budget) + .await; + assert!(matches!( + result, + Err(DbError::ThreadWindowBudgetExceeded("payload byte")) + )); + assert_eq!( + budget.rows, 32, + "stop on first excess payload, not after collecting 201 replies" + ); + assert!(budget.bytes > 8 * 1024 * 1024); + let result = db + .get_thread_window_with_session(community, &req, &mut budget) + .await; + assert!(matches!( + result, + Err(DbError::ThreadWindowBudgetExceeded("payload byte")) + )); + assert_eq!(budget.rows, 33, "retry must preserve consumed allowance"); + + // A smaller page succeeds. The probe is charged but never delivered, and + // the cursor is the last retained raw candidate rather than the probe. + req.limit = 20; + let mut budget = ScanBudget::default(); + let (page, mut session) = db + .get_thread_window_with_session(community, &req, &mut budget) + .await + .unwrap(); + assert_eq!(budget.rows, 21); + assert_eq!(page.rows.len(), 20); + assert!(page.has_more); + assert_eq!( + page.next_cursor.as_ref().unwrap().id, + page.rows[19].event.id.to_hex() + ); + let first_bytes = budget.bytes; + // Even an empty auxiliary page shares the existing reply allowance. + let targets = [root.id.to_hex()]; + session + .thread_window_aux( + &AuxQuery { + community, + targets: &targets, + kinds: &[7], + accessible: &[channel], + cursor: None, + }, + &mut budget, + ) + .await + .unwrap(); + assert_eq!(budget.bytes, first_bytes); + drop(session); + let result = db + .get_thread_window_with_session(community, &req, &mut budget) + .await; + assert!(matches!( + result, + Err(DbError::ThreadWindowBudgetExceeded("payload byte")) + )); + assert_eq!( + budget.rows, 32, + "the next window cannot reset the batch allowance" + ); + + // The probe alone can cross the limit; do not silently omit its charge. + req.limit = 31; + let mut budget = ScanBudget::default(); + let result = db + .get_thread_window_with_session(community, &req, &mut budget) + .await; + assert!(matches!( + result, + Err(DbError::ThreadWindowBudgetExceeded("payload byte")) + )); + assert_eq!(budget.rows, 32); +} diff --git a/crates/buzz-relay/src/api/bridge.rs b/crates/buzz-relay/src/api/bridge.rs index 6327ee166d0..4714ce13438 100644 --- a/crates/buzz-relay/src/api/bridge.rs +++ b/crates/buzz-relay/src/api/bridge.rs @@ -24,6 +24,7 @@ use crate::state::AppState; use super::{api_error, internal_error, not_found}; mod thread_roots; +mod thread_window; pub(crate) async fn enforce_http_admission( state: &AppState, @@ -1191,6 +1192,7 @@ async fn query_events_authed( // depth_limit, feed_types) that nostr::Filter silently drops. let raw_filters: Vec = serde_json::from_slice(body) .map_err(|e| api_error(StatusCode::BAD_REQUEST, &format!("invalid filters: {e}")))?; + let thread_windows = thread_window::parse(&raw_filters)?; let filters: Vec = raw_filters .iter() .map(|v| serde_json::from_value(v.clone())) @@ -1221,6 +1223,26 @@ async fn query_events_authed( )); } + if thread_windows.iter().any(Option::is_some) { + if thread_windows.iter().any(Option::is_none) { + return Err(api_error( + StatusCode::BAD_REQUEST, + "thread_window cannot mix with other query modes", + )); + } + return tokio::time::timeout( + thread_window::DEADLINE, + thread_window::query_batch(state, tenant, &pubkey, thread_windows.iter().flatten()), + ) + .await + .map_err(|_| { + api_error( + StatusCode::SERVICE_UNAVAILABLE, + "thread window deadline exceeded", + ) + })? + .map(|events| Json(Value::Array(events))); + } if read_state_snapshot::requested(&raw_filters) { return read_state_snapshot::query(state, tenant, &pubkey, &raw_filters).await; } diff --git a/crates/buzz-relay/src/api/bridge/thread_window.rs b/crates/buzz-relay/src/api/bridge/thread_window.rs new file mode 100644 index 00000000000..793a541c0d6 --- /dev/null +++ b/crates/buzz-relay/src/api/bridge/thread_window.rs @@ -0,0 +1,244 @@ +//! NIP-CW thread-mode bridge adapter. Legacy thread and channel-window code stays separate. + +use std::{collections::HashSet, time::Duration}; + +use axum::{http::StatusCode, Json}; +use buzz_core::{thread_window::Request, TenantContext}; +use buzz_db::thread_window::{AuxQuery, ScanBudget}; +use serde_json::{json, Value}; + +use super::{event_in_accessible_channel, WINDOW_AUX_DELETE_KINDS, WINDOW_AUX_KINDS}; +use crate::{ + api::{api_error, internal_error}, + state::AppState, +}; + +type Error = (StatusCode, Json); +/// Shared deadline across all thread-window filters in a /query request. +pub(super) const DEADLINE: Duration = Duration::from_secs(8); +const MAX_BYTES: usize = 8 * 1024 * 1024; +/// One ledger shared by all opted-in filters, including replica fallback. +#[derive(Default)] +pub(super) struct Budget { + bytes: usize, + scan: ScanBudget, +} + +/// Validate before search/presence/other extension dispatch can swallow the +/// opt-in. Absent/false preserves legacy handling; any other value is invalid. +pub(super) fn parse(filters: &[Value]) -> Result>, Error> { + let mut count = 0; + filters + .iter() + .map(|raw| match raw.get("thread_window") { + None | Some(Value::Bool(false)) => Ok(None), + Some(_) => { + count += 1; + if count > 4 { + return Err(api_error( + StatusCode::BAD_REQUEST, + "at most four thread windows per query", + )); + } + Request::parse(raw) + .map(Some) + .map_err(|e| api_error(StatusCode::BAD_REQUEST, &e)) + } + }) + .collect() +} + +fn unavailable(message: &str) -> Error { + api_error(StatusCode::SERVICE_UNAVAILABLE, message) +} + +fn database_error(context: &str, error: buzz_db::DbError) -> Error { + if let buzz_db::DbError::ThreadWindowBudgetExceeded(_) = &error { + return unavailable(&format!("{error}; reduce window work before retrying")); + } + // Pool acquisition and PostgreSQL's statement/lock budgets may expire + // before the outer HTTP deadline. They are retryable, not internal faults. + let timed_out = match &error { + buzz_db::DbError::Sqlx(sqlx::Error::PoolTimedOut) => true, + buzz_db::DbError::Sqlx(sqlx::Error::Database(error)) => { + matches!(error.code().as_deref(), Some("57014" | "55P03")) + } + _ => false, + }; + if timed_out { + tracing::warn!(%error, context, "thread window database timeout"); + unavailable("thread database timeout; retry window") + } else { + internal_error(&format!("thread {context}: {error}")) + } +} + +fn append(events: &mut Vec, budget: &mut Budget, event: &nostr::Event) -> Result<(), Error> { + let value = serde_json::to_value(event) + .map_err(|e| internal_error(&format!("thread serialize: {e}")))?; + budget.bytes = budget.bytes.saturating_add(value.to_string().len() + 1); + if budget.bytes > MAX_BYTES { + return Err(unavailable("thread window exceeds response byte budget")); + } + events.push(value); + Ok(()) +} + +/// Authorize the entire batch against one writer access set, then refresh it +/// once before releasing any output. A later window must never suppress only +/// its own rows while releasing an earlier window built before revocation. +pub(super) async fn query_batch<'a>( + state: &AppState, + tenant: &TenantContext, + reader: &nostr::PublicKey, + requests: impl IntoIterator, +) -> Result, Error> { + let accessible = state + .db + .get_accessible_channel_ids(tenant.community(), &reader.to_bytes()) + .await + .map_err(|e| database_error("access", e))?; + let mut budget = Budget::default(); + let mut events = Vec::new(); + for request in requests { + if accessible.contains(&request.channel) { + events.extend(query(state, tenant, reader, request, &accessible, &mut budget).await?); + } + } + let current = state + .db + .get_accessible_channel_ids(tenant.community(), &reader.to_bytes()) + .await + .map_err(|e| database_error("final access", e))?; + // Grants can expose auxiliary events omitted from the original closure; + // revocations can invalidate earlier windows or their cross-channel aux. + if accessible.iter().collect::>() != current.iter().collect::>() { + return Err(unavailable("thread authorization changed; retry query")); + } + Ok(events) +} + +async fn query( + state: &AppState, + tenant: &TenantContext, + reader: &nostr::PublicKey, + request: &Request, + accessible: &[uuid::Uuid], + budget: &mut Budget, +) -> Result, Error> { + let (window, mut session) = state + .db + .get_thread_window_with_session(tenant.community(), request, &mut budget.scan) + .await + .map_err(|e| database_error("window", e))?; + // An unsupported, missing or out-of-scope root is not a served window. + // In particular, never sign false exhaustion for an unsupported root kind. + if !window.root_in_channel { + return Ok(vec![]); + } + let reader_bytes = reader.to_bytes(); + let visible = |se: &buzz_core::StoredEvent| { + event_in_accessible_channel(se, accessible) + && crate::handlers::req::event_visible_to_reader(&se.event, &reader_bytes) + }; + let mut events = Vec::new(); + let page_start_bytes = budget.bytes; + budget.bytes += 2; + let mut targets = vec![request.root.clone()]; + for row in &window.rows { + if !visible(row) { + // This would contradict the SQL's channel and row-kind predicates. + // Do not issue authoritative bounds for a mismatched selection. + return Err(unavailable( + "thread row authorization changed during selection", + )); + } + targets.push(row.event.id.to_hex()); + append(&mut events, budget, &row.event)?; + } + if request.include_aux { + let row_count = events.len(); + let original_targets = targets; + // A writer retry at an older aux cursor alone would miss newer writer + // edits. Restart both closure hops once after permanent degradation; + // preserve the request ledger/deadline, never reset work allowances. + 'closure: loop { + let mut targets = original_targets.clone(); + let mut seen = HashSet::new(); + for kinds in [&WINDOW_AUX_KINDS[..], &WINDOW_AUX_DELETE_KINDS[..]] { + let mut next_targets = HashSet::new(); + // Bound SQL expression size, including deletion-of-aux fanout. + for batch in targets.chunks(200) { + let mut query = AuxQuery { + community: tenant.community(), + targets: batch, + kinds, + accessible, + cursor: None, + }; + loop { + let was_replica = session.is_replica(); + let page = session + .thread_window_aux(&query, &mut budget.scan) + .await + .map_err(|e| database_error("auxiliary closure", e))?; + if was_replica && !session.is_replica() { + events.truncate(row_count); + continue 'closure; + } + next_targets.extend(page.target_ids); + for event in page.events { + if visible(&event) && seen.insert(event.event.id) { + append(&mut events, budget, &event.event)?; + } + } + let Some(cursor) = page.next_cursor else { + break; + }; + if query.cursor.as_ref().is_some_and(|old| { + cursor.created_at > old.created_at + || (cursor.created_at == old.created_at && cursor.id <= old.id) + }) { + return Err(unavailable("thread auxiliary scan did not advance")); + } + query.cursor = Some(cursor); + } + } + targets = next_targets.into_iter().collect(); + targets.sort_unstable(); + if targets.is_empty() { + break; + } + } + break; + } + } + let tags = [ + [ + "d".to_string(), + request.binding(tenant.host(), &reader.to_hex()), + ], + ["h".to_string(), request.channel.to_string()], + ["e".to_string(), request.root.clone()], + ] + .into_iter() + .map(nostr::Tag::parse) + .collect::, _>>() + .map_err(|e| internal_error(&format!("thread bounds tags: {e}")))?; + let bounds = nostr::EventBuilder::new( + nostr::Kind::Custom(buzz_core::kind::KIND_THREAD_WINDOW_BOUNDS as u16), + json!({"version":1,"direction":"older","has_more":window.has_more, + "next_cursor":window.next_cursor}) + .to_string(), + ) + .tags(tags) + .sign_with_keys(&state.relay_keypair) + .map_err(|e| internal_error(&format!("thread bounds sign: {e}")))?; + append(&mut events, budget, &bounds)?; + metrics::histogram!("buzz_thread_window_response_bytes") + .record((budget.bytes - page_start_bytes) as f64); + Ok(events) +} + +#[cfg(test)] +mod postgres_tests; diff --git a/crates/buzz-relay/src/api/bridge/thread_window/postgres_tests.rs b/crates/buzz-relay/src/api/bridge/thread_window/postgres_tests.rs new file mode 100644 index 00000000000..9cee67ea3bd --- /dev/null +++ b/crates/buzz-relay/src/api/bridge/thread_window/postgres_tests.rs @@ -0,0 +1,534 @@ +use super::*; +use axum::{ + body::{to_bytes, Body}, + http::Request as HttpRequest, +}; +use base64::Engine; +use buzz_core::{ + channel::{ChannelType, ChannelVisibility}, + CommunityId, +}; +use nostr::{Event, EventBuilder, Keys, Kind, Tag, Timestamp}; +use sha2::{Digest, Sha256}; +use std::sync::Arc; +use tower::ServiceExt; +use uuid::Uuid; + +struct Fixture { + state: Arc, + host: String, + community: CommunityId, + channel: Uuid, + keys: Keys, + root: Event, + pool: sqlx::PgPool, +} + +impl Fixture { + async fn new() -> Self { + let state = crate::api::bridge::postgres_tests::bridge_handler_test_state() + .await + .unwrap(); + let mut state = (*state).clone(); + // Use production after_connect policy (floor guard, isolation and + // session timeouts), not raw SQLx pools that mask deployed failures. + state.db = production_db(buzz_db::DbConfig::default()).await; + Arc::make_mut(&mut state.config).require_auth_token = true; + state.nip98_replay = Arc::new(buzz_pubsub::RedisNip98ReplayGuard::new( + state.redis_pool.clone(), + )); + let state = Arc::new(state); + let host = format!("tw-{}.local", Uuid::new_v4()); + let community = state + .db + .ensure_configured_community(&host) + .await + .unwrap() + .id; + let channel = Uuid::new_v4(); + let keys = Keys::generate(); + private_channel(&state.db, community, channel, &keys).await; + let root = event(&keys, channel, 9, "root", None, Timestamp::now().as_secs()); + state + .db + .insert_event(community, &root, Some(channel)) + .await + .unwrap(); + let pool = sqlx::PgPool::connect(&crate::test_support::database_url()) + .await + .unwrap(); + Self { + state, + host, + community, + channel, + keys, + root, + pool, + } + } + fn filter(&self) -> Value { + json!({"thread_window":true,"#h":[self.channel],"#e":[self.root.id.to_hex()], + "kinds":[9],"limit":50,"include_aux":true}) + } + async fn post(&self, key: &Keys, path: &str, body: Value) -> (StatusCode, Value) { + post(self.state.clone(), &self.host, key, path, body).await + } + async fn query(&self, filter: &Value) -> Value { + let (status, body) = self.post(&self.keys, "/query", json!([filter])).await; + assert_eq!(status, StatusCode::OK, "{body}"); + body + } + async fn tombstone(&self, event: &Event) { + sqlx::query("UPDATE events SET deleted_at=now() WHERE community_id=$1 AND id=$2") + .bind(self.community.as_uuid()) + .bind(event.id.to_bytes().to_vec()) + .execute(&self.pool) + .await + .unwrap(); + } + // Bulk budget fixtures need structurally readable rows, not valid signatures. + async fn copy_aux(&self, source: &Event, count: i32, content: &str) { + sqlx::query("INSERT INTO events (community_id,id,pubkey,created_at,kind,tags,content,sig,received_at,channel_id) \ + SELECT community_id,decode(md5(n::text)||md5(('aux'||n)::text),'hex'),pubkey,created_at,kind,tags,$4,sig,received_at,channel_id \ + FROM events CROSS JOIN generate_series(1,$3) n WHERE community_id=$1 AND id=$2") + .bind(self.community.as_uuid()).bind(source.id.to_bytes().to_vec()) + .bind(count).bind(content).execute(&self.pool).await.unwrap(); + } + async fn reply(&self, n: usize) -> Event { + let reply = event( + &self.keys, + self.channel, + 9, + &format!("reply {n}"), + Some(&self.root), + self.root.created_at.as_secs() + n as u64, + ); + let ts = chrono::DateTime::from_timestamp(reply.created_at.as_secs() as i64, 0).unwrap(); + let root_ts = + chrono::DateTime::from_timestamp(self.root.created_at.as_secs() as i64, 0).unwrap(); + self.state + .db + .insert_event_with_thread_metadata( + self.community, + &reply, + Some(self.channel), + Some(buzz_db::event::ThreadMetadataParams { + event_id: &reply.id.to_bytes(), + event_created_at: ts, + channel_id: self.channel, + parent_event_id: Some(&self.root.id.to_bytes()), + parent_event_created_at: Some(root_ts), + root_event_id: Some(&self.root.id.to_bytes()), + root_event_created_at: Some(root_ts), + depth: 1, + broadcast: false, + }), + ) + .await + .unwrap(); + reply + } + async fn aux(&self, kind: u16, target: &Event, channel: Option) -> Event { + let aux = event( + &self.keys, + channel.unwrap_or(self.channel), + kind, + &format!("aux {}", Uuid::new_v4()), + Some(target), + self.root.created_at.as_secs() + 60, + ); + self.state + .db + .insert_event(self.community, &aux, channel) + .await + .unwrap(); + aux + } + fn bounds(&self, response: &Value, filter: &Value) -> Value { + self.bounds_on_host(response, filter, &self.host) + } + fn bounds_on_host(&self, response: &Value, filter: &Value, host: &str) -> Value { + let request = Request::parse(filter).unwrap(); + let bounds = response + .as_array() + .unwrap() + .iter() + .filter(|v| v["kind"] == 39007) + .collect::>(); + assert_eq!(bounds.len(), 1, "{response}"); + let event: Event = serde_json::from_value(bounds[0].clone()).unwrap(); + event.verify().unwrap(); + assert_eq!(event.pubkey, self.state.relay_keypair.public_key()); + let tags: Value = serde_json::to_value(&event.tags).unwrap(); + assert_eq!( + tags, + json!([ + ["d", request.binding(host, &self.keys.public_key().to_hex())], + ["h", request.channel.to_string()], + ["e", request.root] + ]) + ); + let content: Value = serde_json::from_str(&event.content).unwrap(); + assert_eq!(content["version"], 1); + assert_eq!(content["direction"], "older"); + assert_eq!( + content["has_more"].as_bool().unwrap(), + !content["next_cursor"].is_null() + ); + content + } +} + +async fn production_db(mut config: buzz_db::DbConfig) -> buzz_db::Db { + config.database_url = crate::test_support::database_url(); + config.max_connections = 5; + buzz_db::Db::new(&config).await.unwrap() +} + +async fn private_channel(db: &buzz_db::Db, community: CommunityId, channel: Uuid, owner: &Keys) { + db.create_channel_with_id( + community, + channel, + "thread-window", + ChannelType::Stream, + ChannelVisibility::Private, + None, + &owner.public_key().to_bytes(), + None, + ) + .await + .unwrap(); +} + +fn ids(response: &Value, kind: Option) -> Vec<&str> { + response + .as_array() + .unwrap() + .iter() + .filter(|e| kind.is_none_or(|kind| e["kind"] == kind)) + .map(|e| e["id"].as_str().unwrap()) + .collect() +} + +fn event( + keys: &Keys, + channel: Uuid, + kind: u16, + content: &str, + target: Option<&Event>, + ts: u64, +) -> Event { + let mut tags = vec![Tag::parse(["h", &channel.to_string()]).unwrap()]; + if let Some(target) = target { + tags.push(Tag::parse(["e", &target.id.to_hex(), "", "reply"]).unwrap()); + } + EventBuilder::new(Kind::Custom(kind), content) + .tags(tags) + .custom_created_at(Timestamp::from(ts)) + .sign_with_keys(keys) + .unwrap() +} + +async fn post( + state: Arc, + host: &str, + keys: &Keys, + path: &str, + value: Value, +) -> (StatusCode, Value) { + let body = serde_json::to_vec(&value).unwrap(); + let proof = EventBuilder::new(Kind::Custom(27235), "") + .tags([ + Tag::parse(["u", &format!("https://{host}{path}")]).unwrap(), + Tag::parse(["method", "POST"]).unwrap(), + Tag::parse(["payload", &hex::encode(Sha256::digest(&body))]).unwrap(), + Tag::parse(["nonce", &Uuid::new_v4().to_string()]).unwrap(), + ]) + .sign_with_keys(keys) + .unwrap(); + let auth = format!( + "Nostr {}", + base64::engine::general_purpose::STANDARD.encode(serde_json::to_vec(&proof).unwrap()) + ); + let response = crate::router::build_router(state) + .oneshot( + HttpRequest::builder() + .method("POST") + .uri(path) + .header("host", host) + .header("authorization", auth) + .body(Body::from(body)) + .unwrap(), + ) + .await + .unwrap(); + let status = response.status(); + let bytes = to_bytes(response.into_body(), 16 * 1024 * 1024) + .await + .unwrap(); + (status, serde_json::from_slice(&bytes).unwrap()) +} + +#[tokio::test] +#[ignore = "requires Postgres"] +async fn thread_window_real_query_signed_bounds_and_sentinel_aux() { + let f = Fixture::new().await; + let filter = f.filter(); + let empty = f.query(&filter).await; + assert_eq!(f.bounds(&empty, &filter)["has_more"], false); + let mut replies = Vec::new(); + for n in 0..51 { + replies.push(f.reply(n).await); + } + let sentinel_aux = f.aux(7, &replies[0], Some(f.channel)).await; + let reaction = f.aux(7, &replies[50], Some(f.channel)).await; + let deletion = f.aux(5, &reaction, None).await; + f.tombstone(&reaction).await; + let root_edit = f.aux(40003, &f.root, Some(f.channel)).await; + let page = f.query(&filter).await; + let rows = ids(&page, Some(9)); + assert_eq!(rows.len(), 50); + assert_eq!(rows[0], replies[50].id.to_hex()); + assert_eq!(rows[49], replies[1].id.to_hex()); + let all = ids(&page, None); + for (event, present) in [ + (&sentinel_aux, false), + (&reaction, false), + (&deletion, true), + (&root_edit, true), + ] { + assert_eq!(all.contains(&event.id.to_hex().as_str()), present); + } + let bounds = f.bounds(&page, &filter); + assert_eq!(bounds["has_more"], true); + assert_eq!(bounds["next_cursor"]["id"], replies[1].id.to_hex()); + let mut next = filter.clone(); + next["until"] = bounds["next_cursor"]["created_at"].clone(); + next["before_id"] = bounds["next_cursor"]["id"].clone(); + let tail = f.query(&next).await; + assert_eq!(f.bounds(&tail, &next)["next_cursor"], Value::Null); + assert_eq!(ids(&tail, Some(9)), [replies[0].id.to_hex()]); + // Actual mobile sentinel + unbounded depth remain valid ONLY in legacy. + // Both cursor spellings and absent/false opt-in preserve ASC/ASC and no bounds. + for flag in [None, Some(false)] { + for (ts_key, id_key) in [ + ("thread_cursor", "thread_cursor_id"), + ("threadCursor", "threadCursorId"), + ] { + let mut legacy = json!({"#h":[f.channel],"#e":[f.root.id.to_hex()],"kinds":[9], + "depth_limit":2147483647,"limit":2}); + legacy[ts_key] = json!(-1); + if let Some(flag) = flag { + legacy["thread_window"] = json!(flag); + } + let old = f.query(&legacy).await; + assert_eq!(old[0]["id"], replies[0].id.to_hex()); + assert_eq!(old[1]["id"], replies[1].id.to_hex()); + assert_eq!(old.as_array().unwrap().len(), 2); + legacy[ts_key] = json!(replies[1].created_at.as_secs()); + legacy[id_key] = json!(replies[1].id.to_hex()); + let next = f.query(&legacy).await; + assert_eq!(next[0]["id"], replies[2].id.to_hex()); + assert_eq!(next[1]["id"], replies[3].id.to_hex()); + assert_eq!(next.as_array().unwrap().len(), 2); + } + } +} + +#[tokio::test] +#[ignore = "requires Postgres"] +async fn thread_window_real_query_denial_revocation_and_colliding_tenants() { + let f = Fixture::new().await; + f.reply(0).await; + let filter = f.filter(); + let outsider = Keys::generate(); + let (status, denied) = f.post(&outsider, "/query", json!([filter])).await; + assert_eq!(status, StatusCode::OK); + assert_eq!(denied, json!([])); + let other_host = format!("other-{}.local", Uuid::new_v4()); + let other = f + .state + .db + .ensure_configured_community(&other_host) + .await + .unwrap() + .id; + private_channel(&f.state.db, other, f.channel, &f.keys).await; + // Same channel ID, same accessible reader, absent root in the other tenant. + let (status, other_page) = post( + f.state.clone(), + &other_host, + &f.keys, + "/query", + json!([filter]), + ) + .await; + assert_eq!(status, StatusCode::OK); + assert_eq!( + other_page, + json!([]), + "absent root must not sign exhaustion" + ); + f.query(&filter).await; + // Prime the usual cached access, then revoke directly on the writer without + // cache invalidation (the shape of cross-node delayed invalidation). + assert!(f + .state + .get_accessible_channel_ids_cached(f.community, &f.keys.public_key().to_bytes()) + .await + .unwrap() + .contains(&f.channel)); + sqlx::query( + "UPDATE channel_members SET removed_at=now() WHERE community_id=$1 AND channel_id=$2", + ) + .bind(f.community.as_uuid()) + .bind(f.channel) + .execute(&f.pool) + .await + .unwrap(); + assert!( + f.state + .get_accessible_channel_ids_cached(f.community, &f.keys.public_key().to_bytes()) + .await + .unwrap() + .contains(&f.channel), + "control: cached access must still be stale" + ); + let (status, revoked) = f.post(&f.keys, "/query", json!([filter])).await; + assert_eq!(status, StatusCode::OK); + assert_eq!(revoked, json!([])); +} + +#[tokio::test] +#[ignore = "requires Postgres"] +async fn thread_window_real_query_validation_forgery_and_sql_failure() { + let f = Fixture::new().await; + for (key, val) in [ + ("until", json!(1)), + ("before_id", json!("z".repeat(64))), + ("top_level", json!(true)), + ("thread_cursor", json!(1)), + ("depth_limit", json!(0)), + ("authors", json!([f.keys.public_key().to_hex()])), + ("page", json!(1)), + ] { + let mut filter = f.filter(); + filter[key] = val; + let (status, body) = f.post(&f.keys, "/query", json!([filter])).await; + assert_eq!(status, StatusCode::BAD_REQUEST, "{body}"); + } + let forged = event( + &f.keys, + f.channel, + 39007, + "{}", + None, + Timestamp::now().as_secs(), + ); + let (status, body) = f + .post(&f.keys, "/events", serde_json::to_value(forged).unwrap()) + .await; + assert!( + !status.is_success() || body["accepted"] == false, + "{status}: {body}" + ); + assert!(body.to_string().contains("relay-only"), "{body}"); + // A required SQL read blocked by DDL must error, not sign empty bounds. + let mut lock = f.pool.begin().await.unwrap(); + sqlx::query("LOCK TABLE thread_metadata IN ACCESS EXCLUSIVE MODE") + .execute(&mut *lock) + .await + .unwrap(); + let (status, body) = f.post(&f.keys, "/query", json!([f.filter()])).await; + assert!(status.is_server_error(), "{status}: {body}"); + lock.rollback().await.unwrap(); +} + +#[tokio::test] +#[ignore = "requires Postgres"] +async fn thread_window_scope_roots_and_aggregate_budgets() { + let f = Fixture::new().await; + let filter = f.filter(); + let (status, body) = f + .post(&f.keys, "/query", json!(vec![filter.clone(); 5])) + .await; + assert_eq!(status, StatusCode::BAD_REQUEST, "{body}"); + for other in [ + json!({"kinds":[20001]}), + json!({"kinds":[9],"search":"x"}), + json!({"kinds":[9],"thread_cursor":-1,"depth_limit":2147483647}), + ] { + let (status, body) = f.post(&f.keys, "/query", json!([filter, other])).await; + assert_eq!(status, StatusCode::BAD_REQUEST, "{body}"); + } + // Neither a private non-conversation root nor a root in another channel + // may authorize aux, even when the aux itself is in the requested channel. + let other_channel = Uuid::new_v4(); + private_channel(&f.state.db, f.community, other_channel, &Keys::generate()).await; + for (channel, kind) in [(f.channel, 30300), (other_channel, 9)] { + let root = event( + &f.keys, + channel, + kind, + "hidden root", + None, + f.root.created_at.as_secs(), + ); + f.state + .db + .insert_event(f.community, &root, Some(channel)) + .await + .unwrap(); + f.aux(40003, &root, Some(f.channel)).await; + let mut hidden_filter = filter.clone(); + hidden_filter["#e"] = json!([root.id.to_hex()]); + let body = f.query(&hidden_filter).await; + assert_eq!( + body, + json!([]), + "hidden roots return neither aux nor bounds" + ); + } + + // Real bounded payloads: one 4.8 MiB page passes, two in the same request + // exceed 8 MiB. A per-window (instead of per-query) ledger fails this test. + let aux = f.aux(40003, &f.root, Some(f.channel)).await; + f.copy_aux(&aux, 80, &"x".repeat(60_000)).await; + f.query(&filter).await; + let (status, two) = f.post(&f.keys, "/query", json!([filter, filter])).await; + assert_eq!(status, StatusCode::SERVICE_UNAVAILABLE, "{two}"); + assert!(two.to_string().contains("byte budget")); +} + +#[tokio::test] +#[ignore = "requires Postgres"] +async fn thread_window_router_aux_row_cap_and_corrupt_page() { + let f = Fixture::new().await; + let aux = f.aux(40003, &f.root, Some(f.channel)).await; + for n in 0..51 { + f.reply(n).await; + } + f.copy_aux(&aux, 8200, "fixture").await; + let (status, body) = f.post(&f.keys, "/query", json!([f.filter()])).await; + assert_eq!(status, StatusCode::SERVICE_UNAVAILABLE, "{body}"); + assert!(body["error"].as_str().unwrap().contains("raw row budget")); + // Keep 1,001 rows and corrupt one guaranteed to lie on the first raw page. + sqlx::query("DELETE FROM events WHERE community_id=$1 AND kind=40003 AND id NOT IN \ + (SELECT id FROM events WHERE community_id=$1 AND kind=40003 ORDER BY created_at DESC,id ASC LIMIT 1001)") + .bind(f.community.as_uuid()).execute(&f.pool).await.unwrap(); + // Positive control: the same 1,001 raw rows succeed before corruption. + let complete = f.query(&f.filter()).await; + assert_eq!(complete.as_array().unwrap().len(), 1052); + assert_eq!(f.bounds(&complete, &f.filter())["has_more"], true); + sqlx::query("UPDATE events SET sig='\\x00' WHERE community_id=$1 AND id = \ + (SELECT id FROM events WHERE community_id=$1 AND kind=40003 ORDER BY created_at DESC,id ASC OFFSET 500 LIMIT 1)") + .bind(f.community.as_uuid()).execute(&f.pool).await.unwrap(); + let (status, body) = f.post(&f.keys, "/query", json!([f.filter()])).await; + assert!(status.is_server_error(), "{status}: {body}"); + assert_eq!(body, json!({"error":"internal server error"})); +} + +mod failure_postgres_tests; + +mod review_postgres_tests; diff --git a/crates/buzz-relay/src/api/bridge/thread_window/postgres_tests/failure_postgres_tests.rs b/crates/buzz-relay/src/api/bridge/thread_window/postgres_tests/failure_postgres_tests.rs new file mode 100644 index 00000000000..c7f7e7bf925 --- /dev/null +++ b/crates/buzz-relay/src/api/bridge/thread_window/postgres_tests/failure_postgres_tests.rs @@ -0,0 +1,347 @@ +use super::*; +use std::{ + future::Future, + sync::atomic::{AtomicUsize, Ordering}, + task::Poll, +}; +use tracing::instrument::WithSubscriber; +use tracing_subscriber::{layer::Context, prelude::*, registry::LookupSpan, Layer}; + +// Pause the HTTP future after the first real auxiliary page completes. This +// observes a production span rather than replacing the database or closure. +struct AuxPages(Arc); +impl Layer for AuxPages +where + S: tracing::Subscriber + for<'a> LookupSpan<'a>, +{ + fn on_close(&self, id: tracing::Id, ctx: Context<'_, S>) { + if ctx.span(&id).unwrap().metadata().name() == "thread_window_aux" { + self.0.fetch_add(1, Ordering::SeqCst); + } + } +} + +async fn pause_after_aux_page( + request: F, +) -> std::pin::Pin>> { + let pages = Arc::new(AtomicUsize::new(0)); + let subscriber = tracing_subscriber::registry().with(AuxPages(pages.clone())); + let mut request = Box::pin(request.with_subscriber(subscriber)); + tokio::time::timeout( + Duration::from_secs(5), + std::future::poll_fn(|cx| { + assert!( + request.as_mut().poll(cx).is_pending(), + "must pause inside closure" + ); + if pages.load(Ordering::SeqCst) > 0 { + Poll::Ready(()) + } else { + Poll::Pending + } + }), + ) + .await + .expect("first auxiliary page must complete before transition"); + assert_eq!( + pages.load(Ordering::SeqCst), + 1, + "barrier must precede closure completion" + ); + request +} + +#[tokio::test] +#[ignore = "requires Postgres"] +async fn thread_window_bridge_restarts_both_aux_hops_after_replica_failure() { + let mut f = Fixture::new().await; + let reply = f.reply(0).await; + let reaction = f.aux(7, &reply, Some(f.channel)).await; + // Two raw pages, so losing the snapshot after page one leaves a meaningful + // old cursor. The writer edit inserted later sorts before that cursor. + f.copy_aux(&reaction, 1000, "fixture").await; + let reader_name = format!("tw-reader-{}", Uuid::new_v4()); + let options: sqlx::postgres::PgConnectOptions = + crate::test_support::database_url().parse().unwrap(); + let reader = sqlx::postgres::PgPoolOptions::new() + .max_connections(1) + .connect_with(options.application_name(&reader_name)) + .await + .unwrap(); + let db = buzz_db::Db::from_pools(f.pool.clone(), reader.clone()); + db.fence() + .force_open_for_tests(chrono::Utc::now() + chrono::Duration::seconds(10)); + Arc::make_mut(&mut f.state).db = db; + let mut filter = f.filter(); + filter["until"] = json!(f.root.created_at.as_secs() + 1); + filter["before_id"] = json!("00".repeat(32)); + let request = pause_after_aux_page(f.post(&f.keys, "/query", json!([filter]))).await; + let held: i64 = sqlx::query_scalar("SELECT count(*) FROM pg_stat_activity WHERE application_name=$1 AND xact_start IS NOT NULL") + .bind(&reader_name).fetch_one(&f.pool).await.unwrap(); + assert_eq!(held, 1, "must actually hold the proved replica transaction"); + let edit = f.aux(40003, &reply, Some(f.channel)).await; + let deletion = f.aux(5, &edit, None).await; + f.tombstone(&edit).await; + sqlx::query("SELECT pg_terminate_backend(pid) FROM pg_stat_activity WHERE application_name=$1") + .bind(&reader_name) + .execute(&f.pool) + .await + .unwrap(); + let (status, body) = request.await; + assert_eq!(status, StatusCode::OK, "{body}"); + assert_eq!(f.bounds(&body, &filter)["has_more"], false); + let ids = ids(&body, None); + assert!( + ids.contains(&deletion.id.to_hex().as_str()), + "restart must discover writer edit tombstone and its deletion" + ); + assert!(!ids.contains(&edit.id.to_hex().as_str())); + assert_eq!( + ids.len(), + ids.iter().collect::>().len(), + "discarded replica output must not duplicate events" + ); + assert_eq!(ids.len(), 1004); // reply, 1001 reactions, deletion, bounds + reader.close().await; +} + +#[tokio::test] +#[ignore = "requires Postgres"] +async fn thread_window_http_deadline_covers_authorization_wait() { + let mut f = Fixture::new().await; + // The shared HTTP deadline still applies when the optional DB lock budget is disabled. + Arc::make_mut(&mut f.state).db = production_db(buzz_db::DbConfig { + lock_timeout_ms: 0, + ..Default::default() + }) + .await; + assert_authorization_timeout( + &f, + "thread window deadline exceeded", + DEADLINE, + DEADLINE + Duration::from_secs(4), + ) + .await; +} + +#[tokio::test] +#[ignore = "requires Postgres"] +async fn thread_window_bounds_rejected_by_ws_event_handler() { + let f = Fixture::new().await; + let forged = event( + &f.keys, + f.channel, + 39007, + "{}", + None, + Timestamp::now().as_secs(), + ); + let auth = crate::connection::AuthState::Authenticated(buzz_auth::AuthContext { + pubkey: f.keys.public_key(), + scopes: vec![], + channel_ids: None, + auth_method: buzz_auth::AuthMethod::Nip42, + agent_owner_pubkey: None, + }); + let (mut conn, mut send_rx) = crate::connection::tests::test_conn_with_auth(auth); + Arc::get_mut(&mut conn).unwrap().tenant = + buzz_core::TenantContext::resolved(f.community, f.host.clone()); + crate::handlers::event::handle_event(forged.clone(), conn, f.state.clone()).await; + let axum::extract::ws::Message::Text(text) = send_rx.try_recv().unwrap() else { + panic!("expected ACK") + }; + let ack: Value = serde_json::from_str(&text).unwrap(); + assert_eq!(ack[0], "OK"); + assert_eq!(ack[1], forged.id.to_hex()); + assert_eq!(ack[2], false); + assert!(ack[3].as_str().unwrap().contains("relay-only"), "{ack}"); +} + +#[tokio::test] +#[ignore = "requires Postgres"] +async fn thread_window_retries_aux_access_grant_during_closure() { + assert_aux_access_change(true).await; +} + +#[tokio::test] +#[ignore = "requires Postgres"] +async fn thread_window_retries_aux_access_revocation_during_closure() { + assert_aux_access_change(false).await; +} + +async fn assert_aux_access_change(grant: bool) { + let f = Fixture::new().await; + let reply = f.reply(0).await; + // A visible first-hop event ensures the closure has a second hop where + // polling can pause, even when the cross-channel edit is initially hidden. + f.aux(7, &reply, Some(f.channel)).await; + let aux_channel = Uuid::new_v4(); + private_channel(&f.state.db, f.community, aux_channel, &f.keys).await; + let edit = f.aux(40003, &reply, Some(aux_channel)).await; + set_access(&f, aux_channel, !grant).await; + let contains_edit = |body: &Value| { + body.as_array() + .unwrap() + .iter() + .any(|e| e["id"] == edit.id.to_hex()) + }; + let filter = f.filter(); + let before = f.query(&filter).await; + f.bounds(&before, &filter); + assert_eq!( + contains_edit(&before), + !grant, + "pre-transition visibility control" + ); + + let request = pause_after_aux_page(f.post(&f.keys, "/query", json!([filter]))).await; + set_access(&f, aux_channel, grant).await; + let current = f + .state + .db + .get_accessible_channel_ids(f.community, &f.keys.public_key().to_bytes()) + .await + .unwrap(); + assert!( + current.contains(&f.channel), + "requested channel stays authorized" + ); + assert_eq!(current.contains(&aux_channel), grant); + let (status, interrupted) = request.await; + assert_eq!(status, StatusCode::SERVICE_UNAVAILABLE, "{interrupted}"); + assert_eq!( + interrupted, + json!({"error":"thread authorization changed; retry query"}) + ); + + let after = f.query(&filter).await; + f.bounds(&after, &filter); + assert_eq!( + contains_edit(&after), + grant, + "retry must use the complete new access set" + ); +} + +#[tokio::test] +#[ignore = "requires Postgres"] +async fn thread_window_production_authorization_lock_timeout_is_retryable() { + let f = Fixture::new().await; + assert_authorization_timeout( + &f, + "thread database timeout; retry window", + Duration::from_millis(buzz_db::DbConfig::default().lock_timeout_ms), + DEADLINE, + ) + .await; +} + +async fn assert_authorization_timeout(f: &Fixture, message: &str, min: Duration, max: Duration) { + let mut lock = f.pool.begin().await.unwrap(); + sqlx::query("LOCK TABLE channel_members IN ACCESS EXCLUSIVE MODE") + .execute(&mut *lock) + .await + .unwrap(); + let started = std::time::Instant::now(); + let (status, body) = f.post(&f.keys, "/query", json!([f.filter()])).await; + assert_eq!(status, StatusCode::SERVICE_UNAVAILABLE, "{body}"); + assert_eq!(body, json!({"error":message})); + assert!((min..max).contains(&started.elapsed())); + lock.rollback().await.unwrap(); + let recovered = f.query(&f.filter()).await; + assert_eq!(f.bounds(&recovered, &f.filter())["has_more"], false); +} + +async fn set_access(f: &Fixture, channel: Uuid, allowed: bool) { + let changed = sqlx::query( + "UPDATE channel_members SET removed_at=CASE WHEN $4 THEN NULL ELSE now() END \ + WHERE community_id=$1 AND channel_id=$2 AND pubkey=$3", + ) + .bind(f.community.as_uuid()) + .bind(channel) + .bind(f.keys.public_key().to_bytes().to_vec()) + .bind(allowed) + .execute(&f.pool) + .await + .unwrap(); + assert_eq!(changed.rows_affected(), 1); +} + +#[tokio::test] +#[ignore = "requires Postgres"] +async fn thread_window_database_timeout_classification_uses_sqlstate() { + let f = Fixture::new().await; + // The adapter handles the same DB failures at initial/final authorization, + // selection and aux closure. Exercise actual PostgreSQL statement errors, + // not string-matched synthetic errors; unrelated faults stay sanitized 500s. + for (sql, expected) in [ + ( + "SET statement_timeout='25ms'; SELECT pg_sleep(1)", + StatusCode::SERVICE_UNAVAILABLE, + ), + ("SELECT 1/0", StatusCode::INTERNAL_SERVER_ERROR), + ] { + let mut conn = f.pool.acquire().await.unwrap(); + let error = sqlx::raw_sql(sql).execute(&mut *conn).await.unwrap_err(); + let (status, body) = database_error("test", error.into()); + assert_eq!(status, expected, "{body:?}"); + } + let (status, body) = database_error("pool", buzz_db::DbError::Sqlx(sqlx::Error::PoolTimedOut)); + assert_eq!(status, StatusCode::SERVICE_UNAVAILABLE); + assert_eq!( + body.0, + json!({"error":"thread database timeout; retry window"}) + ); +} + +#[tokio::test] +#[ignore = "requires Postgres"] +async fn thread_window_batch_discards_earlier_windows_after_revocation() { + // Cover both repeated-root batches and revoking an earlier window's channel + // while the later window remains authorized. All output must fail closed. + for same_channel in [true, false] { + let f = Fixture::new().await; + f.reply(0).await; + let mut first = f.filter(); + first["include_aux"] = json!(false); + let mut later = f.filter(); + if !same_channel { + let channel = Uuid::new_v4(); + private_channel(&f.state.db, f.community, channel, &f.keys).await; + let root = event( + &f.keys, + channel, + 9, + "other root", + None, + f.root.created_at.as_secs(), + ); + f.state + .db + .insert_event(f.community, &root, Some(channel)) + .await + .unwrap(); + f.aux(7, &root, Some(channel)).await; + later["#h"] = json!([channel]); + later["#e"] = json!([root.id.to_hex()]); + } else { + f.aux(7, &f.root, Some(f.channel)).await; + } + let batch = json!([first, later]); + let (status, before) = f.post(&f.keys, "/query", batch.clone()).await; + assert_eq!(status, StatusCode::OK, "{before}"); + assert_eq!(ids(&before, Some(39007)).len(), 2); + let request = pause_after_aux_page(f.post(&f.keys, "/query", batch.clone())).await; + set_access(&f, f.channel, false).await; + let (status, body) = request.await; + assert_eq!(status, StatusCode::SERVICE_UNAVAILABLE, "{body}"); + assert_eq!( + body, + json!({"error":"thread authorization changed; retry query"}) + ); + let (status, after) = f.post(&f.keys, "/query", batch).await; + assert_eq!(status, StatusCode::OK, "{after}"); + assert_eq!(ids(&after, Some(9)).len(), 0); + assert_eq!(ids(&after, Some(39007)).len(), usize::from(!same_channel)); + } +} diff --git a/crates/buzz-relay/src/api/bridge/thread_window/postgres_tests/review_postgres_tests.rs b/crates/buzz-relay/src/api/bridge/thread_window/postgres_tests/review_postgres_tests.rs new file mode 100644 index 00000000000..ba4e4e14370 --- /dev/null +++ b/crates/buzz-relay/src/api/bridge/thread_window/postgres_tests/review_postgres_tests.rs @@ -0,0 +1,162 @@ +use super::*; + +async fn ingest(f: &Fixture, event: &Event) { + let (status, body) = f + .post(&f.keys, "/events", serde_json::to_value(event).unwrap()) + .await; + assert_eq!(status, StatusCode::OK, "{body}"); + assert_eq!(body["accepted"], true, "{body}"); +} + +#[tokio::test] +#[ignore = "requires Postgres"] +async fn thread_window_signed_aux_requires_positional_target() { + let f = Fixture::new().await; + let other = event( + &f.keys, + f.channel, + 9, + "other root", + None, + f.root.created_at.as_secs(), + ); + ingest(&f, &other).await; + let a = f.root.id.to_hex(); + let b = other.id.to_hex(); + let mut edits = Vec::new(); + for (n, tags) in [ + vec![vec!["e", b.as_str(), a.as_str()]], + vec![vec!["e", b.as_str()], vec!["x", "e", a.as_str()]], + vec![vec!["e", a.as_str()]], + vec![vec!["e", b.as_str()], vec!["e", a.as_str()]], + ] + .into_iter() + .enumerate() + { + let edit = EventBuilder::new(Kind::Custom(40003), format!("edit {n}")) + .tags( + std::iter::once(Tag::parse(["h", &f.channel.to_string()]).unwrap()) + .chain(tags.into_iter().map(|tag| Tag::parse(tag).unwrap())), + ) + .sign_with_keys(&f.keys) + .unwrap(); + ingest(&f, &edit).await; + edits.push(edit); + } + let filter = f.filter(); + let page = f.query(&filter).await; + assert_eq!(f.bounds(&page, &filter)["has_more"], false); + let aux = ids(&page, Some(40003)); + for (n, edit) in edits.iter().enumerate() { + assert_eq!(aux.contains(&edit.id.to_hex().as_str()), n >= 2, "case {n}"); + } + // Both non-target shapes really were stored and remain reachable under B. + let mut other_filter = filter; + other_filter["#e"] = json!([b]); + assert_eq!(ids(&f.query(&other_filter).await, Some(40003)).len(), 3); +} + +#[tokio::test] +#[ignore = "requires Postgres"] +async fn thread_window_unsupported_signed_root_has_no_false_bounds() { + let f = Fixture::new().await; + let root = EventBuilder::new(Kind::Custom(40008), "diff --git a/file b/file") + .tags([ + Tag::parse(["h", &f.channel.to_string()]).unwrap(), + Tag::parse(["repo", "https://github.com/block/buzz"]).unwrap(), + Tag::parse(["commit", "1234567"]).unwrap(), + ]) + .sign_with_keys(&f.keys) + .unwrap(); + ingest(&f, &root).await; + let reply = event( + &f.keys, + f.channel, + 9, + "reply to diff", + Some(&root), + root.created_at.as_secs(), + ); + ingest(&f, &reply).await; + let legacy = json!({"#h":[f.channel],"#e":[root.id.to_hex()],"kinds":[9],"depth_limit":100}); + assert_eq!(ids(&f.query(&legacy).await, Some(9)), [reply.id.to_hex()]); + let mut filter = f.filter(); + filter["#e"] = json!([root.id.to_hex()]); + for include_aux in [false, true] { + filter["include_aux"] = json!(include_aux); + assert_eq!( + f.query(&filter).await, + json!([]), + "unsupported is not exhausted" + ); + } +} + +#[tokio::test] +#[ignore = "requires Postgres"] +async fn thread_window_signed_large_aux_fails_with_recoverable_budget_error() { + let f = Fixture::new().await; + // No forged bulk copies: prove these maximum-sized payloads pass ingest. + for n in 0..40 { + let edit = event( + &f.keys, + f.channel, + 40003, + &"x".repeat(256 * 1024), + Some(&f.root), + f.root.created_at.as_secs() + n, + ); + ingest(&f, &edit).await; + } + let filter = f.filter(); + let (status, body) = f.post(&f.keys, "/query", json!([filter])).await; + assert_eq!(status, StatusCode::SERVICE_UNAVAILABLE, "{body}"); + assert!(body["error"] + .as_str() + .unwrap() + .contains("payload byte budget")); + let mut without_aux = filter; + without_aux["include_aux"] = json!(false); + let recovered = f.query(&without_aux).await; + assert_eq!(f.bounds(&recovered, &without_aux)["has_more"], false); +} + +#[tokio::test] +#[ignore = "requires Postgres"] +async fn thread_window_signed_large_replies_without_aux_are_bounded() { + let f = Fixture::new().await; + for n in 0..201 { + let reply = event( + &f.keys, + f.channel, + 9, + &"x".repeat(256 * 1024), + Some(&f.root), + f.root.created_at.as_secs() + n, + ); + ingest(&f, &reply).await; + } + let mut filter = f.filter(); + filter["include_aux"] = json!(false); + filter["limit"] = json!(200); + let (status, body) = f.post(&f.keys, "/query", json!([filter])).await; + assert_eq!(status, StatusCode::SERVICE_UNAVAILABLE, "{body}"); + assert!(body["error"] + .as_str() + .unwrap() + .contains("payload byte budget")); + assert!( + body.as_array().is_none(), + "no partial rows or signed bounds on failure" + ); + filter["limit"] = json!(20); + let page = f.query(&filter).await; + assert_eq!(ids(&page, Some(9)).len(), 20); + assert_eq!(f.bounds(&page, &filter)["has_more"], true); + let (status, body) = f.post(&f.keys, "/query", json!([filter, filter])).await; + assert_eq!(status, StatusCode::SERVICE_UNAVAILABLE, "{body}"); + assert!(body["error"] + .as_str() + .unwrap() + .contains("payload byte budget")); +} diff --git a/docs/nips/NIP-CW.md b/docs/nips/NIP-CW.md index c279cb156ef..9e6ea14e1bd 100644 --- a/docs/nips/NIP-CW.md +++ b/docs/nips/NIP-CW.md @@ -1,8 +1,8 @@ NIP-CW ====== -Channel Window --------------- +Channel and Thread Windows +-------------------------- `draft` `optional` `relay` @@ -10,11 +10,18 @@ Channel Window ## Abstract -This NIP defines the **channel window**: a relay-computed, cursor-paged view of a channel's *top-level* timeline, served as ordinary signed Nostr events through an extended NIP-01 filter. One request returns a page of top-level rows in stable keyset order, optionally accompanied by the aux closure and two relay-signed overlay families: +This NIP defines two relay-computed, cursor-paged views served as ordinary signed Nostr events through extended NIP-01 filters: + +- **channel mode** pages a channel's top-level timeline; +- **thread mode** pages one root's descendants newest-first. + +Both modes use stable keyset order and explicit relay-signed exhaustion bounds. A channel-mode request may also include the aux closure and thread-summary overlays: - the **aux closure** — stored reactions, deletions, and edits targeting the returned rows, with their original authors and signatures (`include_aux`), - **thread summaries** — one relay-signed `kind:39005` per row that has replies (`include_summaries`), -- **window bounds** — exactly one relay-signed `kind:39006` carrying the authoritative `has_more` fact and the next-page cursor. +- **channel bounds** — exactly one relay-signed `kind:39006` carrying the authoritative `has_more` fact and the next-page cursor. + +Thread mode returns reply rows, optional bounded aux closure, and exactly one request-bound `kind:39007` thread-bounds overlay. `39006` and `39007` are deliberately distinct: the existing channel/cursor identity of `39006` cannot disambiguate concurrent roots and MUST NOT be reinterpreted. The extension adds no endpoint and no envelope. The wire format is the flat array of signed events the query surface already returns; a client that ignores this NIP receives standard behavior everywhere. @@ -30,7 +37,7 @@ A relay that computes thread structure at ingest already knows which events are This NIP does not change ingest, storage, or fan-out. Rows returned in a window are ordinary stored events; the overlays are computed per query and never stored. -This NIP does not define thread *reading*. Replies never appear as window rows; fetching a thread's contents is out of scope. +This NIP does not define around-target retrieval or cross-page snapshot isolation. It defines protocol compatibility and fallback rules, but the Buzz thread-mode implementation ships no client opt-in or fallback implementation. Thread mode starts at the newest reply and continues toward older replies. This NIP does not require WebSocket REQ support. A relay MAY serve window filters only on an HTTP query surface and ignore the extension fields on REQ (see §Degradation). @@ -40,14 +47,14 @@ This document uses MUST, MUST NOT, SHOULD, MAY, and RECOMMENDED as defined in RF - **relay identity**: The keypair whose pubkey the relay advertises (e.g. NIP-11 `self`). All overlay events are signed with it. - **row**: A stored, signed event returned as part of the page proper (usually client-authored; Buzz also stores relay-signed events carrying actor provenance). Rows are the only events that count against `limit`. -- **top-level**: An event that opens a thread rather than replying into one — defined by wire tags in §Top-level Classification. -- **overlay**: A relay-signed event (`kind:39005`, `kind:39006`) synthesized at query time. Overlays are metadata *about* rows: never a row, never a cursor input, never durable history. +- **top-level**: An event that opens a thread rather than replying into one — defined by wire tags in §Channel-mode Top-level Classification. +- **overlay**: A relay-signed event (`kind:39005`, `kind:39006`, or `kind:39007`) synthesized at query time. Overlays are metadata *about* rows: never a row, never a cursor input, never durable history. - **composite cursor**: The pair `(created_at, id)` identifying a position in the total order. `created_at` is unix seconds; `id` is a 64-character lowercase hex event id. -- **scan position**: The composite cursor of the last event the relay's query *retained*, whether or not that event was ultimately delivered as a row (see §Relay Processing step 3). The cursor tracks where the scan stopped, not what the client received. +- **scan position**: The composite cursor of the last event the relay's query *retained*, whether or not that event was ultimately delivered as a row (see §Channel-mode Relay Processing Algorithm step 3). The cursor tracks where the scan stopped, not what the client received. -## Request +## Channel-mode Request -A window request is a standard filter plus extension fields, submitted wherever the relay accepts filters (for Buzz: the NIP-98-authenticated HTTP bridge `POST /query`): +A channel-mode window request is a standard filter plus extension fields, submitted wherever the relay accepts filters (for Buzz: the NIP-98-authenticated HTTP bridge `POST /query`): ```jsonc { @@ -63,7 +70,7 @@ A window request is a standard filter plus extension fields, submitted wherever ``` - `top_level` — MUST be boolean `true` to select the window path. Any other value (absent, `false`, string, number) means the filter is served as a normal filter. -- `#h` — the window MUST target exactly one channel. Zero or multiple channels: reject with an error (Buzz: HTTP `400`). A channel the requester cannot access is handled by §Access Scoping, not by an error that confirms the channel exists. +- `#h` — the window MUST target exactly one channel. Zero or multiple channels: reject with an error (Buzz: HTTP `400`). A channel the requester cannot access is handled by §Channel-mode Access Scoping, not by an error that confirms the channel exists. - `limit` — the row budget. Overlays and aux events MUST NOT count against it. Relays SHOULD clamp it to a documented range (Buzz: default 50, maximum 200, minimum 1). - `until` + `before_id` — the request cursor: the `next_cursor` from the previous page's `kind:39006` overlay, echoed verbatim — `until` = `next_cursor.created_at`, `before_id` = `next_cursor.id`. **Both present or both absent.** Exactly one present MUST be rejected: a timestamp-only cursor silently loses or duplicates same-second rows, which is the failure mode this NIP exists to remove. Both absent = head-of-channel request. - `kinds` — optional; restricts which kinds may be rows. It does not affect overlay or aux kinds. @@ -72,7 +79,7 @@ Cursor grammar: `until` MUST be a non-negative integer of unix seconds represent Offset/page-number pagination MUST NOT be honored on the window path. -## Top-level Classification +## Channel-mode Top-level Classification The row set must be reproducible from wire data alone, so the reply/top-level distinction is defined by tags, not by any relay's storage schema. @@ -87,11 +94,11 @@ An event is **top-level** — eligible to be a window row — iff its depth is 0 Storage fallback (fail-open): a relay that indexes this classification at ingest may hold events stored before the index existed, whose depth is unknown. Such events MUST be treated as top-level rather than vanishing from every window. This is a compatibility rule for pre-index data, not a third protocol state — an interoperating implementation classifying from tags alone has no unknown case. -## Relay Processing Algorithm +## Channel-mode Relay Processing Algorithm -For a valid window filter on an accessible channel (§Access Scoping) the relay MUST: +For a valid window filter on an accessible channel (§Channel-mode Access Scoping) the relay MUST: -1. **Select rows.** From the target channel, take events that are top-level (§Top-level Classification), not deleted, and matching `kinds` if present, in the total order `created_at DESC, id ASC` (`id` compared bytewise). With a cursor `(ts, id)`, retain only events where `created_at < ts OR (created_at = ts AND id > id)`. +1. **Select rows.** From the target channel, take events that are top-level (§Channel-mode Top-level Classification), not deleted, and matching `kinds` if present, in the total order `created_at DESC, id ASC` (`id` compared bytewise). With a cursor `(ts, id)`, retain only events where `created_at < ts OR (created_at = ts AND id > id)`. 2. **Probe exhaustion.** Evaluate the query with an internal budget of `limit + 1` rows *after all predicates*. If `limit + 1` rows match, `has_more = true` and the sentinel row is discarded — it MUST NOT appear on the wire, in overlays, or in the aux closure. Otherwise `has_more = false`. 3. **Derive the next cursor.** If `has_more`, `next_cursor` is the **scan position**: the composite cursor of the last retained candidate, captured *before* any serving-time reconstruction or filtering of individual events. Otherwise `next_cursor = null`. The invariant `next_cursor = null ⇔ has_more = false` MUST hold. Because it is a scan position, `next_cursor` MAY reference an event that does not appear in the response (e.g. one skipped by the relay as unreconstructable); it is authoritative regardless, and deriving it from delivered rows instead would stall pagination on every skipped event. 4. **Append the aux closure** (if `include_aux` and at least one row): two hops of events referencing the rows by `e` tag. Hop 1: reactions (`kind:7`), deletions (`kind:5`, `kind:9005`), and edits (Buzz `kind:40003`) whose `e` tag is a row id. Hop 2: deletions whose `e` tag is a hop-1 event id (a delete-of-a-reaction). Each event appears at most once; access-scoped events the requester cannot read are omitted. Relays MAY cap each hop (Buzz: 1000 events per hop). @@ -100,7 +107,7 @@ For a valid window filter on an accessible channel (§Access Scoping) the relay The response is the surface's ordinary flat array of signed events — rows first in keyset order, then aux, then summaries, then bounds. Clients MUST partition by kind and MUST NOT rely on array position beyond the ordering of rows. -## Access Scoping +## Channel-mode Access Scoping Access is evaluated before any of the steps above. A syntactically valid window request for a channel the requester cannot access — including a channel that does not exist — MUST produce the relay's ordinary access-scoped result for that surface, with **no rows and no overlays**. For Buzz's query surface that ordinary result is an empty array, exactly as any other filter against an inaccessible channel produces. @@ -109,7 +116,7 @@ Two consequences implementers MUST NOT miss: - The "exactly one `kind:39006`" guarantee applies only to *served* windows — responses where access succeeded. The absence of a bounds overlay is therefore meaningful: it tells an extension-aware client that no window was served (access-scoped, or the relay does not implement this NIP — see §Degradation). - An inaccessible channel is thereby indistinguishable from a nonexistent one, but *not* from an accessible empty channel: the latter is a served window and does return a `39006` (`has_more: false`). This is the same existence-disclosure posture as the relay's ordinary reads — a requester who can query a channel at all was already entitled to know it exists. -## Overlay Event Formats +## Channel-mode Overlay Event Formats Overlays are signed by the relay identity and synthesized per response. Both kinds sit in the parameterized-replaceable range, so a client that caches them gets replace-by-`d`-tag semantics from NIP-01 with no special handling. Relays MUST reject client-submitted events of either kind at ingest. @@ -155,18 +162,185 @@ Exactly one per served window response. The **only** authority on exhaustion. Ta - `next_cursor` — the composite cursor to echo as `until` + `before_id` for the next page, or `null` iff `has_more` is `false`. - Reserved: an `oldest_retained` content field may be added (retention gap signaling) without a wire break. Clients MUST ignore unknown content fields. -## Client Behavior +## Channel-mode Client Behavior 1. **Head request**: send the window filter with no cursor. Render rows in received order. 2. **Continue**: read `kind:39006`; if `has_more`, send the same filter with `until = next_cursor.created_at`, `before_id = next_cursor.id`. Repeat until `has_more = false`. 3. **Exhaustion**: `39006.has_more` is the only exhaustion signal. `rows < limit` proves nothing — an exact-multiple final page returns `limit` rows with `has_more = false`, and predicate filtering can shrink any page. A client MUST NOT stop paging on row count, and MUST NOT treat a full page as "more available." 4. **Immutability**: fetched pages are immutable history chained cursor→cursor. New live events MUST NOT be spliced into fetched pages; deliver them through a separate live subscription (`since: now`) and merge at render time. On reconnect, refetch the head page and re-arm the live subscription; deeper pages need no repair. -5. **Bounds integrity**: a window response missing its `kind:39006`, or carrying more than one, or carrying one whose `d`-tag binding does not echo the request cursor, whose content is not parseable JSON, or whose content violates `has_more = true ⇔ next_cursor ≠ null`, is not a usable page — the client MUST discard it (and MAY retry) rather than guess at exhaustion. Clients SHOULD additionally reject overlays that violate the exact tag cardinality of §Overlay Event Formats or whose content fields have the wrong runtime types (hardening against a malformed or hostile serializer). Cryptographic verification is governed by §Overlay Trust. +5. **Bounds integrity**: a window response missing its `kind:39006`, or carrying more than one, or carrying one whose `d`-tag binding does not echo the request cursor, whose content is not parseable JSON, or whose content violates `has_more = true ⇔ next_cursor ≠ null`, is not a usable page — the client MUST discard it (and MAY retry) rather than guess at exhaustion. Clients SHOULD additionally reject overlays that violate the exact tag cardinality of §Channel-mode Overlay Event Formats or whose content fields have the wrong runtime types (hardening against a malformed or hostile serializer). Cryptographic verification is governed by §Overlay Trust. 6. **Overlays are metadata**: never render a `39005`/`39006` as a message, never feed one into cursor math, and key cached summaries by their `d` tag (latest wins). +## Legacy Oldest-first Threads + +Buzz's authenticated `POST /query` also supports an older thread path, separate +from both window modes. It is selected by a single `#e` root and `depth_limit`, +with `thread_window` absent or false: + +```jsonc +{ + "#h": [""], + "#e": [""], + "kinds": [9, 40002], + "depth_limit": 100, + "limit": 100, + "include_aux": true, + "thread_cursor": 1751500000, + "thread_cursor_id": "" +} +``` + +- Omit both cursor fields to start at the oldest reply. The historical + `thread_cursor: -1` start sentinel remains accepted. `depth_limit` is an + explicit maximum depth; `2147483647` is the existing unbounded-depth sentinel. + The row limit defaults to 100 and is capped at 500. +- Replies are ordered by `created_at ASC, id ASC`. To continue, derive + `thread_cursor` and `thread_cursor_id` from the **last loaded reply**, not + from an auxiliary event. The next page satisfies `created_at > cursor OR + (created_at = cursor AND id > cursor_id)`. The camel-case aliases + `threadCursor` and `threadCursorId` are also accepted. +- Timestamp-only continuation remains accepted for compatibility but skips + other replies in the same second; clients SHOULD send the composite pair. +- This path traverses stored thread metadata, not the newest-first mode's strict + row-kind filter. Clients SHOULD supply explicit `kinds` for query authorization + but MUST NOT assume this legacy thread path restricts its reply rows by them. +- `include_aux` appends root/reply reactions, edits and deletions, followed by + deletions of auxiliary events. These do not count against the reply limit. +- There is **no signed bounds event or server-issued continuation cursor**. + A full reply page can be the final page, so continuation may require one more + request. A short/empty reply page is the legacy stop heuristic, not a signed + exhaustion fact: access filtering can also shorten a response. + +A client explicitly falling back from thread windows MUST restart this path +from the oldest reply with clean pagination state. `until`/`before_id` are not +legacy thread continuation fields and MUST NOT be reused as such. + +## Thread Mode + +`thread_window: true` requests a newest-first page of replies to one root. It is +served through Buzz's NIP-98-authenticated `POST /query`; channel mode, +`kind:39006`, and the legacy oldest-first thread path are unchanged. The response +is a flat array of reply events, optional auxiliary events, and exactly one +relay-signed `kind:39007` bounds event. This specification adds no endpoint, +subscription, or around-target query. + +### Request + +```jsonc +{ + "thread_window": true, + "#h": [""], + "#e": [""], + "kinds": [9, 40002], + "depth_limit": 100, + "limit": 50, + "include_aux": true, + "until": 1751500000, + "before_id": "<64-hex id>" +} +``` + +`#h` and `#e` MUST each contain exactly one value. `kinds` MUST contain one to +four entries from `9`, `40002`, `45001`, and `45003`; relays normalize them to +a sorted, distinct list. `limit` defaults to 50 and MUST be 1–200; +`depth_limit` defaults to 100 and MUST be 1–100; `include_aux` defaults to +false. A continuation MUST provide both `until` and `before_id`, or neither for +the head page. Unknown fields, malformed values, and mixing thread windows with +another query mode MUST be rejected. A query MAY contain at most four window +filters. + +### Relay Processing + +For an authorized request the relay MUST: + +1. Verify that the root is a supported conversation event (`9`, `40002`, + `45001`, or `45003`) in the requested community and channel. Retained root + tombstones are eligible. A missing, out-of-scope, or unsupported root + (including a `40008` diff) returns no events or bounds; it MUST NOT receive + a signed exhausted page even if ingest has threaded replies beneath it. +2. Select non-deleted replies in that channel at depths 1 through + `depth_limit`, restricted by `kinds`, ordered by `created_at DESC, id ASC`. + A continuation retains rows where `created_at < until OR (created_at = until + AND id > before_id)`. +3. Evaluate `limit + 1` after all predicates. Discard the extra row and set + `has_more`; when more rows exist, `next_cursor` is the last retained scan + candidate, captured before event reconstruction. Otherwise it is `null`. +4. If `include_aux` is true, append reactions (`kind:7`), deletions + (`kind:5`/`9005`), and edits (`kind:40003`) targeting the root or returned + replies, followed by deletions targeting those auxiliary events. Preserve + original signatures, deduplicate by event ID, and apply access control to + every event. A target reference requires an `e` tag with the target ID + in its second position; containing both strings elsewhere is insufficient. +5. Append one bounds event to each served page, including empty and exhausted + pages. Refresh access for the **entire query batch** before releasing any + events or bounds. If the reader's accessible-channel set changed while any + window was built, discard all accumulated output and return a retryable + error (Buzz: HTTP `503`). This includes cross-channel auxiliary access. + +Rows alone count against `limit`. Auxiliary events and bounds do not. +Inaccessible and nonexistent channels return no events or bounds. + +### Resource Limits and Recovery + +Buzz shares resource limits across all thread-window filters in one query, +including replica retries: 64 auxiliary SQL scans, 8,192 raw reply and auxiliary +rows (including probes), 8 MiB of their combined content plus serialized tags, 8 MiB +of serialized output, and an eight-second overall deadline. Reply and auxiliary +payloads are consumed incrementally and charged before reconstruction; tombstones and +probes consume the raw allowance too. No partial page or bounds is returned +when a budget is exceeded. + +Resource exhaustion and database timeouts return HTTP `503`; malformed stored +auxiliary data remains a separate HTTP `500`. Clients MUST discard the entire +failed batch without advancing any cursor. For work/byte exhaustion, reduce +`limit`, split the batch, or explicitly request `include_aux: false` and fetch +needed auxiliary data separately. A root's auxiliary history can exceed the +allowance even at `limit: 1`; blind retries of the same query will not help. +For timeouts or changed access, retry with bounded backoff and fresh access. +These errors MUST NOT trigger legacy compatibility fallback. + +### Thread Bounds: `kind:39007` + +Bounds are synthesized per query and MUST NOT be accepted at ingest. Tags are +exactly one `d`, one `h`, and one `e`: + +```jsonc +{ + "kind": 39007, + "pubkey": "", + "tags": [ + ["d", "tw:1:"], + ["h", ""], + ["e", ""] + ], + "content": "{\"version\":1,\"direction\":\"older\",\"has_more\":true,\"next_cursor\":{\"created_at\":1751500000,\"id\":\"<64-hex id>\"}}" +} +``` + +The binding is SHA-256 over compact UTF-8 JSON of this ordered array: + +```jsonc +["tw",1,"older","","","","",50,100,[9,40002],null,true] +// host, reader, channel, root, limit, depth, sorted kinds, request cursor, include_aux +// cursor is null or [,""] +``` + +`host` is the server-resolved normalized authority and `reader` is the +authenticated lowercase pubkey. Clients MUST verify the expected relay signer +and signature, exact tags and request binding, version, direction, and +`next_cursor == null` iff `has_more == false`. Only validated bounds determine +exhaustion; clients MUST echo `next_cursor` as `until` and `before_id`. + +An extension-unaware relay may return legacy oldest-first history without +bounds. A client that explicitly falls back MUST restart with clean legacy +pagination state and MUST NOT reuse a descending cursor. Invalid signatures, +authorization failures, timeouts, corruption, or incomplete auxiliary closure +MUST NOT trigger compatibility fallback. Buzz currently ships no thread-mode +client opt-in or fallback implementation. + ## Degradation -Every extension field in this NIP is an *additional* key on a standard filter, and clients and relays that do not implement it need no changes: +The channel-mode extension fields in this NIP are *additional* keys on a standard filter, and clients and relays that do not implement it need no changes: - **Extension-unaware relay**: a tolerant filter parser (one that ignores unknown keys, as common NIP-01 implementations do) serves the filter as a plain `kinds` + `#h` query — a complete, correct, standard event stream. A strict parser may instead reject the filter outright. Both are safe: neither produces a wrong-but-plausible top-level timeline. A client MUST treat *either* signal — a response with no valid `kind:39006`, or an error/unsupported-filter response — as a downgrade, and fall back by reissuing a clean standard filter with all extension keys removed and assembling threads client-side. (Buzz's own WebSocket REQ path is such a tolerant parser: the filter deserializer drops the extension fields, so a window filter on REQ serves the standard query.) - **Extension-unaware client**: never sends `top_level`, never sees an overlay kind, and observes a completely standard relay. @@ -175,27 +349,27 @@ A relay implementing this NIP MAY advertise it in its NIP-11 relay information d ## Security and Privacy Considerations -Overlays are relay-authored facts about data the requester can already read. A relay MUST apply its normal access scoping to rows and to every aux-closure event, and §Access Scoping governs inaccessible channels: no rows, no overlays, no distinguishable error. +Overlays are relay-authored facts about data the requester can already read. A relay MUST apply the applicable mode's access rules to rows and to every aux-closure event. Inaccessible channels produce no rows, no overlays, and no distinguishable existence error; thread mode additionally validates the requested root within the host-derived community and channel before serving descendants. `kind:39005` aggregates thread activity (participant pubkeys, counts, recency) into one event. It only ever describes threads rooted in a channel the requester can read, so it reveals nothing a client could not compute from readable events — it saves round trips, not permissions. -Client-submitted `39005`/`39006` MUST be rejected at ingest (relay-only kinds); a forged overlay accepted into storage could later masquerade as relay-signed state. +Client-submitted `39005`/`39006`/`39007` events MUST be rejected at ingest (relay-only kinds); a forged overlay accepted into storage could later masquerade as relay-signed state. ### Overlay Trust -Because `kind:39006` is the pagination authority, a client MUST adopt exactly one of these trust profiles before using the window fast path: +`kind:39006` and `kind:39007` are the pagination authorities for channel and thread mode respectively. Thread-mode clients MUST perform the signer, signature, tag, and request-binding verification specified in §Thread Bounds. Before using the channel-window fast path, a client MUST adopt exactly one of these trust profiles: -- **Authenticated-transport profile** (what Buzz desktop ships): the client speaks to a relay it deliberately configured as its source of truth, over TLS (HTTPS/WSS) to that configured origin — server-origin authentication comes from the TLS certificate chain, which is what proves the response bytes came from the relay. (NIP-98 request signing and NIP-42 auth run over this channel too, but they authenticate the *requester* to the relay for access control; they are not evidence of response provenance.) The MUST-level structural checks of §Client Behavior step 5 — exactly one bounds, request binding, parseable content, `has_more`/`next_cursor` agreement — are still mandatory and are what #1500 enforces. The SHOULD-level checks of step 5 (exact tag cardinality, runtime field-type validation) and cryptographically binding overlay signatures to the advertised NIP-11 identity are future hardening, to be applied uniformly across all relay-signed reads (with NIP-DV, NIP-IA), not a current guarantee. Under this profile, "relay-signed" is a TLS-origin claim, not a client-verified cryptographic one. -- **Identity-verified profile**: the client has obtained and trusts the relay identity pubkey out-of-band or via NIP-11. It MUST verify each overlay's event id, Schnorr signature, and signer against that identity, and treat any failure as the §step-5 discard. This is the profile for clients that cannot or do not authenticate their transport end-to-end. +- **Authenticated-transport profile** (what Buzz desktop ships): the client speaks to a relay it deliberately configured as its source of truth, over TLS (HTTPS/WSS) to that configured origin — server-origin authentication comes from the TLS certificate chain, which is what proves the response bytes came from the relay. (NIP-98 request signing and NIP-42 auth run over this channel too, but they authenticate the *requester* to the relay for access control; they are not evidence of response provenance.) The MUST-level structural checks of §Channel-mode Client Behavior step 5 — exactly one bounds, request binding, parseable content, `has_more`/`next_cursor` agreement — are still mandatory and are what #1500 enforces. The SHOULD-level checks of step 5 (exact tag cardinality, runtime field-type validation) and cryptographically binding channel overlay signatures to the advertised NIP-11 identity are future hardening, to be applied uniformly across relay-signed reads (with NIP-DV, NIP-IA), not a current channel-mode guarantee. Under this profile, "relay-signed" is a TLS-origin claim, not a client-verified cryptographic one. +- **Identity-verified profile**: the client has obtained and trusts the relay identity pubkey out-of-band or via NIP-11. It MUST verify each overlay's event id, Schnorr signature, and signer against that identity, and treat any failure as the §Channel-mode Client Behavior step-5 discard. This is the profile for clients that cannot or do not authenticate their transport end-to-end. -A client with neither an authenticated transport nor a verifiable relay identity MUST NOT use the window fast path: it falls back to the standard filter (§Degradation), where it verifies every event signature itself. +A channel-mode client with neither an authenticated transport nor a verifiable relay identity MUST NOT use the channel-window fast path: it falls back to the standard filter (§Degradation), where it verifies every event signature itself. Thread mode has the stricter verification and clean-legacy-restart rules in §Thread Bounds; it MUST NOT downgrade on an invalid signed response. ## Implementation Gotchas -- The `limit + 1` probe MUST run after *all* predicates (access, deletion, top-level, `kinds`). A probe over a superset produces false `has_more = true` on the last page. +- The `limit + 1` probe MUST run after *all* predicates: access, deletion and `kinds` in both modes; top-level classification in channel mode only; root, channel and depth restrictions in thread mode. A probe over a superset produces false `has_more = true` on the last page. - The cursor comparison uses `id > $id` (bytewise ascending) because the total order is `created_at DESC, id ASC`. Getting the id inequality backwards drops or duplicates same-second rows — precisely the bug the composite cursor removes. - `next_cursor` is the last retained *scan candidate*, not the last delivered row: capture the scan position before per-event reconstruction so a skipped event cannot stall pagination. Clients echo it verbatim and never derive or validate it against the rows they received. -- Events ingested before the relay computed thread metadata have no depth; they MUST be treated as top-level rather than vanishing from every window. +- **Channel mode only:** events ingested before the relay computed thread metadata have no depth; they MUST be treated as top-level rather than vanishing from channel windows. Thread mode instead requires metadata at depths 1..`depth_limit`. - The `d` tag on `39006` differs per request cursor by design: concurrent pages of one channel coexist in a replaceable-event cache instead of clobbering each other. The per-channel-singleton alternative would make page N overwrite page N+1's bounds. ## Relation to Other NIPs diff --git a/docs/thread-window-deployment.md b/docs/thread-window-deployment.md new file mode 100644 index 00000000000..2b35f5a1725 --- /dev/null +++ b/docs/thread-window-deployment.md @@ -0,0 +1,70 @@ +# Thread-window index deployment + +Migration `0049_thread_window_index.sql` adds the index used by opt-in, +newest-first thread windows: + +```sql +CREATE INDEX idx_thread_metadata_window + ON public.thread_metadata + (community_id, root_event_id, event_created_at DESC, event_id ASC); +``` + +The migration bounds both lock acquisition and execution time. That is safe for +fresh or small databases, but it deliberately fails instead of holding a +write-conflicting table lock while building the index on a populated database. +Brownfield deployments must prebuild the exact index concurrently before +starting a relay version that includes migration 0049. + +## Brownfield procedure + +Run the following against the relay database with a role allowed to create an +index. `CREATE INDEX CONCURRENTLY` cannot run inside a transaction block. + +```sql +CREATE INDEX CONCURRENTLY idx_thread_metadata_window + ON public.thread_metadata + (community_id, root_event_id, event_created_at DESC, event_id ASC); +``` + +Then verify that PostgreSQL considers the index ready, live, valid, and an exact +match for the definition expected by migration 0049: + +```sql +SELECT + i.indisvalid, + i.indisready, + i.indislive, + pg_get_indexdef(i.indexrelid) AS definition +FROM pg_index AS i +JOIN pg_class AS c ON c.oid = i.indexrelid +JOIN pg_namespace AS n ON n.oid = c.relnamespace +WHERE n.nspname = 'public' + AND c.relname = 'idx_thread_metadata_window'; +``` + +Expected flags are all `true`. The expected definition is: + +```text +CREATE INDEX idx_thread_metadata_window ON public.thread_metadata USING btree (community_id, root_event_id, event_created_at DESC, event_id) +``` + +After verification, deploy normally. Migration 0049 detects the prebuilt index, +skips the write-conflicting `CREATE INDEX`, validates the catalog shape again, +and records the migration. + +## Recovery + +A cancelled or failed concurrent build can leave an invalid index behind. Do not +retry deployment against that remnant: migration 0049 rejects invalid, unready, +non-live, or differently defined indexes. + +Remove only the failed index, outside a transaction, then repeat the prebuild and +verification steps: + +```sql +DROP INDEX CONCURRENTLY IF EXISTS public.idx_thread_metadata_window; +``` + +Do not replace this with a non-concurrent build on a populated production table. +The bounded startup failure is intentional; ingestion availability takes +precedence over completing the rollout in one attempt. diff --git a/migrations/0049_thread_window_index.sql b/migrations/0049_thread_window_index.sql new file mode 100644 index 00000000000..7bcdc82bd0f --- /dev/null +++ b/migrations/0049_thread_window_index.sql @@ -0,0 +1,29 @@ +-- Newest-first thread keyset index. Additive: legacy root/parent indexes remain. +-- Brownfield operators MUST prebuild concurrently as documented in +-- docs/thread-window-deployment.md. +-- Startup is deliberately bounded: a busy/large table fails deployment rather +-- than blocking ingestion for an unbounded index build. Retry after prebuild. +SET LOCAL lock_timeout = '1s'; +SET LOCAL statement_timeout = '5s'; +-- IF NOT EXISTS still requests a writer-conflicting ShareLock before checking +-- whether the index exists. Bypass CREATE entirely for prebuilt indexes, then +-- validate them below; a failed concurrent build must never count as success. +DO $$ +BEGIN + IF to_regclass('public.idx_thread_metadata_window') IS NULL THEN + CREATE INDEX idx_thread_metadata_window + ON public.thread_metadata (community_id, root_event_id, event_created_at DESC, event_id ASC); + END IF; + + IF NOT EXISTS ( + SELECT 1 FROM pg_index i + JOIN pg_class c ON c.oid = i.indexrelid + JOIN pg_namespace n ON n.oid = c.relnamespace + WHERE n.nspname = 'public' AND c.relname = 'idx_thread_metadata_window' + AND i.indisvalid AND i.indisready AND i.indislive + AND pg_get_indexdef(i.indexrelid) = + 'CREATE INDEX idx_thread_metadata_window ON public.thread_metadata USING btree (community_id, root_event_id, event_created_at DESC, event_id)' + ) THEN + RAISE EXCEPTION 'idx_thread_metadata_window invalid or wrong definition; see docs/thread-window-deployment.md'; + END IF; +END $$; diff --git a/schema/schema.sql b/schema/schema.sql index 797a83f5d10..1a0a849f7ac 100644 --- a/schema/schema.sql +++ b/schema/schema.sql @@ -531,6 +531,8 @@ CREATE TABLE thread_metadata ( CREATE INDEX idx_thread_metadata_parent ON thread_metadata (community_id, parent_event_id); CREATE INDEX idx_thread_metadata_root ON thread_metadata (community_id, root_event_id); +CREATE INDEX idx_thread_metadata_window + ON thread_metadata (community_id, root_event_id, event_created_at DESC, event_id ASC); CREATE INDEX idx_thread_metadata_channel_depth ON thread_metadata (community_id, channel_id, depth, event_created_at); CREATE INDEX idx_thread_metadata_event_id ON thread_metadata (community_id, event_id);