diff --git a/CHANGELOG.md b/CHANGELOG.md index 41a27dd..07ecdd6 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -7,6 +7,16 @@ and this project adheres to [Semantic Versioning](https://semver.org/spec/v2.0.0 ## [Unreleased] +### Added + +- Draft NIP-91 `&` tag filters now require every listed tag value for + historical queries, live subscriptions, COUNT, search, router post-filtering, + and negentropy. NIP-11 advertises NIP-91 plus the distinct-tag and aggregate + AND-entry limits. Outbound router `REQ`s and `wok sync` `NEG-OPEN`s keep `&` + clauses only for upstreams whose NIP-11 document advertises NIP-91, and + otherwise fold each `&x` value set into the `#x` compatibility clause and + apply the exact AND semantics locally. + ## [0.3.1] - 2026-08-15 ### Fixed diff --git a/Cargo.lock b/Cargo.lock index 7f67a1f..90d7452 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -223,6 +223,17 @@ version = "0.2.2" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "f079e83a288787bcd14a6aea84cee5c87a67c5a3e660c30f557a3d24761b3527" +[[package]] +name = "chacha20" +version = "0.10.1" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "d524456ba66e72eb8b115ff89e01e497f8e6d11d78b70b1aa13c0fbd97540a81" +dependencies = [ + "cfg-if", + "cpufeatures 0.3.0", + "rand_core 0.10.1", +] + [[package]] name = "chrono" version = "0.4.45" @@ -305,6 +316,15 @@ dependencies = [ "libc", ] +[[package]] +name = "cpufeatures" +version = "0.3.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "8b2a41393f66f16b0823bb79094d54ac5fbd34ab292ddafb9a0456ac9f87d201" +dependencies = [ + "libc", +] + [[package]] name = "crc32fast" version = "1.5.0" @@ -555,8 +575,10 @@ source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "ff2abc00be7fca6ebc474524697ae276ad847ad0a6b3faa4bcb027e9a4614ad0" dependencies = [ "cfg-if", + "js-sys", "libc", "wasi", + "wasm-bindgen", ] [[package]] @@ -578,8 +600,11 @@ source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "300e883d756b2e4ec94e02791f39b04b522276138852cfc41d9fb7e904106099" dependencies = [ "cfg-if", + "js-sys", "libc", "r-efi 6.0.0", + "rand_core 0.10.1", + "wasm-bindgen", ] [[package]] @@ -695,6 +720,23 @@ dependencies = [ "pin-project-lite", "smallvec", "tokio", + "want", +] + +[[package]] +name = "hyper-rustls" +version = "0.27.9" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "33ca68d021ef39cf6463ab54c1d0f5daf03377b70561305bb89a8f83aab66e0f" +dependencies = [ + "http", + "hyper", + "hyper-util", + "rustls", + "rustls-native-certs", + "tokio", + "tokio-rustls", + "tower-service", ] [[package]] @@ -703,12 +745,21 @@ version = "0.1.20" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "96547c2556ec9d12fb1578c4eaf448b04993e7fb79cbaad930a656880a6bdfa0" dependencies = [ + "base64", "bytes", + "futures-channel", + "futures-util", "http", "http-body", "hyper", + "ipnet", + "libc", + "percent-encoding", "pin-project-lite", + "socket2 0.6.5", "tokio", + "tower-service", + "tracing", ] [[package]] @@ -868,6 +919,12 @@ dependencies = [ "libc", ] +[[package]] +name = "ipnet" +version = "2.12.1" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "6a756c3fac73139e83f14c2d742155dd2b78d3ee56597b419a0579b7bdd6dd78" + [[package]] name = "is_terminal_polyfill" version = "1.70.2" @@ -971,6 +1028,12 @@ version = "0.4.33" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "0ceec5bc11778974d1bcb055b18002eba7f4b3518b6a0081b3af5f21666da9ad" +[[package]] +name = "lru-slab" +version = "0.1.2" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "112b39cec0b298b6c1999fee3e31427f74f676e4cb9879ed1a121b43661a4154" + [[package]] name = "matchers" version = "0.2.0" @@ -1231,6 +1294,62 @@ version = "1.2.3" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "a1d01941d82fa2ab50be1e79e6714289dd7cde78eba4c074bc5a4374f650dfe0" +[[package]] +name = "quinn" +version = "0.11.11" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "0c1a41e437b6bbd489372cd4971de128e85c855f56c57f283d20ff016cf7c0a8" +dependencies = [ + "bytes", + "cfg_aliases", + "pin-project-lite", + "quinn-proto", + "quinn-udp", + "rustc-hash", + "rustls", + "socket2 0.5.10", + "thiserror", + "tokio", + "tracing", + "web-time", +] + +[[package]] +name = "quinn-proto" +version = "0.11.16" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "2f4bfc015262b9df63c8845072ce59068853ff5872180c2ce2f13038b970e560" +dependencies = [ + "bytes", + "getrandom 0.4.3", + "lru-slab", + "rand 0.10.2", + "rand_pcg", + "ring", + "rustc-hash", + "rustls", + "rustls-pki-types", + "slab", + "thiserror", + "tinyvec", + "tracing", + "web-time", +] + +[[package]] +name = "quinn-udp" +version = "0.5.15" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "35a133f956daabe89a61a685c2649f13d82d5aa4bd5d12d1277e1072a21c0694" +dependencies = [ + "cfg_aliases", + "libc", + "once_cell", + "socket2 0.5.10", + "tracing", + "windows-sys 0.52.0", +] + [[package]] name = "quote" version = "1.0.47" @@ -1273,6 +1392,17 @@ dependencies = [ "rand_core 0.9.5", ] +[[package]] +name = "rand" +version = "0.10.2" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "c7f5fa3a058cd35567ef9bfa5e75732bee0f9e4c55fa90477bef2dfcdbc4be80" +dependencies = [ + "chacha20", + "getrandom 0.4.3", + "rand_core 0.10.1", +] + [[package]] name = "rand_chacha" version = "0.3.1" @@ -1311,6 +1441,21 @@ dependencies = [ "getrandom 0.3.4", ] +[[package]] +name = "rand_core" +version = "0.10.1" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "63b8176103e19a2643978565ca18b50549f6101881c443590420e4dc998a3c69" + +[[package]] +name = "rand_pcg" +version = "0.10.2" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "caa0f4137e1c0a72f4c651489402276c8e8e1cf081f3b0ba156d2cbeef09e86a" +dependencies = [ + "rand_core 0.10.1", +] + [[package]] name = "rand_xorshift" version = "0.4.0" @@ -1366,6 +1511,44 @@ version = "0.8.11" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "d6f6ff9a378485b298a5286656da665ba74413d36db0979633275d2e708145d4" +[[package]] +name = "reqwest" +version = "0.12.28" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "eddd3ca559203180a307f12d114c268abf583f59b03cb906fd0b3ff8646c1147" +dependencies = [ + "base64", + "bytes", + "futures-core", + "http", + "http-body", + "http-body-util", + "hyper", + "hyper-rustls", + "hyper-util", + "js-sys", + "log", + "percent-encoding", + "pin-project-lite", + "quinn", + "rustls", + "rustls-native-certs", + "rustls-pki-types", + "serde", + "serde_json", + "serde_urlencoded", + "sync_wrapper", + "tokio", + "tokio-rustls", + "tower", + "tower-http", + "tower-service", + "url", + "wasm-bindgen", + "wasm-bindgen-futures", + "web-sys", +] + [[package]] name = "ring" version = "0.17.14" @@ -1380,6 +1563,12 @@ dependencies = [ "windows-sys 0.52.0", ] +[[package]] +name = "rustc-hash" +version = "2.1.3" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "6b1e7f9a428571be2dc5bc0505c13fb6bf936822b894ec87abf8a08a4e51742d" + [[package]] name = "rustc_version" version = "0.4.1" @@ -1434,6 +1623,7 @@ version = "1.15.1" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "2f4925028c7eb5d1fcdaf196971378ed9d2c1c4efc7dc5d011256f76c99c0a96" dependencies = [ + "web-time", "zeroize", ] @@ -1466,6 +1656,12 @@ dependencies = [ "wait-timeout", ] +[[package]] +name = "ryu" +version = "1.0.23" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "9774ba4a74de5f7b1c1451ed6cd5285a32eddb5cccb8cc655a4e50009e06477f" + [[package]] name = "same-file" version = "1.0.6" @@ -1591,6 +1787,18 @@ dependencies = [ "serde", ] +[[package]] +name = "serde_urlencoded" +version = "0.7.1" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "d3491c14715ca2294c4d6a88f15e84739788c1d030eed8c110436aafdaa2f3fd" +dependencies = [ + "form_urlencoded", + "itoa", + "ryu", + "serde", +] + [[package]] name = "sha1" version = "0.10.7" @@ -1598,7 +1806,7 @@ source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "a978451301f4db1d02937a4ab3ccce137717b81826e79b7d49ffe3244a13c3b8" dependencies = [ "cfg-if", - "cpufeatures", + "cpufeatures 0.2.17", "digest", ] @@ -1609,7 +1817,7 @@ source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "a7507d819769d01a365ab707794a4084392c824f54a7a6a7862f8c3d0892b283" dependencies = [ "cfg-if", - "cpufeatures", + "cpufeatures 0.2.17", "digest", ] @@ -1716,6 +1924,15 @@ dependencies = [ "unicode-ident", ] +[[package]] +name = "sync_wrapper" +version = "1.0.2" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "0bf256ce5efdfa370213c1dabab5935a12e49f2c58d15e9eac2870d3b4f27263" +dependencies = [ + "futures-core", +] + [[package]] name = "synstructure" version = "0.13.2" @@ -1793,6 +2010,21 @@ dependencies = [ "zerovec", ] +[[package]] +name = "tinyvec" +version = "1.12.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "bb4ebadaa0af04fab11ae01eb5f9fdb5f9c5b875506e210e71c07873528baa7f" +dependencies = [ + "tinyvec_macros", +] + +[[package]] +name = "tinyvec_macros" +version = "0.1.1" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "1f3ccbac311fea05f86f61904b462b55fb3df8837a366dfc601a0161d0532f20" + [[package]] name = "tokio" version = "1.53.1" @@ -1888,6 +2120,51 @@ version = "0.1.2" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "5d99f8c9a7727884afe522e9bd5edbfc91a3312b36a77b5fb8926e4c31a41801" +[[package]] +name = "tower" +version = "0.5.3" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "ebe5ef63511595f1344e2d5cfa636d973292adc0eec1f0ad45fae9f0851ab1d4" +dependencies = [ + "futures-core", + "futures-util", + "pin-project-lite", + "sync_wrapper", + "tokio", + "tower-layer", + "tower-service", +] + +[[package]] +name = "tower-http" +version = "0.6.11" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "4cfcf7e2740e6fc6d4d688b4ef00650406bb94adf4731e43c096c3a19fe40840" +dependencies = [ + "bitflags", + "bytes", + "futures-util", + "http", + "http-body", + "pin-project-lite", + "tower", + "tower-layer", + "tower-service", + "url", +] + +[[package]] +name = "tower-layer" +version = "0.3.3" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "121c2a6cda46980bb0fcd1647ffaf6cd3fc79a013de288782836f6df9c48780e" + +[[package]] +name = "tower-service" +version = "0.3.3" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "8df9b6e13f2d32c91b9bd719c00d1958837bc7dec474d94952798cc8e69eeec3" + [[package]] name = "tracing" version = "0.1.44" @@ -1962,6 +2239,12 @@ dependencies = [ "tracing-serde", ] +[[package]] +name = "try-lock" +version = "0.2.5" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "e421abadd41a4225275504ea4d6566923418b7f05506fbc9c0fe86ba7396114b" + [[package]] name = "tungstenite" version = "0.26.2" @@ -2066,6 +2349,15 @@ dependencies = [ "winapi-util", ] +[[package]] +name = "want" +version = "0.3.1" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "bfa7760aed19e106de2c7c0b581b509f2f25d3dacaf737cb82ac61bc6d760b0e" +dependencies = [ + "try-lock", +] + [[package]] name = "wasi" version = "0.11.1+wasi-snapshot-preview1" @@ -2094,6 +2386,16 @@ dependencies = [ "wasm-bindgen-shared", ] +[[package]] +name = "wasm-bindgen-futures" +version = "0.4.77" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "6b7777d5cc23d0e91404e53ce2d5e8ec7acae3026b16233dba62cd3246457950" +dependencies = [ + "js-sys", + "wasm-bindgen", +] + [[package]] name = "wasm-bindgen-macro" version = "0.2.127" @@ -2126,6 +2428,26 @@ dependencies = [ "unicode-ident", ] +[[package]] +name = "web-sys" +version = "0.3.104" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "c435338968042f4f59a557f690a253676d47ce13ceb55d70100e7facf6620a30" +dependencies = [ + "js-sys", + "wasm-bindgen", +] + +[[package]] +name = "web-time" +version = "1.1.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "5a6580f308b1fad9207618087a65c04e7a10bc77e02c8e84e9b00dd4b12fa0bb" +dependencies = [ + "js-sys", + "wasm-bindgen", +] + [[package]] name = "winapi" version = "0.3.9" @@ -2481,6 +2803,7 @@ dependencies = [ "notify", "parking_lot", "rand 0.8.7", + "reqwest", "rustls", "secp256k1", "serde", diff --git a/Cargo.toml b/Cargo.toml index cc8eaec..ba797ac 100644 --- a/Cargo.toml +++ b/Cargo.toml @@ -40,6 +40,7 @@ parking_lot = "0.12" proptest = "1" rand = "0.8" regex = "1" +reqwest = { version = "0.12", default-features = false, features = ["rustls-tls-native-roots"] } # Select the crypto provider used by tokio-tungstenite's Rustls client. rustls = { version = "0.23", default-features = false, features = ["ring", "std", "tls12"] } secp256k1 = { version = "0.30", features = ["global-context", "hashes", "rand"] } diff --git a/README.md b/README.md index 4b6cbed..14babad 100644 --- a/README.md +++ b/README.md @@ -32,7 +32,7 @@ quirk. It also provides an additional Unix-domain socket transport. ranked content search, NIP-62 restart-safe Request to Vanish, NIP-70 protected events, NIP-59 gift-wrap deletion semantics, NIP-77 negentropy set reconciliation** (persistent LMDB B-tree, tree-backed multi-round sync - sessions). + sessions), and **draft NIP-91 AND tag filters**. - **Standards-first ephemeral delivery**: ephemeral kinds are live-only by default, with an explicit persisted TTL compatibility mode. - **permessage-deflate** via an in-house RFC 6455/7692 codec (no Rust WS diff --git a/crates/wok-cli/Cargo.toml b/crates/wok-cli/Cargo.toml index 207812b..344c261 100644 --- a/crates/wok-cli/Cargo.toml +++ b/crates/wok-cli/Cargo.toml @@ -19,6 +19,7 @@ libc.workspace = true notify.workspace = true parking_lot.workspace = true rand.workspace = true +reqwest.workspace = true rustls.workspace = true serde.workspace = true serde_json.workspace = true diff --git a/crates/wok-cli/src/main.rs b/crates/wok-cli/src/main.rs index 5e38d12..44b7f73 100644 --- a/crates/wok-cli/src/main.rs +++ b/crates/wok-cli/src/main.rs @@ -25,9 +25,10 @@ fn foreach_by_filter_scan( filter: &serde_json::Value, max_limit: u64, max_tags: usize, + max_and_entries: usize, cb: impl FnMut(u64), ) -> Result<(), wok_query::QueryError> { - wok_query::foreach_by_filter(txn, filter, max_limit, max_tags, cb) + wok_query::foreach_by_filter(txn, filter, max_limit, max_tags, max_and_entries, cb) } /// Read one line from `reader` into `buf` (newline stripped, like @@ -773,6 +774,7 @@ fn cmd_scan(cfg: &Config, filter: &str, count: bool) -> Result<()> { &filter, cfg.relay.max_filter_limit, cfg.relay.max_tags_per_filter, + cfg.relay.max_and_entries, |lev| { n += 1; if !count { @@ -837,6 +839,7 @@ fn cmd_delete(cfg: &Config, age: Option, filter: Option, dry_run: b &filter, u64::MAX, cfg.relay.max_tags_per_filter, + cfg.relay.max_and_entries, |lev| levs.push(lev), )?; } @@ -888,6 +891,7 @@ fn cmd_monitor(cfg: &Config) -> Result<()> { &arr[3], cfg.relay.max_filter_limit, cfg.relay.max_tags_per_filter, + cfg.relay.max_and_entries, )?; let mut s = wok_query::Subscription::new(conn, wok_query::SubId::new(sub)?, fg, false); @@ -961,6 +965,7 @@ fn cmd_dict(cfg: &Config, cmd: DictCmd) -> Result<()> { &filter, u64::MAX, cfg.relay.max_tags_per_filter, + cfg.relay.max_and_entries, |lev| levs.push(lev), )?; } @@ -1120,7 +1125,12 @@ fn cmd_neg(cfg: &Config, cmd: NegCmd) -> Result<()> { } NegCmd::Add { filter } => { let v: serde_json::Value = serde_json::from_str(&filter)?; - let compiled = wok_query::NostrFilterGroup::from_value(&v, u64::MAX, 64)?; + let compiled = wok_query::NostrFilterGroup::from_value( + &v, + u64::MAX, + 64, + cfg.relay.max_and_entries, + )?; if compiled.filters.is_empty() { bail!("filter will never match"); } @@ -1181,7 +1191,8 @@ fn build_negentropy_tree(env: &Env, tree_id: u64, batch_size: usize) -> Result<( })?; let filter: serde_json::Value = serde_json::from_str(&filter_str.context("couldn't find treeId")?)?; - let compiled = wok_query::NostrFilterGroup::from_value(&filter, u64::MAX, 64)?; + let compiled = + wok_query::NostrFilterGroup::from_value(&filter, u64::MAX, usize::MAX, usize::MAX)?; if compiled.requires_content() { bail!("negentropy filters do not support content search"); } @@ -1343,6 +1354,26 @@ fn process_range_option(range: &str, filter: &mut serde_json::Value) -> Result<( Ok(()) } +fn downloaded_event_matches( + filter_group: &wok_query::NostrFilterGroup, + event: &serde_json::Value, + limits: &EventLimits, +) -> bool { + wok_event::nostr_json_to_packed_event(event, limits) + .map(|packed| { + if filter_group.requires_content() { + event + .get("content") + .and_then(serde_json::Value::as_str) + .map(|content| filter_group.does_match_with_content(packed.view(), content)) + .unwrap_or(false) + } else { + filter_group.does_match(packed.view()) + } + }) + .unwrap_or(false) +} + /// Verify and write a batch of downloaded events, updating negentropy trees /// like C++ WriterPipeline (verifyMsg + verifyTime). fn write_downloaded( @@ -1421,6 +1452,7 @@ async fn cmd_sync( &filter_json, u64::MAX, cfg.relay.max_tags_per_filter, + cfg.relay.max_and_entries, )?; let env = open_env(cfg)?; @@ -1457,6 +1489,7 @@ async fn cmd_sync( &filter_json, u64::MAX, cfg.relay.max_tags_per_filter, + cfg.relay.max_and_entries, |lev| levs.push(lev), )?; levs.sort_unstable(); @@ -1541,9 +1574,13 @@ async fn cmd_sync( } }; - let mut ws = mesh::connect_mesh(&url, cfg.events.max_event_size).await?; + let connect = mesh::connect_mesh(&url, cfg.events.max_event_size); + let capability_filter = mesh::outbound_filter_for_relay(&url, &filter_json); + let (ws, remote_filter) = tokio::join!(connect, capability_filter); + let mut ws = ws?; + let remote_filter = remote_filter?; let init = initiate(&env)?; - let open = serde_json::json!(["NEG-OPEN", "N", filter_json, hex::encode(init)]); + let open = serde_json::json!(["NEG-OPEN", "N", remote_filter, hex::encode(init)]); ws.send(Message::Text(open.to_string().into())).await?; const HIGH_WATER_UP: usize = 100; @@ -1657,6 +1694,11 @@ async fn cmd_sync( } "EVENT" => { if let Some(ev) = v.get(2) { + let matches_filter = + downloaded_event_matches(&filter_group, ev, &cfg.event_limits()); + if !matches_filter { + continue; + } batch.push(ev.clone()); if batch.len() >= 1000 { write_downloaded(&env, cfg, &mut batch, &mut written)?; @@ -2107,6 +2149,57 @@ mod main_tests { assert_eq!(next_reconnect_delay(maximum, maximum), maximum); } + #[test] + fn nip91_sync_post_filters_legacy_relay_results() { + let filter = wok_query::NostrFilterGroup::from_value( + &json!({"&t":["meme", "cat"], "#t":["meme", "cat", "black"]}), + 500, + 3, + 16, + ) + .unwrap(); + let event = |tags: Vec| { + json!({ + "id":"11".repeat(32), + "pubkey":"22".repeat(32), + "created_at":1, + "kind":1, + "tags":tags, + "content":"", + "sig":"33".repeat(64) + }) + }; + assert!(downloaded_event_matches( + &filter, + &event(vec![ + json!(["t", "meme"]), + json!(["t", "cat"]), + json!(["t", "black"]) + ]), + &EventLimits::default(), + )); + assert!(!downloaded_event_matches( + &filter, + &event(vec![json!(["t", "meme"]), json!(["t", "cat"])]), + &EventLimits::default(), + )); + + let search = wok_query::NostrFilterGroup::from_value( + &json!({"search":"needle", "&t":["meme", "cat"]}), + 500, + 3, + 16, + ) + .unwrap(); + let mut searchable = event(vec![json!(["t", "meme"]), json!(["t", "cat"])]); + searchable["content"] = json!("a sync needle"); + assert!(downloaded_event_matches( + &search, + &searchable, + &EventLimits::default(), + )); + } + #[tokio::test] async fn stream_reconnects_after_remote_close_without_blocking_runtime() { let listener = tokio::net::TcpListener::bind("127.0.0.1:0").await.unwrap(); diff --git a/crates/wok-cli/src/mesh.rs b/crates/wok-cli/src/mesh.rs index 250940b..bf639ce 100644 --- a/crates/wok-cli/src/mesh.rs +++ b/crates/wok-cli/src/mesh.rs @@ -1,10 +1,284 @@ //! Shared helpers for outbound mesh (relay-to-relay) websocket connections. +use futures_util::future::{BoxFuture, FutureExt, Shared}; +use serde_json::Value; +use std::collections::{HashMap, HashSet}; +use std::sync::Arc; +use std::time::{Duration, Instant}; use tokio::net::TcpStream; use tokio_tungstenite::tungstenite::protocol::WebSocketConfig; use tokio_tungstenite::tungstenite::Error as WsError; use tokio_tungstenite::{MaybeTlsStream, WebSocketStream}; +const NIP11_TIMEOUT: Duration = Duration::from_secs(3); +const NIP11_VALID_TTL: Duration = Duration::from_secs(60 * 60); +const NIP11_FAILURE_TTL: Duration = Duration::from_secs(5 * 60); +const NIP11_MAX_BYTES: usize = 256 * 1024; +const NIP11_CACHE_CAPACITY: usize = 256; + +#[derive(Clone)] +struct CapabilityResult { + supported_nips: Result>, Arc>, + expires_at: Instant, +} + +type CapabilityFuture = Shared>>; + +struct CacheEntry { + generation: u64, + last_used: Instant, + result: CapabilityFuture, +} + +#[derive(Default)] +struct CapabilityCache { + entries: HashMap, + next_generation: u64, +} + +static NIP11_CLIENT: std::sync::LazyLock = std::sync::LazyLock::new(|| { + reqwest::Client::builder() + .timeout(NIP11_TIMEOUT) + .redirect(reqwest::redirect::Policy::none()) + .build() + .expect("NIP-11 HTTP client configuration is valid") +}); + +static CAPABILITY_CACHE: std::sync::LazyLock> = + std::sync::LazyLock::new(|| tokio::sync::Mutex::new(CapabilityCache::default())); + +fn nip11_url(relay_url: &str) -> Option { + let mut url = url::Url::parse(relay_url).ok()?; + let scheme = match url.scheme() { + "ws" => "http", + "wss" => "https", + _ => return None, + }; + url.set_scheme(scheme).ok()?; + url.set_fragment(None); + Some(url.to_string()) +} + +async fn fetch_nip11(http_url: String) -> Arc { + let fetched_at = Instant::now(); + let fetched = async { + let mut response = NIP11_CLIENT + .get(&http_url) + .header(reqwest::header::ACCEPT, "application/nostr+json") + .send() + .await?; + if !response.status().is_success() { + anyhow::bail!("HTTP {}", response.status()); + } + if response + .content_length() + .is_some_and(|length| length > NIP11_MAX_BYTES as u64) + { + anyhow::bail!("response exceeds {NIP11_MAX_BYTES} bytes"); + } + let mut bytes = Vec::new(); + while let Some(chunk) = response.chunk().await? { + if bytes.len().saturating_add(chunk.len()) > NIP11_MAX_BYTES { + anyhow::bail!("response exceeds {NIP11_MAX_BYTES} bytes"); + } + bytes.extend_from_slice(&chunk); + } + let document: Value = serde_json::from_slice(&bytes)?; + parse_supported_nips(&document) + } + .await; + + match fetched { + Ok(supported_nips) => { + tracing::debug!(url = %http_url, ?supported_nips, "cached relay NIP-11 capabilities"); + Arc::new(CapabilityResult { + supported_nips: Ok(Arc::new(supported_nips)), + expires_at: fetched_at + NIP11_VALID_TTL, + }) + } + Err(error) => { + tracing::warn!(url = %http_url, %error, "NIP-11 capability lookup failed"); + Arc::new(CapabilityResult { + supported_nips: Err(Arc::from(error.to_string())), + expires_at: fetched_at + NIP11_FAILURE_TTL, + }) + } + } +} + +fn parse_supported_nips(document: &Value) -> anyhow::Result> { + let Some(nips) = document.get("supported_nips") else { + return Ok(HashSet::new()); + }; + let nips = nips + .as_array() + .ok_or_else(|| anyhow::anyhow!("supported_nips is not an array"))?; + nips.iter() + .map(|nip| { + nip.as_u64() + .ok_or_else(|| anyhow::anyhow!("supported_nips contains a non-integer")) + }) + .collect() +} + +async fn relay_supported_nips(relay_url: &str) -> Result>, Arc> { + let Some(http_url) = nip11_url(relay_url) else { + return Err(Arc::from("invalid relay URL")); + }; + + loop { + let now = Instant::now(); + let (generation, result) = { + let mut cache = CAPABILITY_CACHE.lock().await; + if let Some(entry) = cache.entries.get_mut(&http_url) { + entry.last_used = now; + (entry.generation, entry.result.clone()) + } else { + if cache.entries.len() >= NIP11_CACHE_CAPACITY { + let evict = cache + .entries + .iter() + .min_by_key(|(_, entry)| entry.last_used) + .map(|(url, _)| url.clone()); + if let Some(evict) = evict { + cache.entries.remove(&evict); + } + } + let generation = cache.next_generation; + cache.next_generation = cache.next_generation.wrapping_add(1); + let result = fetch_nip11(http_url.clone()).boxed().shared(); + cache.entries.insert( + http_url.clone(), + CacheEntry { + generation, + last_used: now, + result: result.clone(), + }, + ); + (generation, result) + } + }; + + let capability = result.await; + if capability.expires_at > Instant::now() { + return capability.supported_nips.clone(); + } + + let mut cache = CAPABILITY_CACHE.lock().await; + if cache + .entries + .get(&http_url) + .is_some_and(|entry| entry.generation == generation) + { + cache.entries.remove(&http_url); + } + } +} + +/// True when a filter object carries at least one NIP-91 `&x` query key. +fn carries_and_tags(filter: &Value) -> bool { + filter + .as_object() + .is_some_and(|object| object.keys().any(|key| key.starts_with('&'))) +} + +/// Rewrite one filter object into the NIP-91 compatibility form: every `&x` +/// value set is folded into the `#x` OR clause (creating it when absent) and +/// the `&x` key is removed. +/// +/// The fold, rather than a plain strip, is what the draft asks clients to send: +/// "Tag values used in `AND` by libraries and clients `MUST` include standard +/// `OR` tags [`#`] for compatibility with relays that do not support NIP-91." +/// Every event matching the AND clause carries all of its required values, so +/// the rewritten filter is a strict superset — it can never collapse into a +/// match-all the way dropping the keys outright would. +fn compatibility_filter(filter: &Value) -> anyhow::Result { + let Some(object) = filter.as_object() else { + return Ok(filter.clone()); + }; + + let mut compatible = serde_json::Map::new(); + let mut required: Vec<(String, &Vec)> = Vec::new(); + for (key, value) in object { + let Some(tag) = key.strip_prefix('&') else { + compatible.insert(key.clone(), value.clone()); + continue; + }; + if tag.len() != 1 || !tag.as_bytes()[0].is_ascii_alphabetic() { + anyhow::bail!("unindexed AND tag filter: {key}"); + } + let values = value + .as_array() + .ok_or_else(|| anyhow::anyhow!("{key} not an array"))?; + if values.is_empty() { + anyhow::bail!("{key} array must not be empty"); + } + required.push((format!("#{tag}"), values)); + } + + for (compat_key, values) in required { + let alternatives = compatible + .entry(compat_key.clone()) + .or_insert_with(|| Value::Array(Vec::new())) + .as_array_mut() + .ok_or_else(|| anyhow::anyhow!("{compat_key} not an array"))?; + for value in values { + if !alternatives.contains(value) { + alternatives.push(value.clone()); + } + } + } + + Ok(Value::Object(compatible)) +} + +/// Adapt an outbound REQ or NEG-OPEN filter to what the upstream relay +/// understands. Relays reject unrecognised filter keys, so NIP-91 `&x` clauses +/// are only forwarded when the upstream advertises NIP 91 in its NIP-11 +/// document; otherwise they are folded into the `#x` compatibility clause and +/// the exact AND semantics are re-applied locally by the caller. +/// +/// NIP-11 governs the NIP-91 decision alone. An unreachable or malformed +/// document means "unknown", which degrades to the compatibility filter rather +/// than failing — no other protocol support is inferred from it. +pub(crate) async fn outbound_filter_for_relay( + relay_url: &str, + filter: &Value, +) -> anyhow::Result { + let carries_and = match filter { + Value::Array(entries) => entries.iter().any(carries_and_tags), + other => carries_and_tags(other), + }; + if !carries_and { + return Ok(filter.clone()); + } + + // Reject malformed `&` clauses before the network round trip so the error + // never depends on upstream reachability. + let compatible = match filter { + Value::Array(entries) => Value::Array( + entries + .iter() + .map(compatibility_filter) + .collect::>>()?, + ), + other => compatibility_filter(other)?, + }; + + match relay_supported_nips(relay_url).await { + Ok(nips) if nips.contains(&91) => return Ok(filter.clone()), + Ok(_) => {} + Err(error) => { + tracing::warn!(url = relay_url, %error, "NIP-11 unavailable; assuming no NIP-91 support") + } + } + tracing::warn!( + url = relay_url, + "upstream does not advertise NIP-91; sending the `#` compatibility filter and applying AND tags locally" + ); + Ok(compatible) +} + /// Inbound messages on mesh connections carry at most one full event /// (<= max_event_size) plus envelope overhead. Cap at 2x (+ slack) so a /// malicious peer can't exploit tungstenite's 64 MiB default max message @@ -29,3 +303,203 @@ pub(crate) async fn connect_mesh( .await?; Ok(ws) } + +#[cfg(test)] +mod tests { + use super::*; + use std::sync::atomic::{AtomicUsize, Ordering}; + use tokio::io::{AsyncReadExt, AsyncWriteExt}; + + async fn nip11_server( + document: &'static str, + ) -> (String, Arc, tokio::task::JoinHandle<()>) { + let listener = tokio::net::TcpListener::bind("127.0.0.1:0").await.unwrap(); + let address = listener.local_addr().unwrap(); + let requests = Arc::new(AtomicUsize::new(0)); + let request_count = requests.clone(); + let task = tokio::spawn(async move { + while let Ok((mut stream, _)) = listener.accept().await { + request_count.fetch_add(1, Ordering::SeqCst); + let mut request = vec![0; 4096]; + let _ = stream.read(&mut request).await; + let response = format!( + "HTTP/1.1 200 OK\r\nContent-Type: application/nostr+json\r\nContent-Length: {}\r\nConnection: close\r\n\r\n{}", + document.len(), + document + ); + let _ = stream.write_all(response.as_bytes()).await; + } + }); + (format!("ws://{address}/relay"), requests, task) + } + + #[test] + fn nip11_uses_the_relay_websocket_uri() { + assert_eq!( + nip11_url("wss://relay.example:7447/nostr?network=main#fragment").as_deref(), + Some("https://relay.example:7447/nostr?network=main") + ); + assert_eq!( + nip11_url("ws://relay.example/").as_deref(), + Some("http://relay.example/") + ); + assert!(nip11_url("https://relay.example/").is_none()); + } + + #[tokio::test] + async fn advertised_nip91_preserves_and_fields() { + let (url, requests, server) = nip11_server(r#"{"supported_nips":[1,11,91]}"#).await; + let filter = serde_json::json!({"&t":["a","b"],"#t":["a","b","c"]}); + assert_eq!( + outbound_filter_for_relay(&url, &filter).await.unwrap(), + filter + ); + assert_eq!(requests.load(Ordering::SeqCst), 1); + server.abort(); + } + + #[tokio::test] + async fn unsupported_nip91_folds_and_values_into_the_compatibility_clause() { + let (url, requests, server) = nip11_server(r#"{"supported_nips":[1,11]}"#).await; + let filter = serde_json::json!({ + "&t":["a","b"], + "#t":["b","c"], + "kinds":[1] + }); + let checks = (0..8).map(|_| outbound_filter_for_relay(&url, &filter)); + let results = futures_util::future::join_all(checks).await; + for result in results { + assert_eq!( + result.unwrap(), + serde_json::json!({"#t":["b","c","a"],"kinds":[1]}) + ); + } + assert_eq!(requests.load(Ordering::SeqCst), 1); + + let _ = outbound_filter_for_relay(&url, &filter).await; + assert_eq!(requests.load(Ordering::SeqCst), 1); + server.abort(); + } + + /// An `&`-only filter must not degrade into a match-all: stripping the key + /// outright would subscribe the router to the whole remote firehose and + /// make `wok sync` reconcile the entire remote set. + #[tokio::test] + async fn and_only_filter_never_widens_to_the_remote_firehose() { + let (url, _, server) = nip11_server(r#"{"supported_nips":[1,11]}"#).await; + let filter = serde_json::json!({"&t":["a","b"],"kinds":[1],"limit":10}); + assert_eq!( + outbound_filter_for_relay(&url, &filter).await.unwrap(), + serde_json::json!({"#t":["a","b"],"kinds":[1],"limit":10}) + ); + server.abort(); + } + + #[tokio::test] + async fn every_and_key_gets_its_own_compatibility_clause() { + let (url, _, server) = nip11_server(r#"{"supported_nips":[1,11]}"#).await; + let author = "11".repeat(32); + let filter = serde_json::json!({"&t":["a"],"&p":[author.clone()]}); + assert_eq!( + outbound_filter_for_relay(&url, &filter).await.unwrap(), + serde_json::json!({"#t":["a"],"#p":[author]}) + ); + server.abort(); + } + + #[tokio::test] + async fn filter_arrays_are_rewritten_per_element() { + let (url, _, server) = nip11_server(r#"{"supported_nips":[1,11]}"#).await; + let filter = serde_json::json!([{"&t":["a"]},{"#e":["ff".repeat(32)]}]); + assert_eq!( + outbound_filter_for_relay(&url, &filter).await.unwrap(), + serde_json::json!([{"#t":["a"]},{"#e":["ff".repeat(32)]}]) + ); + server.abort(); + } + + #[tokio::test] + async fn malformed_nip11_uses_and_caches_the_compatibility_filter() { + let (url, requests, server) = nip11_server("not-json").await; + let filter = serde_json::json!({"&t":["a"]}); + let expected = serde_json::json!({"#t":["a"]}); + assert_eq!( + outbound_filter_for_relay(&url, &filter).await.unwrap(), + expected + ); + assert_eq!( + outbound_filter_for_relay(&url, &filter).await.unwrap(), + expected + ); + assert_eq!(requests.load(Ordering::SeqCst), 1); + server.abort(); + } + + /// NIP-11 governs the NIP-91 decision only. `wok sync` sends NEG-OPEN and + /// lets the relay answer for NIP-77 itself, so an upstream that omits 77 + /// from `supported_nips` must not be rejected before the websocket opens. + #[tokio::test] + async fn nip77_is_not_inferred_from_the_nip11_document() { + let (url, _, server) = nip11_server(r#"{"supported_nips":[1,11,91]}"#).await; + let filter = serde_json::json!({"&t":["a"],"#t":["a","b"]}); + assert_eq!( + outbound_filter_for_relay(&url, &filter).await.unwrap(), + filter + ); + server.abort(); + + let (url, _, server) = nip11_server(r#"{"supported_nips":[1,11]}"#).await; + assert_eq!( + outbound_filter_for_relay(&url, &filter).await.unwrap(), + serde_json::json!({"#t":["a","b"]}) + ); + server.abort(); + } + + #[tokio::test] + async fn an_unavailable_relay_falls_back_to_the_compatibility_filter() { + let filter = serde_json::json!({"&t":["a"]}); + assert_eq!( + outbound_filter_for_relay("not a relay URL", &filter) + .await + .unwrap(), + serde_json::json!({"#t":["a"]}) + ); + } + + #[tokio::test] + async fn malformed_and_clauses_fail_before_any_capability_lookup() { + for filter in [ + serde_json::json!({"&t":"a"}), + serde_json::json!({"&t":[]}), + serde_json::json!({"&foo":["a"]}), + serde_json::json!({"&t":["a"],"#t":"a"}), + ] { + assert!( + outbound_filter_for_relay("not a relay URL", &filter) + .await + .is_err(), + "{filter} should be rejected" + ); + } + } + + #[tokio::test] + async fn filters_without_nip91_fields_do_not_trigger_discovery() { + let filter = serde_json::json!({"#t":["a","b"]}); + assert_eq!( + outbound_filter_for_relay("not a relay URL", &filter) + .await + .unwrap(), + filter + ); + } + + #[test] + fn supported_nips_require_integer_entries() { + let nips = parse_supported_nips(&serde_json::json!({"supported_nips":[77,91]})).unwrap(); + assert!(nips.contains(&77)); + assert!(nips.contains(&91)); + assert!(parse_supported_nips(&serde_json::json!({"supported_nips":["91"]})).is_err()); + } +} diff --git a/crates/wok-cli/src/migrate.rs b/crates/wok-cli/src/migrate.rs index 15a124a..f9434a6 100644 --- a/crates/wok-cli/src/migrate.rs +++ b/crates/wok-cli/src/migrate.rs @@ -700,6 +700,7 @@ mod tests { &json!({"search":"migration fixture"}), 100, 3, + 16, |lev_id| search_hits.push(lev_id), ) .unwrap(); diff --git a/crates/wok-cli/src/reindex.rs b/crates/wok-cli/src/reindex.rs index 831cb2f..2a6a9c8 100644 --- a/crates/wok-cli/src/reindex.rs +++ b/crates/wok-cli/src/reindex.rs @@ -262,13 +262,20 @@ fn rebuild_negentropy(env: &Env) -> Result<(u64, u64)> { let txn = env.begin_ro()?; let filter: serde_json::Value = serde_json::from_str(filter)?; let mut records = Vec::new(); - wok_query::foreach_by_filter(&txn, &filter, u64::MAX, 64, |lev_id| { - if let Ok(Some(packed)) = wok_db::get_packed_ro(&txn, lev_id) { - if let Ok(packed) = wok_event::PackedEventView::new(&packed) { - records.push((packed.created_at(), packed.id().to_vec())); + wok_query::foreach_by_filter( + &txn, + &filter, + u64::MAX, + usize::MAX, + usize::MAX, + |lev_id| { + if let Ok(Some(packed)) = wok_db::get_packed_ro(&txn, lev_id) { + if let Ok(packed) = wok_event::PackedEventView::new(&packed) { + records.push((packed.created_at(), packed.id().to_vec())); + } } - } - })?; + }, + )?; records }; let mut txn = env.begin_rw()?; @@ -390,7 +397,7 @@ mod tests { assert!(check_integrity(&repaired.begin_ro().unwrap()).unwrap().ok()); let txn = repaired.begin_ro().unwrap(); let mut search_hits = Vec::new(); - wok_query::foreach_by_filter(&txn, &json!({"search":"repair me"}), 100, 3, |lev_id| { + wok_query::foreach_by_filter(&txn, &json!({"search":"repair me"}), 100, 3, 16, |lev_id| { search_hits.push(lev_id) }) .unwrap(); diff --git a/crates/wok-cli/src/router.rs b/crates/wok-cli/src/router.rs index 0e847f7..d48f279 100644 --- a/crates/wok-cli/src/router.rs +++ b/crates/wok-cli/src/router.rs @@ -12,7 +12,7 @@ use std::path::{Path, PathBuf}; use std::time::{Duration, Instant}; use tokio::sync::mpsc; use wok_db::{Decompressor, Env, EventToWrite}; -use wok_event::{parse_and_verify_event, PackedEventView}; +use wok_event::{parse_and_verify_event, EventLimits, PackedEventView}; use wok_query::NostrFilterGroup; use wok_relay::plugin::{PluginEventSifter, PluginResult}; use wok_relay::Config; @@ -43,8 +43,13 @@ pub struct StreamSpec { pub urls: Vec, } -fn compile_router_filter(name: &str, filter: &Value) -> Result { - let filter_group = NostrFilterGroup::from_value(filter, u64::MAX, 64) +fn compile_router_filter( + name: &str, + filter: &Value, + max_tags: usize, + max_and_entries: usize, +) -> Result { + let filter_group = NostrFilterGroup::from_value(filter, u64::MAX, max_tags, max_and_entries) .map_err(|error| anyhow::anyhow!("stream {name}: bad filter: {error}"))?; if filter_group.requires_content() { bail!("stream {name}: router filters do not support content search"); @@ -52,6 +57,16 @@ fn compile_router_filter(name: &str, filter: &Value) -> Result Ok(filter_group) } +fn router_event_matches( + filter_group: &NostrFilterGroup, + event: &Value, + limits: &EventLimits, +) -> bool { + wok_event::nostr_json_to_packed_event(event, limits) + .map(|packed| filter_group.does_match(packed.view())) + .unwrap_or(false) +} + enum Stmt { Open(String), Close, @@ -325,7 +340,8 @@ pub fn parse_router_config(text: &str) -> Result { .map(|_| ()) .and_then(|_| { // Validate the filter compiles. - compile_router_filter(name, &cfg.streams[name].filter).map(|_| ()) + compile_router_filter(name, &cfg.streams[name].filter, usize::MAX, usize::MAX) + .map(|_| ()) }) } @@ -521,6 +537,14 @@ pub async fn run_router(cfg: Config, router_path: PathBuf) -> Result<()> { queue_budget.add_permits(approx_bytes); if let Some(g) = groups.get_mut(&group) { if g.spec.dir != "up" { + let matches_filter = router_event_matches( + &g.filter_group, + &event, + &cfg.event_limits(), + ); + if !matches_filter { + continue; + } let cmd = g.spec.plugin_down.clone(); let ev = event.clone(); let res = tokio::task::block_in_place(|| { @@ -668,7 +692,12 @@ fn reconcile(cfg: &Config, router_cfg: &RouterConfig, groups: &mut HashMap fg, Err(e) => { tracing::error!("{e}; skipped"); @@ -747,7 +776,15 @@ async fn run_conn( .map_err(|_| anyhow::anyhow!("timeout"))? .map_err(|e| anyhow::anyhow!(e.to_string())) }; - let (ws, _) = match connect.await { + let remote_filter = async { + if dir == "down" || dir == "both" { + crate::mesh::outbound_filter_for_relay(&url, &filter).await + } else { + Ok(filter.clone()) + } + }; + let (connect_result, remote_filter) = tokio::join!(connect, remote_filter); + let (ws, _) = match connect_result { Ok(x) => x, Err(e) => { let _ = tx @@ -761,6 +798,20 @@ async fn run_conn( return; } }; + let remote_filter = match remote_filter { + Ok(f) => f, + Err(e) => { + let _ = tx + .send(ManagerMsg::Log { + group: group.clone(), + url: url.clone(), + text: format!("bad filter: {e}"), + }) + .await; + let _ = tx.send(ManagerMsg::Disconnected { group, url }).await; + return; + } + }; let (conn_tx, mut conn_rx) = mpsc::channel::(256); CONN_REGISTRY .lock() @@ -774,7 +825,7 @@ async fn run_conn( let (mut wtx, mut wrx) = ws.split(); if dir == "down" || dir == "both" { - let mut f = filter.clone(); + let mut f = remote_filter; f["limit"] = json!(0); let msg = json!(["REQ", "X", f]).to_string(); if wtx.send(Message::Text(msg.into())).await.is_err() { @@ -1054,4 +1105,40 @@ mod tests { "{\"authors\":[\"aaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaa\"],\"kinds\":[1]}" ); } + + #[test] + fn nip91_post_filters_fallback_results_from_older_relays() { + let filter = compile_router_filter( + "nip91", + &json!({"&t":["meme", "cat"], "#t":["meme", "cat", "black"]}), + 3, + 16, + ) + .unwrap(); + let event = |tags: Vec| { + json!({ + "id":"11".repeat(32), + "pubkey":"22".repeat(32), + "created_at":1, + "kind":1, + "tags":tags, + "content":"", + "sig":"33".repeat(64) + }) + }; + assert!(router_event_matches( + &filter, + &event(vec![ + json!(["t", "meme"]), + json!(["t", "cat"]), + json!(["t", "black"]) + ]), + &EventLimits::default(), + )); + assert!(!router_event_matches( + &filter, + &event(vec![json!(["t", "meme"]), json!(["t", "cat"])]), + &EventLimits::default(), + )); + } } diff --git a/crates/wok-compat/tests/cpp_export.rs b/crates/wok-compat/tests/cpp_export.rs index bcc1595..afe900d 100644 --- a/crates/wok-compat/tests/cpp_export.rs +++ b/crates/wok-compat/tests/cpp_export.rs @@ -87,7 +87,7 @@ fn rust_open_own_db_and_scan() { let txn = env.begin_ro().unwrap(); let mut decomp = Decompressor::new(); let mut found = false; - wok_query::foreach_by_filter(&txn, &json!({"kinds":[1]}), 500, 3, |lev| { + wok_query::foreach_by_filter(&txn, &json!({"kinds":[1]}), 500, 3, 16, |lev| { let json = event_json_owned(&txn, &mut decomp, lev, 65536).unwrap(); if json.contains("scan-me") { found = true; @@ -124,6 +124,7 @@ fn cpp_write_rust_read_query() { &json!({"kinds":[1], "#t":["compat"]}), 500, 3, + 16, |lev| { ids.push(event_json_owned(&txn, &mut decomp, lev, 65536).unwrap()); }, @@ -160,7 +161,7 @@ fn wok_replace_keeps_only_newest_event() { let txn = env.begin_ro().unwrap(); let mut decomp = Decompressor::new(); let mut found = Vec::new(); - wok_query::foreach_by_filter(&txn, &json!({"kinds":[0]}), 500, 3, |lev| { + wok_query::foreach_by_filter(&txn, &json!({"kinds":[0]}), 500, 3, 16, |lev| { found.push(event_json_owned(&txn, &mut decomp, lev, 65536).unwrap()); }) .unwrap(); @@ -182,7 +183,7 @@ fn wok_delete_removes_event_from_queries() { { let txn = env.begin_ro().unwrap(); let mut levs = Vec::new(); - wok_query::foreach_by_filter(&txn, &json!({"ids":[id]}), 500, 3, |lev| levs.push(lev)) + wok_query::foreach_by_filter(&txn, &json!({"ids":[id]}), 500, 3, 16, |lev| levs.push(lev)) .unwrap(); drop(txn); let mut txn = env.begin_rw().unwrap(); @@ -191,8 +192,10 @@ fn wok_delete_removes_event_from_queries() { } let txn = env.begin_ro().unwrap(); let mut found = Vec::new(); - wok_query::foreach_by_filter(&txn, &json!({"ids":[id]}), 500, 3, |lev| found.push(lev)) - .unwrap(); + wok_query::foreach_by_filter(&txn, &json!({"ids":[id]}), 500, 3, 16, |lev| { + found.push(lev) + }) + .unwrap(); assert!( found.is_empty(), "deleted event remained queryable: {found:?}" @@ -241,7 +244,7 @@ fn rust_scan_order_matches_created_at_desc() { let txn = env.begin_ro().unwrap(); let mut decomp = Decompressor::new(); let mut contents = Vec::new(); - wok_query::foreach_by_filter(&txn, &json!({"kinds":[1], "limit": 5}), 500, 3, |lev| { + wok_query::foreach_by_filter(&txn, &json!({"kinds":[1], "limit": 5}), 500, 3, 16, |lev| { let j = event_json_owned(&txn, &mut decomp, lev, 65536).unwrap(); contents.push(j); }) diff --git a/crates/wok-compat/tests/cpp_negentropy.rs b/crates/wok-compat/tests/cpp_negentropy.rs index f33a407..b60017e 100644 --- a/crates/wok-compat/tests/cpp_negentropy.rs +++ b/crates/wok-compat/tests/cpp_negentropy.rs @@ -120,7 +120,7 @@ fn strfry_refuses_wok_owned_tree_database() { let mut recs = Vec::new(); { let txn = env.begin_ro().unwrap(); - wok_query::foreach_by_filter(&txn, &json!({}), u64::MAX, 64, |lev| { + wok_query::foreach_by_filter(&txn, &json!({}), u64::MAX, 64, 16, |lev| { if let Ok(Some(buf)) = wok_db::get_packed_ro(&txn, lev) { let p = PackedEventView::new(&buf).unwrap(); recs.push((p.created_at(), p.id().to_vec())); diff --git a/crates/wok-compat/tests/e2e_transports.rs b/crates/wok-compat/tests/e2e_transports.rs index 899edaa..080706a 100644 --- a/crates/wok-compat/tests/e2e_transports.rs +++ b/crates/wok-compat/tests/e2e_transports.rs @@ -118,6 +118,92 @@ async fn ws_publish_and_subscribe() { handle.request_shutdown(); } +#[tokio::test(flavor = "multi_thread", worker_threads = 2)] +async fn nip91_historical_live_and_count_share_and_semantics() { + let dir = tempfile::tempdir().unwrap(); + let env = Env::open(dir.path(), EnvOptions::default()).unwrap(); + env.ensure_initialized().unwrap(); + let cfg = test_cfg(dir.path()); + let handle = wok_relay::start(env, cfg).unwrap(); + let listener = tokio::net::TcpListener::bind("127.0.0.1:0").await.unwrap(); + let addr = listener.local_addr().unwrap(); + let relay = handle.clone(); + tokio::spawn(async move { + let _ = wok_ws::serve_listener(relay, listener).await; + }); + tokio::time::sleep(Duration::from_millis(50)).await; + + let (mut ws, _) = tokio_tungstenite::connect_async(format!("ws://{addr}/")) + .await + .unwrap(); + let now = now_secs(); + for (content, tags) in [ + ("nip91-black", vec!["meme", "cat", "black"]), + ("nip91-white", vec!["meme", "cat", "white"]), + ("nip91-missing-or", vec!["meme", "cat"]), + ("nip91-missing-and", vec!["meme", "black"]), + ] { + let event = sign_event(json!({ + "created_at":now, + "kind":1, + "tags":tags.into_iter().map(|value| json!(["t", value])).collect::>(), + "content":content, + })); + ws.send(Message::Text(json!(["EVENT", event]).to_string().into())) + .await + .unwrap(); + let replies = recv_until(&mut ws, |text| text.contains("\"OK\"")).await; + assert!(replies.iter().any(|text| text.contains("true"))); + } + + let filter = json!({ + "kinds":[1], + "&t":["meme", "cat"], + "#t":["meme", "cat", "black", "white"] + }); + ws.send(Message::Text( + json!(["REQ", "nip91-live", filter.clone()]) + .to_string() + .into(), + )) + .await + .unwrap(); + let history = recv_until(&mut ws, |text| text.contains("EOSE")).await; + assert!(history.iter().any(|text| text.contains("nip91-black"))); + assert!(history.iter().any(|text| text.contains("nip91-white"))); + assert!(history + .iter() + .all(|text| !text.contains("nip91-missing-or") && !text.contains("nip91-missing-and"))); + + ws.send(Message::Text( + json!(["COUNT", "nip91-count", filter]).to_string().into(), + )) + .await + .unwrap(); + let count = recv_until(&mut ws, |text| text.contains("COUNT")).await; + assert!( + count.iter().any(|text| text.contains("\"count\":2")), + "expected exact NIP-91 count: {count:?}" + ); + assert!(count.iter().all(|text| !text.contains("\"hll\""))); + + let live = sign_event(json!({ + "created_at":now + 1, + "kind":1, + "tags":[["t","meme"],["t","cat"],["t","black"]], + "content":"nip91-live-match", + })); + ws.send(Message::Text(json!(["EVENT", live]).to_string().into())) + .await + .unwrap(); + let live_replies = recv_until(&mut ws, |text| text.contains("nip91-live-match")).await; + assert!(live_replies + .iter() + .any(|text| text.contains("nip91-live-match"))); + + handle.request_shutdown(); +} + #[tokio::test(flavor = "multi_thread", worker_threads = 2)] async fn nip62_vanish_is_immediate_and_blocks_rebroadcast() { let dir = tempfile::tempdir().unwrap(); @@ -327,7 +413,7 @@ async fn nip45_count_returns_mergeable_hll_for_canonical_tag_query() { let actual = body["hll"].as_str().expect("HLL response"); assert_eq!(actual.len(), 512); - let parsed_filter = NostrFilter::parse(&count_filter, 500, 3).unwrap(); + let parsed_filter = NostrFilter::parse(&count_filter, 500, 3, 16).unwrap(); let mut expected = HyperLogLog::for_filter(&parsed_filter).unwrap(); for key in &keys { let (pubkey, _) = key.x_only_public_key(); @@ -845,7 +931,7 @@ async fn nip11_http_document() { .any(|n| n == 1)); assert_eq!( client["supported_nips"], - json!([1, 9, 11, 13, 40, 45, 50, 62, 70, 77]) + json!([1, 9, 11, 13, 40, 45, 50, 62, 70, 77, 91]) ); assert_eq!(client["limitation"]["max_event_tags"], 2000); assert_eq!(client["limitation"]["created_at_lower_limit"], u64::MAX / 4); @@ -854,6 +940,8 @@ async fn nip11_http_document() { assert_eq!(client["limitation"]["max_total_events_per_req"], 2000); assert_eq!(client["limitation"]["min_pow_difficulty"], 20); assert_eq!(client["limitation"]["max_query_cost"], 1000); + assert_eq!(client["limitation"]["max_tags_per_filter"], 3); + assert_eq!(client["limitation"]["max_and_entries"], 16); assert_eq!( client["software"], "git+https://github.com/erskingardner/wok.git" diff --git a/crates/wok-compat/tests/nip_conformance.rs b/crates/wok-compat/tests/nip_conformance.rs index 30a60a5..30488d2 100644 --- a/crates/wok-compat/tests/nip_conformance.rs +++ b/crates/wok-compat/tests/nip_conformance.rs @@ -66,6 +66,7 @@ fn nip01_filter_kinds_since_until_limit() { &json!({"kinds":[1],"since":10,"until":20,"limit":5}), 500, 3, + 16, ) .unwrap(); assert_eq!(fg.filters[0].limit, 5); @@ -73,21 +74,22 @@ fn nip01_filter_kinds_since_until_limit() { #[test] fn nip01_filter_ids_and_kinds_use_event_field_grammar() { - assert!(NostrFilterGroup::from_value(&json!({"ids":["AA".repeat(32)]}), 500, 3).is_err()); + assert!(NostrFilterGroup::from_value(&json!({"ids":["AA".repeat(32)]}), 500, 3, 16).is_err()); assert!(NostrFilterGroup::from_value( &json!({"authors":[format!("0x{}", "11".repeat(32))]}), 500, - 3 + 3, + 16, ) .is_err()); - assert!(NostrFilterGroup::from_value(&json!({"kinds":[65536]}), 500, 3).is_err()); - assert!(NostrFilterGroup::from_value(&json!({"ids":[]}), 500, 3).is_err()); - assert!(NostrFilterGroup::from_value(&json!({"#1":["value"]}), 500, 3).is_err()); + assert!(NostrFilterGroup::from_value(&json!({"kinds":[65536]}), 500, 3, 16).is_err()); + assert!(NostrFilterGroup::from_value(&json!({"ids":[]}), 500, 3, 16).is_err()); + assert!(NostrFilterGroup::from_value(&json!({"#1":["value"]}), 500, 3, 16).is_err()); } #[test] fn nip01_unknown_filter_field_rejected() { - assert!(NostrFilterGroup::from_value(&json!({"foo":1}), 500, 3).is_err()); + assert!(NostrFilterGroup::from_value(&json!({"foo":1}), 500, 3, 16).is_err()); } #[test] @@ -161,20 +163,21 @@ fn nip50_search_filter_and_extensions() { }), 500, 3, + 16, ) .unwrap(); let search = filters.filters[0].search.as_ref().unwrap(); assert_eq!(search.terms, vec!["apps", "best", "nostr"]); assert_eq!(search.phrase, "best nostr apps"); assert_eq!(filters.filters[0].limit, 20); - assert!(NostrFilterGroup::from_value(&json!({"search":7}), 500, 3).is_err()); + assert!(NostrFilterGroup::from_value(&json!({"search":7}), 500, 3, 16).is_err()); } #[test] fn nip11_software_not_strfry_when_unconfigured() { let cfg = wok_relay::Config::default(); let nips = wok_relay::supported_nips(&cfg); - assert_eq!(nips, vec![1, 9, 11, 40, 45, 50, 62, 70, 77]); + assert_eq!(nips, vec![1, 9, 11, 40, 45, 50, 62, 70, 77, 91]); assert!(!nips.contains(&2), "client-side NIP-02 is not advertised"); assert!(!nips.contains(&4), "client-side NIP-04 is not advertised"); assert!(!nips.contains(&28), "client-side NIP-28 is not advertised"); @@ -194,7 +197,7 @@ fn nip01_invalid_signature_rejected() { #[test] fn nip01_filter_ids_are_exact_length() { - assert!(NostrFilterGroup::from_value(&json!({"ids":["aabb"]}), 500, 3).is_err()); + assert!(NostrFilterGroup::from_value(&json!({"ids":["aabb"]}), 500, 3, 16).is_err()); } #[test] @@ -215,6 +218,7 @@ fn nip01_duplicate_filters_still_parse() { ], 500, 3, + 16, ) .unwrap(); assert_eq!(fg.size(), 2); @@ -252,10 +256,28 @@ fn nip77_payload_hex_has_no_prefix_or_half_byte() { assert!(wok_event::from_hex_strict("1").is_err()); } +#[test] +fn nip91_and_tags_take_precedence_over_compatibility_or_tags() { + let filters = NostrFilterGroup::from_value( + &json!({ + "&t":["meme", "cat"], + "#t":["meme", "cat", "black", "white"] + }), + 500, + 3, + 16, + ) + .unwrap(); + let filter = &filters.filters[0]; + assert_eq!(filter.and_tags[&'t'].size(), 2); + assert_eq!(filter.tags[&'t'].size(), 2); + assert!(!filter.index_only_scans); +} + #[test] fn advertised_nips_are_subset_of_tested() { // NIP-86 coverage lives in e2e_transports.rs (management API e2e tests). - let tested = [1u64, 9, 11, 13, 40, 42, 45, 50, 59, 62, 70, 77, 86]; + let tested = [1u64, 9, 11, 13, 40, 42, 45, 50, 59, 62, 70, 77, 86, 91]; assert_eq!( wok_relay::RELAY_CAPABILITY_CATALOG .iter() diff --git a/crates/wok-negentropy/src/cache.rs b/crates/wok-negentropy/src/cache.rs index 5e2372e..b878e0d 100644 --- a/crates/wok-negentropy/src/cache.rs +++ b/crates/wok-negentropy/src/cache.rs @@ -9,10 +9,11 @@ use wok_query::NostrFilter; fn parse_negentropy_filter( filter_str: &str, max_tags_per_filter: usize, + max_and_entries: usize, ) -> Result { let value = wok_event::json::parse_strict(filter_str) .map_err(|error| NegError::msg(error.to_string()))?; - let filter = NostrFilter::parse(&value, u64::MAX, max_tags_per_filter) + let filter = NostrFilter::parse(&value, u64::MAX, max_tags_per_filter, max_and_entries) .map_err(|error| NegError::msg(error.to_string()))?; if filter.search.is_some() { return Err(NegError::msg( @@ -63,7 +64,7 @@ impl NegentropyFilterCache { lookup::foreach_negentropy_filter(txn, |id, filter_str| { // C++ tao::json::from_string throws on a corrupt filter row; // propagate instead of silently substituting a match-all {}. - match parse_negentropy_filter(filter_str, max_tags) { + match parse_negentropy_filter(filter_str, max_tags, usize::MAX) { Ok(f) => self.filters.push(FilterInfo { filter: f, tree_id: id, @@ -94,7 +95,7 @@ impl NegentropyFilterCache { let mut parse_err: Option = None; let max_tags = self.max_tags_per_filter; lookup::foreach_negentropy_filter_rw(txn, |id, filter_str| { - match parse_negentropy_filter(filter_str, max_tags) { + match parse_negentropy_filter(filter_str, max_tags, usize::MAX) { Ok(f) => self.filters.push(FilterInfo { filter: f, tree_id: id, @@ -198,9 +199,21 @@ mod tests { #[test] fn rejects_search_filters_without_event_content() { - let error = parse_negentropy_filter(r#"{"search":"nostr"}"#, 3).unwrap_err(); + let error = parse_negentropy_filter(r#"{"search":"nostr"}"#, 3, 16).unwrap_err(); assert!(error .to_string() .contains("negentropy filters do not support content search")); } + + #[test] + fn accepts_nip91_filters_with_overlap_normalization() { + let filter = parse_negentropy_filter( + r##"{"&t":["meme","cat"],"#t":["meme","cat","black"]}"##, + 3, + 16, + ) + .unwrap(); + assert_eq!(filter.and_tags[&'t'].size(), 2); + assert_eq!(filter.tags[&'t'].size(), 1); + } } diff --git a/crates/wok-query/src/filter.rs b/crates/wok-query/src/filter.rs index 4111f28..1152851 100644 --- a/crates/wok-query/src/filter.rs +++ b/crates/wok-query/src/filter.rs @@ -6,6 +6,8 @@ use std::collections::{BTreeMap, HashSet}; use wok_db::{parse_search_query, search_term_set, SearchQuery, SearchTermSet}; use wok_event::{from_lower_hex_exact, PackedEventView, MAX_INDEXED_TAG_VAL_SIZE}; +pub const DEFAULT_MAX_AND_ENTRIES: usize = 16; + #[derive(Debug, Clone)] pub struct FilterSetBytes { items: Vec>, @@ -68,6 +70,10 @@ impl FilterSetBytes { pub fn iter(&self) -> impl Iterator { self.items.iter().map(|v| v.as_slice()) } + + fn remove_all(&mut self, excluded: &Self) { + self.items.retain(|item| !excluded.does_match(item)); + } } #[derive(Debug, Clone)] @@ -118,6 +124,7 @@ pub struct NostrFilter { pub authors: Option, pub kinds: Option, pub tags: BTreeMap, + pub and_tags: BTreeMap, pub search: Option, pub since: u64, pub until: u64, @@ -130,6 +137,7 @@ impl NostrFilter { filter_obj: &Value, max_filter_limit: u64, max_tags_per_filter: usize, + max_and_entries: usize, ) -> Result { if !filter_obj.is_object() { return Err(QueryError::msg("provided filter is not an object")); @@ -139,13 +147,14 @@ impl NostrFilter { authors: None, kinds: None, tags: BTreeMap::new(), + and_tags: BTreeMap::new(), search: None, since: 0, until: u64::MAX, limit: u64::MAX, index_only_scans: false, }; - let mut num_major = 0u64; + let mut num_major_non_tag = 0u64; let obj = filter_obj.as_object().unwrap(); for (k, v) in obj { if k == "ids" { @@ -155,7 +164,7 @@ impl NostrFilter { if v.as_array().unwrap().is_empty() { return Err(QueryError::msg("ids array must not be empty")); } - num_major += 1; + num_major_non_tag += 1; f.ids = Some( FilterSetBytes::parse(v, true, 32, 32) .map_err(|e| QueryError::msg(format!("error parsing ids: {e}")))?, @@ -167,7 +176,7 @@ impl NostrFilter { if v.as_array().unwrap().is_empty() { return Err(QueryError::msg("authors array must not be empty")); } - num_major += 1; + num_major_non_tag += 1; f.authors = Some( FilterSetBytes::parse(v, true, 32, 32) .map_err(|e| QueryError::msg(format!("error parsing authors: {e}")))?, @@ -179,7 +188,7 @@ impl NostrFilter { if v.as_array().unwrap().is_empty() { return Err(QueryError::msg("kinds array must not be empty")); } - num_major += 1; + num_major_non_tag += 1; f.kinds = Some( FilterSetUint::parse(v) .map_err(|e| QueryError::msg(format!("error parsing kinds: {e}")))?, @@ -191,7 +200,6 @@ impl NostrFilter { if v.as_array().unwrap().is_empty() { return Err(QueryError::msg(format!("{k} array must not be empty"))); } - num_major += 1; if k.len() == 2 && k.as_bytes()[1].is_ascii_alphabetic() { let tag = k.as_bytes()[1] as char; let set = if tag == 'p' || tag == 'e' { @@ -204,6 +212,25 @@ impl NostrFilter { } else { return Err(QueryError::msg("unindexed tag filter")); } + } else if k.starts_with('&') { + if !v.is_array() { + return Err(QueryError::msg(format!("{k} not an array"))); + } + if v.as_array().unwrap().is_empty() { + return Err(QueryError::msg(format!("{k} array must not be empty"))); + } + if k.len() == 2 && k.as_bytes()[1].is_ascii_alphabetic() { + let tag = k.as_bytes()[1] as char; + let set = if tag == 'p' || tag == 'e' { + FilterSetBytes::parse(v, true, 32, 32) + } else { + FilterSetBytes::parse(v, false, 0, MAX_INDEXED_TAG_VAL_SIZE) + } + .map_err(|e| QueryError::msg(format!("error parsing {k}: {e}")))?; + f.and_tags.insert(tag, set); + } else { + return Err(QueryError::msg("unindexed AND tag filter")); + } } else if k == "since" { f.since = v .as_u64() @@ -224,19 +251,38 @@ impl NostrFilter { Some(parse_search_query(query).map_err(|error| { QueryError::msg(format!("error parsing search: {error}")) })?); - num_major += 1; + num_major_non_tag += 1; } else { return Err(QueryError::msg(format!("unrecognised filter item: {k}"))); } } - if f.tags.len() > max_tags_per_filter { + for (tag, required) in &f.and_tags { + if let Some(alternatives) = f.tags.get_mut(tag) { + alternatives.remove_all(required); + } + } + f.tags.retain(|_, alternatives| alternatives.size() > 0); + + let tag_key_count = f + .tags + .keys() + .chain(f.and_tags.keys()) + .copied() + .collect::>() + .len(); + if tag_key_count > max_tags_per_filter { return Err(QueryError::msg("too many tags in filter")); } + let and_entry_count: usize = f.and_tags.values().map(FilterSetBytes::size).sum(); + if and_entry_count > max_and_entries { + return Err(QueryError::msg("too many AND entries in filter")); + } if f.limit > max_filter_limit { f.limit = max_filter_limit; } - f.index_only_scans = - num_major <= 1 || (num_major == 2 && f.authors.is_some() && f.kinds.is_some()); + let num_major = num_major_non_tag + tag_key_count as u64; + f.index_only_scans = f.and_tags.is_empty() + && (num_major <= 1 || (num_major == 2 && f.authors.is_some() && f.kinds.is_some())); Ok(f) } @@ -283,6 +329,21 @@ impl NostrFilter { return false; } } + for (tag, required) in &self.and_tags { + for value in required.iter() { + let mut found = false; + ev.foreach_tag(|name, candidate| { + if name == *tag && candidate == value { + found = true; + return false; + } + true + }); + if !found { + return false; + } + } + } true } @@ -313,6 +374,7 @@ impl NostrFilter { && self.authors.is_none() && self.kinds.is_none() && self.tags.is_empty() + && self.and_tags.is_empty() && self.search.is_none() } } @@ -327,13 +389,14 @@ impl NostrFilterGroup { arr: &[Value], max_filter_limit: u64, max_tags: usize, + max_and_entries: usize, ) -> Result { if arr.len() < 3 { return Err(QueryError::msg("too small")); } let mut fg = Self::default(); for item in arr.iter().skip(2) { - fg.add_filter(item, max_filter_limit, max_tags)?; + fg.add_filter(item, max_filter_limit, max_tags, max_and_entries)?; } Ok(fg) } @@ -342,14 +405,15 @@ impl NostrFilterGroup { filter: &Value, max_filter_limit: u64, max_tags: usize, + max_and_entries: usize, ) -> Result { let mut fg = Self::default(); if filter.is_array() { for e in filter.as_array().unwrap() { - fg.add_filter(e, max_filter_limit, max_tags)?; + fg.add_filter(e, max_filter_limit, max_tags, max_and_entries)?; } } else { - fg.add_filter(filter, max_filter_limit, max_tags)?; + fg.add_filter(filter, max_filter_limit, max_tags, max_and_entries)?; } Ok(fg) } @@ -359,8 +423,9 @@ impl NostrFilterGroup { item: &Value, max_filter_limit: u64, max_tags: usize, + max_and_entries: usize, ) -> Result<(), QueryError> { - let f = NostrFilter::parse(item, max_filter_limit, max_tags)?; + let f = NostrFilter::parse(item, max_filter_limit, max_tags, max_and_entries)?; self.filters.push(f); Ok(()) } @@ -402,16 +467,28 @@ impl NostrFilterGroup { pub fn estimated_cost(&self, count_only: bool) -> u64 { let mut cost = 0u64; for filter in &self.filters { - let filter_cost = if filter.search.is_some() { + let and_surcharge = 25u64.saturating_mul( + filter + .and_tags + .values() + .map(|values| values.size() as u64) + .sum(), + ); + let base_cost = if filter.search.is_some() { 250 } else if let Some(ids) = &filter.ids { ids.size() as u64 + } else if !filter.and_tags.is_empty() { + // `DbScan` seeds a single cursor from one required tag value, + // whatever else the filter carries; `and_surcharge` covers the + // per-event verification the seed then requires. + 25 + } else if !filter.tags.is_empty() { + 25u64.saturating_mul(filter.tags.len() as u64) } else if filter.authors.is_some() && filter.kinds.is_some() { 5 } else if let Some(authors) = &filter.authors { 10u64.saturating_mul(authors.size() as u64) - } else if !filter.tags.is_empty() { - 25u64.saturating_mul(filter.tags.len() as u64) } else if let Some(kinds) = &filter.kinds { 100u64.saturating_mul(kinds.size() as u64) } else if filter.since != 0 || filter.until != u64::MAX { @@ -419,7 +496,7 @@ impl NostrFilterGroup { } else { 1_000 }; - cost = cost.saturating_add(filter_cost.max(1)); + cost = cost.saturating_add(base_cost.saturating_add(and_surcharge).max(1)); } if count_only { cost.saturating_mul(2) @@ -475,12 +552,22 @@ impl FilterValidator { .get(&'p') .map(|t| t.size() == 1) .unwrap_or(false); + let has_and_p = filter + .and_tags + .get(&'p') + .map(|t| t.size() == 1) + .unwrap_or(false); let has_e = filter .tags .get(&'e') .map(|t| t.size() == 1) .unwrap_or(false); - if !has_author && !has_p && !has_e { + let has_and_e = filter + .and_tags + .get(&'e') + .map(|t| t.size() == 1) + .unwrap_or(false); + if !has_author && !has_p && !has_e && !has_and_p && !has_and_e { return Err(QueryError::msg( "filter must have exactly one author, p tag, or e tag", )); @@ -518,16 +605,16 @@ mod tests { #[test] fn empty_filter_lists_are_rejected() { - assert!(NostrFilterGroup::from_value(&json!({"ids": []}), 500, 3).is_err()); - assert!(NostrFilterGroup::from_value(&json!({"authors": []}), 500, 3).is_err()); - assert!(NostrFilterGroup::from_value(&json!({"kinds": []}), 500, 3).is_err()); - assert!(NostrFilterGroup::from_value(&json!({"#p": []}), 500, 3).is_err()); + assert!(NostrFilterGroup::from_value(&json!({"ids": []}), 500, 3, 16).is_err()); + assert!(NostrFilterGroup::from_value(&json!({"authors": []}), 500, 3, 16).is_err()); + assert!(NostrFilterGroup::from_value(&json!({"kinds": []}), 500, 3, 16).is_err()); + assert!(NostrFilterGroup::from_value(&json!({"#p": []}), 500, 3, 16).is_err()); } #[test] fn kind_and_time_match() { - let f = - NostrFilter::parse(&json!({"kinds":[1], "since": 10, "until": 20}), 500, 3).unwrap(); + let f = NostrFilter::parse(&json!({"kinds":[1], "since": 10, "until": 20}), 500, 3, 16) + .unwrap(); let ev = packed_note(1, 2, 15, 1, &[]); assert!(f.does_match(ev.view())); let ev = packed_note(1, 2, 9, 1, &[]); @@ -538,32 +625,39 @@ mod tests { #[test] fn rejects_unknown_field() { - assert!(NostrFilter::parse(&json!({"foo": 1}), 500, 3).is_err()); + assert!(NostrFilter::parse(&json!({"foo": 1}), 500, 3, 16).is_err()); } #[test] fn ids_must_be_32_bytes() { - assert!(NostrFilter::parse(&json!({"ids":["aa"]}), 500, 3).is_err()); + assert!(NostrFilter::parse(&json!({"ids":["aa"]}), 500, 3, 16).is_err()); } #[test] fn index_only_heuristic() { - let f = NostrFilter::parse(&json!({"kinds":[1]}), 500, 3).unwrap(); + let f = NostrFilter::parse(&json!({"kinds":[1]}), 500, 3, 16).unwrap(); assert!(f.index_only_scans); - let f = - NostrFilter::parse(&json!({"authors":["11".repeat(32)], "kinds":[1]}), 500, 3).unwrap(); + let f = NostrFilter::parse( + &json!({"authors":["11".repeat(32)], "kinds":[1]}), + 500, + 3, + 16, + ) + .unwrap(); assert!(f.index_only_scans); - let f = NostrFilter::parse(&json!({"ids":["11".repeat(32)], "kinds":[1]}), 500, 3).unwrap(); + let f = + NostrFilter::parse(&json!({"ids":["11".repeat(32)], "kinds":[1]}), 500, 3, 16).unwrap(); assert!(!f.index_only_scans); } #[test] fn estimates_broad_queries_before_opening_a_cursor() { - let broad = NostrFilterGroup::from_value(&json!({}), 500, 3).unwrap(); + let broad = NostrFilterGroup::from_value(&json!({}), 500, 3, 16).unwrap(); let targeted = NostrFilterGroup::from_value( &json!({"ids":["11".repeat(32), "22".repeat(32)]}), 500, 3, + 16, ) .unwrap(); assert_eq!(broad.estimated_cost(false), 1_000); @@ -573,7 +667,7 @@ mod tests { #[test] fn search_matching_accepts_precomputed_event_terms() { - let filter = NostrFilter::parse(&json!({"search":"café nostr"}), 500, 3).unwrap(); + let filter = NostrFilter::parse(&json!({"search":"café nostr"}), 500, 3, 16).unwrap(); let event = packed_note(1, 2, 15, 1, &[]); let terms = search_term_set("A NOSTR relay with Café support"); @@ -581,4 +675,112 @@ mod tests { assert!(!filter.does_match_with_search_terms(event.view(), None)); assert!(filter.does_match_with_search_terms(event.view(), Some(&terms))); } + + #[test] + fn nip91_requires_every_and_value_and_preserves_residual_or() { + let filter = NostrFilter::parse( + &json!({ + "&t":["meme", "cat"], + "#t":["meme", "cat", "black", "white"] + }), + 500, + 3, + 16, + ) + .unwrap(); + assert!(!filter.index_only_scans); + assert_eq!(filter.and_tags[&'t'].size(), 2); + assert_eq!(filter.tags[&'t'].size(), 2); + + let matching = packed_note( + 1, + 2, + 15, + 1, + &[("t", b"meme"), ("t", b"cat"), ("t", b"black")], + ); + let missing_and = packed_note(2, 2, 15, 1, &[("t", b"meme"), ("t", b"black")]); + let missing_or = packed_note(3, 2, 15, 1, &[("t", b"meme"), ("t", b"cat")]); + assert!(filter.does_match(matching.view())); + assert!(!filter.does_match(missing_and.view())); + assert!(!filter.does_match(missing_or.view())); + } + + #[test] + fn nip91_complete_overlap_removes_the_or_clause() { + let filter = NostrFilter::parse( + &json!({"&t":["meme", "cat"], "#t":["cat", "meme"]}), + 500, + 3, + 16, + ) + .unwrap(); + assert!(filter.tags.is_empty()); + let event = packed_note(1, 2, 15, 1, &[("t", b"meme"), ("t", b"cat")]); + assert!(filter.does_match(event.view())); + assert!(!filter.is_full_db_query()); + } + + #[test] + fn nip91_validates_shape_hex_and_limits() { + assert!(NostrFilter::parse(&json!({"&t":[]}), 500, 3, 16).is_err()); + assert!(NostrFilter::parse(&json!({"&t":"meme"}), 500, 3, 16).is_err()); + assert!(NostrFilter::parse(&json!({"&1":["x"]}), 500, 3, 16).is_err()); + assert!(NostrFilter::parse(&json!({"&tag":["x"]}), 500, 3, 16).is_err()); + assert!(NostrFilter::parse(&json!({"&p":["AA".repeat(32)]}), 500, 3, 16).is_err()); + assert!(NostrFilter::parse(&json!({"&e":["aa"]}), 500, 3, 16).is_err()); + + let duplicate = + NostrFilter::parse(&json!({"&t":["a", "a", "b"], "#t":["a", "b"]}), 500, 1, 2).unwrap(); + assert_eq!(duplicate.and_tags[&'t'].size(), 2); + assert!(NostrFilter::parse(&json!({"&t":["a", "b", "c"]}), 500, 1, 2,).is_err()); + assert!( + NostrFilter::parse(&json!({"&t":["a"], "#p":["11".repeat(32)]}), 500, 1, 2,).is_err() + ); + } + + #[test] + fn nip91_entries_are_always_charged_to_query_cost() { + let tags = NostrFilterGroup::from_value(&json!({"&t":["a", "b"]}), 500, 3, 16).unwrap(); + assert_eq!(tags.estimated_cost(false), 75); + assert_eq!(tags.estimated_cost(true), 150); + + let id = NostrFilterGroup::from_value( + &json!({"ids":["11".repeat(32)], "&t":["a", "b"]}), + 500, + 3, + 16, + ) + .unwrap(); + assert_eq!(id.estimated_cost(false), 51); + + let author_and = NostrFilterGroup::from_value( + &json!({"authors":["11".repeat(32)], "&t":["a", "b"]}), + 500, + 3, + 16, + ) + .unwrap(); + assert_eq!(author_and.estimated_cost(false), 75); + + let author_kind_and = NostrFilterGroup::from_value( + &json!({"authors":["11".repeat(32)], "kinds":[1], "&t":["a"]}), + 500, + 3, + 16, + ) + .unwrap(); + assert_eq!(author_kind_and.estimated_cost(false), 50); + + // A required tag seeds one cursor, so the OR alternatives `DbScan` + // never opens are not charged; without it they are. + let or_and = + NostrFilterGroup::from_value(&json!({"#e":["11".repeat(32)], "&t":["a"]}), 500, 3, 16) + .unwrap(); + assert_eq!(or_and.estimated_cost(false), 50); + let or_only = + NostrFilterGroup::from_value(&json!({"#e":["11".repeat(32)], "#t":["a"]}), 500, 3, 16) + .unwrap(); + assert_eq!(or_only.estimated_cost(false), 50); + } } diff --git a/crates/wok-query/src/hll.rs b/crates/wok-query/src/hll.rs index 2934aba..0207ee6 100644 --- a/crates/wok-query/src/hll.rs +++ b/crates/wok-query/src/hll.rs @@ -55,7 +55,7 @@ impl HyperLogLog { /// values have ambiguous merge semantics, so callers deliberately omit HLL /// for those shapes. pub fn offset_for_filter(filter: &NostrFilter) -> Option { - if filter.tags.len() != 1 { + if !filter.and_tags.is_empty() || filter.tags.len() != 1 { return None; } let (tag, values) = filter.tags.first_key_value()?; @@ -91,7 +91,7 @@ mod tests { use serde_json::json; fn filter(value: serde_json::Value) -> NostrFilter { - NostrFilter::parse(&value, 500, 3).unwrap() + NostrFilter::parse(&value, 500, 3, 16).unwrap() } #[test] @@ -141,5 +141,10 @@ mod tests { "#p":["11".repeat(32)] }))) .is_none()); + assert!(offset_for_filter(&filter(json!({ + "&p":["11".repeat(32)], + "#p":["11".repeat(32)] + }))) + .is_none()); } } diff --git a/crates/wok-query/src/lib.rs b/crates/wok-query/src/lib.rs index 9264618..5b5e4e5 100644 --- a/crates/wok-query/src/lib.rs +++ b/crates/wok-query/src/lib.rs @@ -7,7 +7,9 @@ pub mod scan; pub mod scheduler; pub mod subid; -pub use filter::{dumb_match, FilterValidator, NostrFilter, NostrFilterGroup}; +pub use filter::{ + dumb_match, FilterValidator, NostrFilter, NostrFilterGroup, DEFAULT_MAX_AND_ENTRIES, +}; pub use hll::{offset_for_filter as nip45_hll_offset, HyperLogLog}; pub use monitor::{ActiveMonitors, Recipient}; pub use scan::{foreach_by_filter, DbQuery, DbScan}; diff --git a/crates/wok-query/src/monitor.rs b/crates/wok-query/src/monitor.rs index 3679a62..9bd3946 100644 --- a/crates/wok-query/src/monitor.rs +++ b/crates/wok-query/src/monitor.rs @@ -1,6 +1,9 @@ //! Live subscription inverted index matching `src/ActiveMonitors.h`. -use crate::subid::{SubId, Subscription}; +use crate::{ + subid::{SubId, Subscription}, + NostrFilter, +}; use std::collections::HashMap; use wok_db::SearchTermSet; use wok_event::PackedEventView; @@ -192,13 +195,9 @@ impl ActiveMonitors { a.copy_from_slice(authors.at(i)); self.all_authors.entry(a).or_default().push(item()); } - } else if !f.tags.is_empty() { - for (name, set) in &f.tags { - for i in 0..set.size() { - let mut spec = vec![*name as u8]; - spec.extend_from_slice(set.at(i)); - self.all_tags.entry(spec).or_default().push(item()); - } + } else if !f.tags.is_empty() || !f.and_tags.is_empty() { + for spec in tag_lookup_specs(f) { + self.all_tags.entry(spec).or_default().push(item()); } } else if let Some(kinds) = &f.kinds { for i in 0..kinds.size() { @@ -227,13 +226,9 @@ impl ActiveMonitors { a.copy_from_slice(authors.at(i)); remove_where(&mut self.all_authors, &a, &pred); } - } else if !f.tags.is_empty() { - for (name, set) in &f.tags { - for i in 0..set.size() { - let mut spec = vec![*name as u8]; - spec.extend_from_slice(set.at(i)); - remove_where(&mut self.all_tags, &spec, &pred); - } + } else if !f.tags.is_empty() || !f.and_tags.is_empty() { + for spec in tag_lookup_specs(f) { + remove_where(&mut self.all_tags, &spec, &pred); } } else if let Some(kinds) = &f.kinds { for i in 0..kinds.size() { @@ -246,6 +241,26 @@ impl ActiveMonitors { } } +fn tag_lookup_specs(filter: &NostrFilter) -> Vec> { + if let Some((name, values)) = filter.and_tags.first_key_value() { + let mut spec = vec![*name as u8]; + spec.extend_from_slice(values.at(0)); + return vec![spec]; + } + + filter + .tags + .iter() + .flat_map(|(name, values)| { + (0..values.size()).map(move |i| { + let mut spec = vec![*name as u8]; + spec.extend_from_slice(values.at(i)); + spec + }) + }) + .collect() +} + fn remove_where( map: &mut HashMap>, key: &K, @@ -258,3 +273,49 @@ fn remove_where( } } } + +#[cfg(test)] +mod tests { + use super::*; + use crate::NostrFilterGroup; + use serde_json::json; + use wok_event::{PackedEventBuilder, PackedEventTagBuilder}; + + fn event(id: u8, tags: &[&[u8]]) -> wok_event::PackedEvent { + let mut builder = PackedEventTagBuilder::default(); + for value in tags { + builder.add('t', value).unwrap(); + } + PackedEventBuilder::build(&[id; 32], &[2; 32], id as u64, 1, 0, &builder).unwrap() + } + + #[test] + fn nip91_live_monitor_uses_a_seed_then_checks_every_constraint() { + let group = NostrFilterGroup::from_value( + &json!({ + "&t":["meme", "cat"], + "#t":["meme", "cat", "black"] + }), + 500, + 3, + 16, + ) + .unwrap(); + let sub_id = SubId::new("nip91").unwrap(); + let mut sub = Subscription::new(1, sub_id.clone(), group, false); + sub.latest_event_id = 0; + let mut monitors = ActiveMonitors::new(10); + assert!(monitors.add_sub(sub, 0)); + + let matching = event(1, &[b"meme", b"cat", b"black"]); + assert_eq!(monitors.process(1, matching.view(), None).len(), 1); + let missing_or = event(2, &[b"meme", b"cat"]); + assert!(monitors.process(2, missing_or.view(), None).is_empty()); + let missing_and = event(3, &[b"meme", b"black"]); + assert!(monitors.process(3, missing_and.view(), None).is_empty()); + + monitors.remove_sub(1, &sub_id); + let later = event(4, &[b"meme", b"cat", b"black"]); + assert!(monitors.process(4, later.view(), None).is_empty()); + } +} diff --git a/crates/wok-query/src/scan.rs b/crates/wok-query/src/scan.rs index 13876aa..656535f 100644 --- a/crates/wok-query/src/scan.rs +++ b/crates/wok-query/src/scan.rs @@ -482,13 +482,19 @@ enum QueryScanner { } impl DbScan { + /// Number of index cursors this scan opens. A `&` (AND) tag filter always + /// seeds exactly one, whatever the size of its value set. + pub fn cursor_count(&self) -> usize { + self.cursors.len() + } + pub fn new(f: &NostrFilter, txn: &RoTxn<'_>) -> Self { let dbis = txn.env().dbis(); let mut index_only = f.index_only_scans; let mut cursors = Vec::new(); let (index_dbi, desc) = if f.ids.is_some() { (dbis.event_id, "ID") - } else if !f.tags.is_empty() { + } else if !f.tags.is_empty() || !f.and_tags.is_empty() { (dbis.event_tag, "Tag") } else if f.authors.is_some() && f.kinds.is_some() @@ -519,14 +525,24 @@ impl DbScan { outstanding: 0, }); } - } else if !f.tags.is_empty() { - let (tag_name, filter_set) = f - .tags - .iter() - .min_by_key(|(_, s)| s.size()) - .map(|(k, v)| (*k, v)) - .unwrap(); - for i in 0..filter_set.size() { + } else if !f.tags.is_empty() || !f.and_tags.is_empty() { + // A required (`&`) tag seeds a single cursor because every listed + // value must be present, while an OR tag needs one cursor per + // alternative, so a required tag always wins on cursor count when + // one is available. Index cardinality is not known at planning + // time, so cursor count is the only comparable unit: a narrow + // `#e:[]` can still lose to a broad `&t:["nostr"]`. That costs + // extra scanning, never correctness, because a non-empty `and_tags` + // forces every candidate through full-event verification. + let and_choice = f.and_tags.iter().min_by_key(|(_, values)| values.size()); + let or_choice = f.tags.iter().min_by_key(|(_, values)| values.size()); + let (tag_name, filter_set, from_and) = match (and_choice, or_choice) { + (Some((tag, values)), _) => (*tag, values, true), + (None, Some((tag, values))) => (*tag, values, false), + (None, None) => unreachable!(), + }; + let cursor_count = if from_and { 1 } else { filter_set.size() }; + for i in 0..cursor_count { let mut search = vec![tag_name as u8]; search.extend_from_slice(filter_set.at(i)); let mut resume = search.clone(); @@ -957,12 +973,14 @@ pub fn foreach_by_filter( filter: &serde_json::Value, max_limit: u64, max_tags: usize, + max_and_entries: usize, mut cb: F, ) -> Result<(), crate::subid::QueryError> where F: FnMut(u64), { - let fg = crate::filter::NostrFilterGroup::from_value(filter, max_limit, max_tags)?; + let fg = + crate::filter::NostrFilterGroup::from_value(filter, max_limit, max_tags, max_and_entries)?; let sub = Subscription::new(1, crate::subid::SubId::new(".").unwrap(), fg, false); let mut q = DbQuery::new(sub, 0, 0); q.process(txn, |_, lev| cb(lev), u64::MAX) diff --git a/crates/wok-query/tests/filter_prop.rs b/crates/wok-query/tests/filter_prop.rs index 581ceae..d287985 100644 --- a/crates/wok-query/tests/filter_prop.rs +++ b/crates/wok-query/tests/filter_prop.rs @@ -3,6 +3,18 @@ use serde_json::json; use wok_event::{PackedEventBuilder, PackedEventTagBuilder, PackedEventView}; use wok_query::{dumb_match, NostrFilter}; +fn tag_value(value: u8) -> String { + format!("v{value}") +} + +fn packed_tags(values: &[u8]) -> wok_event::PackedEvent { + let mut tags = PackedEventTagBuilder::default(); + for value in values { + tags.add('t', tag_value(*value).as_bytes()).unwrap(); + } + PackedEventBuilder::build(&[1; 32], &[2; 32], 50, 1, 0, &tags).unwrap() +} + fn packed( id: u8, pk: u8, @@ -69,11 +81,47 @@ proptest! { "kinds": [filter_kind], "since": since, "until": until, - }), 500, 3).unwrap(); + }), 500, 3, 16).unwrap(); let a = f.does_match(ev.view()); let b = naive_match(&f, ev.view()); let c = dumb_match(&f, ev.view()); prop_assert_eq!(a, b); prop_assert_eq!(a, c); } + + + #[test] + fn nip91_matcher_agrees_with_wire_semantics( + event_values in prop::collection::vec(0u8..6, 0..8), + and_values in prop::collection::vec(0u8..6, 1..6), + or_values in prop::collection::vec(0u8..6, 0..6), + ) { + let event = packed_tags(&event_values); + let and_json: Vec<_> = and_values.iter().copied().map(tag_value).collect(); + let or_json: Vec<_> = and_values + .iter() + .chain(or_values.iter()) + .copied() + .map(tag_value) + .collect(); + let filter = NostrFilter::parse( + &json!({"&t":and_json, "#t":or_json}), + 500, + 3, + 16, + ) + .unwrap(); + + let required: std::collections::HashSet<_> = and_values.iter().copied().collect(); + let alternatives: std::collections::HashSet<_> = or_values + .iter() + .copied() + .filter(|value| !required.contains(value)) + .collect(); + let present: std::collections::HashSet<_> = event_values.iter().copied().collect(); + let expected = required.iter().all(|value| present.contains(value)) + && (alternatives.is_empty() + || alternatives.iter().any(|value| present.contains(value))); + prop_assert_eq!(filter.does_match(event.view()), expected); + } } diff --git a/crates/wok-query/tests/scan_kinds.rs b/crates/wok-query/tests/scan_kinds.rs index 6003a2a..9c7b56c 100644 --- a/crates/wok-query/tests/scan_kinds.rs +++ b/crates/wok-query/tests/scan_kinds.rs @@ -4,7 +4,7 @@ use tempfile::TempDir; use wok_db::{write_events, Env, EnvOptions, EventToWrite, NoopNegentropy}; use wok_event::{parse_and_verify_event, EventLimits}; use wok_query::{ - foreach_by_filter, DbQuery, NostrFilterGroup, QueryScheduler, SubId, Subscription, + foreach_by_filter, DbQuery, DbScan, NostrFilterGroup, QueryScheduler, SubId, Subscription, }; fn sign(kind: u64, content: &str, created: u64) -> (Vec, String, String) { @@ -35,6 +35,25 @@ fn sign_with_key( (parsed.packed.into_bytes(), parsed.json, hex::encode(id)) } +fn sign_with_tags(content: &str, created: u64, tags: &[&str]) -> (Vec, String, String) { + let mut rng = rand::thread_rng(); + let kp = Keypair::new(SECP256K1, &mut rng); + let (xonly, _) = kp.x_only_public_key(); + let mut ev = json!({ + "created_at": created, + "kind": 1, + "tags": tags.iter().map(|value| json!(["t", value])).collect::>(), + "content": content, + "pubkey": hex::encode(xonly.serialize()), + }); + let id = wok_event::event_id_hash(&ev).unwrap(); + ev["id"] = json!(hex::encode(id)); + let sig = SECP256K1.sign_schnorr(&id, &kp); + ev["sig"] = json!(hex::encode(sig.as_ref())); + let parsed = parse_and_verify_event(&ev, &EventLimits::default(), None, true, false).unwrap(); + (parsed.packed.into_bytes(), parsed.json, hex::encode(id)) +} + #[test] fn scan_kinds_and_limit() { let tmp = TempDir::new().unwrap(); @@ -56,15 +75,73 @@ fn scan_kinds_and_limit() { } let txn = env.begin_ro().unwrap(); let mut got = Vec::new(); - foreach_by_filter(&txn, &json!({"kinds":[1], "limit": 5}), 500, 3, |lev| { + foreach_by_filter(&txn, &json!({"kinds":[1], "limit": 5}), 500, 3, 16, |lev| { got.push(lev); }) .unwrap(); assert_eq!(got.len(), 5); - let fg = NostrFilterGroup::from_value(&json!({"kinds":[0]}), 500, 3).unwrap(); + let fg = NostrFilterGroup::from_value(&json!({"kinds":[0]}), 500, 3, 16).unwrap(); assert_eq!(fg.size(), 1); } +#[test] +fn nip91_historical_scan_uses_one_seed_and_verifies_the_full_event() { + let tmp = TempDir::new().unwrap(); + let env = Env::open(tmp.path(), EnvOptions::default()).unwrap(); + env.ensure_initialized().unwrap(); + let inputs = [ + ("black", 100, &["meme", "cat", "black"][..]), + ("white", 200, &["meme", "cat", "white"][..]), + ("missing-or", 300, &["meme", "cat"][..]), + ("missing-and", 400, &["meme", "black"][..]), + ]; + let mut events: Vec<_> = inputs + .iter() + .map(|(content, created, tags)| { + let (packed, json, _) = sign_with_tags(content, *created, tags); + EventToWrite::new(packed, json) + }) + .collect(); + let mut write = env.begin_rw().unwrap(); + write_events(&mut write, &mut NoopNegentropy, &mut events, false).unwrap(); + write.commit().unwrap(); + + let expected = vec![events[1].lev_id, events[0].lev_id]; + let txn = env.begin_ro().unwrap(); + let mut hits = Vec::new(); + foreach_by_filter( + &txn, + &json!({ + "&t":["meme", "cat"], + "#t":["meme", "cat", "black", "white"] + }), + 500, + 3, + 16, + |lev| hits.push(lev), + ) + .unwrap(); + assert_eq!(hits, expected); + + let group = NostrFilterGroup::from_value( + &json!({ + "&t":["meme", "cat"], + "#t":["meme", "cat", "black", "white"] + }), + 500, + 3, + 16, + ) + .unwrap(); + let mut count = DbQuery::new( + Subscription::new(1, SubId::new("nip91-count").unwrap(), group, true), + 0, + 100, + ); + assert!(count.process(&txn, |_, _| {}, u64::MAX).unwrap()); + assert_eq!(count.sent_count(), 2); +} + #[test] fn request_wide_limit_caps_deduplicated_multi_filter_results_but_not_count() { let tmp = TempDir::new().unwrap(); @@ -82,7 +159,7 @@ fn request_wide_limit_caps_deduplicated_multi_filter_results_but_not_count() { write.commit().unwrap(); let request = json!(["REQ", "multi", {"kinds":[1], "limit":5}, {"kinds":[0], "limit":5}]); - let group = NostrFilterGroup::from_req(request.as_array().unwrap(), 500, 3).unwrap(); + let group = NostrFilterGroup::from_req(request.as_array().unwrap(), 500, 3, 16).unwrap(); let txn = env.begin_ro().unwrap(); let mut delivery = DbQuery::new( @@ -127,7 +204,7 @@ fn count_dedup_budget_caps_multi_filter_scans_as_limited() { // Two broad filters, each matching 5 events; budget below the total. let request = json!(["COUNT", "c", {"kinds":[1]}, {"kinds":[3]}]); - let group = NostrFilterGroup::from_req(request.as_array().unwrap(), 500, 3).unwrap(); + let group = NostrFilterGroup::from_req(request.as_array().unwrap(), 500, 3, 16).unwrap(); let txn = env.begin_ro().unwrap(); let mut count = DbQuery::new( Subscription::new(1, SubId::new("c").unwrap(), group, true), @@ -164,6 +241,7 @@ fn deep_author_kind_pages_are_complete_and_non_overlapping() { &json!({"authors":[pubkey], "kinds":[1], "limit":500}), 500, 3, + 16, |lev| first.push(lev), ) .unwrap(); @@ -186,6 +264,7 @@ fn deep_author_kind_pages_are_complete_and_non_overlapping() { }), 500, 3, + 16, |lev| second.push(lev), ) .unwrap(); @@ -195,7 +274,7 @@ fn deep_author_kind_pages_are_complete_and_non_overlapping() { let request = json!(["REQ", "scheduled", { "authors":[pubkey], "kinds":[1], "limit":500 }]); - let group = NostrFilterGroup::from_req(request.as_array().unwrap(), 500, 3).unwrap(); + let group = NostrFilterGroup::from_req(request.as_array().unwrap(), 500, 3, 16).unwrap(); let mut scheduler = QueryScheduler::new(8, 2_000, 0); assert!(scheduler .add_sub( @@ -221,3 +300,37 @@ fn deep_author_kind_pages_are_complete_and_non_overlapping() { assert_eq!(scheduled.len(), 500); assert_eq!(completions, vec![500]); } + +/// A required (`&`) tag needs one cursor no matter how many values it lists, +/// while an OR tag needs one per alternative, so the required tag seeds even +/// when a smaller OR set is present. Regression guard for a seed comparison +/// that mixed value counts with cursor counts. +#[test] +fn nip91_and_tag_seeds_a_single_cursor() { + let tmp = TempDir::new().unwrap(); + let env = Env::open(tmp.path(), EnvOptions::default()).unwrap(); + env.ensure_initialized().unwrap(); + let txn = env.begin_ro().unwrap(); + + let seed_count = |filter: serde_json::Value| { + let group = NostrFilterGroup::from_value(&filter, 500, 3, 16).unwrap(); + DbScan::new(&group.filters[0], &txn).cursor_count() + }; + + assert_eq!(seed_count(json!({"&t":["a", "b", "c"]})), 1); + // One narrow OR alternative does not outbid the required tag. + assert_eq!( + seed_count(json!({"&t":["a", "b", "c"], "#e":["11".repeat(32)]})), + 1 + ); + // Several required keys: the smallest set decides which one seeds, and it + // is still a single cursor. + assert_eq!(seed_count(json!({"&t":["a", "b"], "&d":["x"]})), 1); + // Without a required tag the smallest OR set seeds, one cursor per value. + assert_eq!( + seed_count( + json!({"#t":["a", "b"], "#e":["11".repeat(32), "22".repeat(32), "33".repeat(32)]}) + ), + 2 + ); +} diff --git a/crates/wok-query/tests/search.rs b/crates/wok-query/tests/search.rs index c5436a9..eb258c7 100644 --- a/crates/wok-query/tests/search.rs +++ b/crates/wok-query/tests/search.rs @@ -65,6 +65,7 @@ fn search_is_ranked_then_limited_and_honors_structured_filters() { &json!({"search":"nostr search", "kinds":[1], "limit":1}), 100, 3, + 16, |lev_id| hits.push(lev_id), ) .unwrap(); @@ -92,6 +93,7 @@ fn search_ignores_extensions_and_is_unicode_case_insensitive() { &json!({"search":"café domain:example.com include:spam"}), 100, 3, + 16, |lev_id| hits.push(lev_id), ) .unwrap(); @@ -121,6 +123,7 @@ fn delete_removes_search_postings_and_missing_index_is_backfilled() { &json!({"search":"backfill sentinel"}), 100, 3, + 16, |hit| hits.push(hit), ) .unwrap(); @@ -137,6 +140,7 @@ fn delete_removes_search_postings_and_missing_index_is_backfilled() { &json!({"search":"backfill sentinel"}), 100, 3, + 16, |hit| hits.push(hit), ) .unwrap(); @@ -167,6 +171,7 @@ fn multiple_search_filters_are_merged_by_quality() { ]), 100, 3, + 16, |lev_id| hits.push(lev_id), ) .unwrap(); @@ -183,6 +188,7 @@ fn multiple_search_filters_are_merged_by_quality() { ]), 100, 3, + 16, ) .unwrap(); let txn = env.begin_ro().unwrap(); diff --git a/crates/wok-relay/src/capabilities.rs b/crates/wok-relay/src/capabilities.rs index 06b318c..635affb 100644 --- a/crates/wok-relay/src/capabilities.rs +++ b/crates/wok-relay/src/capabilities.rs @@ -23,6 +23,7 @@ pub enum CapabilityCondition { NegentropyEnabled, Nip59Safe, Nip62Enabled, + Nip91Enabled, PowRequired, } @@ -44,6 +45,9 @@ impl CapabilityCondition { && cfg.relay.auth.restrict_read_to_involved_pubkey } Self::Nip62Enabled => cfg.relay.nip62.enabled, + Self::Nip91Enabled => { + cfg.relay.max_tags_per_filter > 0 && cfg.relay.max_and_entries > 0 + } Self::PowRequired => cfg.relay.abuse.enabled && cfg.relay.abuse.min_pow_difficulty > 0, } } @@ -116,6 +120,11 @@ pub const RELAY_CAPABILITY_CATALOG: &[RelayCapability] = &[ name: "Relay management API", enabled_when: CapabilityCondition::AdminEnabled, }, + RelayCapability { + nip: 91, + name: "AND operator in filters", + enabled_when: CapabilityCondition::Nip91Enabled, + }, ]; pub fn relay_capabilities(cfg: &Config) -> Vec { @@ -147,16 +156,25 @@ mod tests { #[test] fn conditional_capabilities_follow_runtime_configuration() { let mut cfg = Config::default(); - assert_eq!(supported_nips(&cfg), vec![1, 9, 11, 40, 45, 50, 62, 70, 77]); + assert_eq!( + supported_nips(&cfg), + vec![1, 9, 11, 40, 45, 50, 62, 70, 77, 91] + ); cfg.relay.auth.service_url = "wss://relay.example.com/".into(); cfg.relay.max_filter_limit_count = 0; cfg.relay.negentropy_enabled = false; cfg.events.ephemeral_persistence = EphemeralPersistence::Ttl; - assert_eq!(supported_nips(&cfg), vec![1, 9, 11, 40, 42, 50, 62, 70]); + assert_eq!(supported_nips(&cfg), vec![1, 9, 11, 40, 42, 50, 62, 70, 91]); cfg.relay.abuse.min_pow_difficulty = 20; - assert_eq!(supported_nips(&cfg), vec![1, 9, 11, 13, 40, 42, 50, 62, 70]); + assert_eq!( + supported_nips(&cfg), + vec![1, 9, 11, 13, 40, 42, 50, 62, 70, 91] + ); + + cfg.relay.max_and_entries = 0; + assert!(!supported_nips(&cfg).contains(&91)); } #[test] diff --git a/crates/wok-relay/src/config.rs b/crates/wok-relay/src/config.rs index 495f7ee..6575da2 100644 --- a/crates/wok-relay/src/config.rs +++ b/crates/wok-relay/src/config.rs @@ -192,6 +192,8 @@ pub struct RelayConfig { pub query_timeslice_budget_us: u64, pub max_filter_limit: u64, pub max_tags_per_filter: usize, + #[serde(default = "default_max_and_entries")] + pub max_and_entries: usize, pub max_filter_limit_count: u64, pub max_total_events_per_req: u64, pub max_subs_per_connection: usize, @@ -364,6 +366,7 @@ impl Default for Config { query_timeslice_budget_us: 10000, max_filter_limit: 500, max_tags_per_filter: 3, + max_and_entries: default_max_and_entries(), max_filter_limit_count: 1_000_000, max_total_events_per_req: 2_000, max_subs_per_connection: 200, @@ -430,6 +433,10 @@ impl Default for Config { } } +fn default_max_and_entries() -> usize { + wok_query::DEFAULT_MAX_AND_ENTRIES +} + impl From for TomlConfig { fn from(config: Config) -> Self { Self { diff --git a/crates/wok-relay/src/restrict.rs b/crates/wok-relay/src/restrict.rs index da104b1..e1321be 100644 --- a/crates/wok-relay/src/restrict.rs +++ b/crates/wok-relay/src/restrict.rs @@ -69,7 +69,12 @@ impl ReadRestrictor { .get(&'p') .map(|t| t.size() > 0 && (0..t.size()).all(|i| t.at(i) == pk)) .unwrap_or(false); - if !author_scoped && !p_scoped { + let and_p_scoped = f + .and_tags + .get(&'p') + .map(|t| t.size() > 0 && (0..t.size()).all(|i| t.at(i) == pk)) + .unwrap_or(false); + if !author_scoped && !p_scoped && !and_p_scoped { return false; } } @@ -111,9 +116,9 @@ mod tests { #[test] fn fully_restricted_kinds() { let r = ReadRestrictor::new(vec![4, 1059], true); - let fg = NostrFilterGroup::from_value(&json!({"kinds":[4]}), 500, 3).unwrap(); + let fg = NostrFilterGroup::from_value(&json!({"kinds":[4]}), 500, 3, 16).unwrap(); assert!(r.is_filter_group_fully_restricted(&fg)); - let fg = NostrFilterGroup::from_value(&json!({"kinds":[1,4]}), 500, 3).unwrap(); + let fg = NostrFilterGroup::from_value(&json!({"kinds":[1,4]}), 500, 3, 16).unwrap(); assert!(!r.is_filter_group_fully_restricted(&fg)); } @@ -132,13 +137,55 @@ mod tests { #[test] fn count_without_kinds_cannot_leak_restricted_population() { let r = ReadRestrictor::new(vec![4, 1059], true); - let broad = NostrFilterGroup::from_value(&json!({}), 500, 3).unwrap(); + let broad = NostrFilterGroup::from_value(&json!({}), 500, 3, 16).unwrap(); assert!(!r.is_filter_allowed_to_count(&broad, None)); assert!(!r.is_filter_allowed_to_count(&broad, Some(&[9u8; 32]))); let scoped = - NostrFilterGroup::from_value(&json!({"authors":[hex::encode([9u8; 32])]}), 500, 3) + NostrFilterGroup::from_value(&json!({"authors":[hex::encode([9u8; 32])]}), 500, 3, 16) .unwrap(); assert!(r.is_filter_allowed_to_count(&scoped, Some(&[9u8; 32]))); } + + #[test] + fn nip91_required_p_tag_safely_scopes_restricted_count() { + let r = ReadRestrictor::new(vec![4], true); + let authed = [9u8; 32]; + let other = [8u8; 32]; + let scoped = NostrFilterGroup::from_value( + &json!({ + "kinds":[4], + "&p":[hex::encode(authed)], + "#p":[hex::encode(authed)] + }), + 500, + 3, + 16, + ) + .unwrap(); + assert!(r.is_filter_allowed_to_count(&scoped, Some(&authed))); + assert!(!r.is_filter_allowed_to_count(&scoped, Some(&[7u8; 32]))); + + let ambiguous_and = NostrFilterGroup::from_value( + &json!({ + "kinds":[4], + "&p":[hex::encode(authed), hex::encode(other)], + "#p":[hex::encode(authed), hex::encode(other)] + }), + 500, + 3, + 16, + ) + .unwrap(); + assert!(!r.is_filter_allowed_to_count(&ambiguous_and, Some(&authed))); + + let unsafe_or = NostrFilterGroup::from_value( + &json!({"kinds":[4], "#p":[hex::encode(authed), hex::encode(other)]}), + 500, + 3, + 16, + ) + .unwrap(); + assert!(!r.is_filter_allowed_to_count(&unsafe_or, Some(&authed))); + } } diff --git a/crates/wok-relay/src/server.rs b/crates/wok-relay/src/server.rs index d486ee6..b4a59e2 100644 --- a/crates/wok-relay/src/server.rs +++ b/crates/wok-relay/src/server.rs @@ -1489,7 +1489,12 @@ fn ingest_req( }; let mut arr = vec![json!("REQ"), json!(sub_id)]; arr.extend(filters); - let fg = match NostrFilterGroup::from_req(&arr, max_limit, cfg.relay.max_tags_per_filter) { + let fg = match NostrFilterGroup::from_req( + &arr, + max_limit, + cfg.relay.max_tags_per_filter, + cfg.relay.max_and_entries, + ) { Ok(fg) => fg, Err(e) => { fail_closed(e.to_string()); @@ -1580,8 +1585,13 @@ fn ingest_neg( return Err("negentropy filter must be an object".into()); } let max_limit = cfg.relay.max_sync_events + 1; - let fg = NostrFilterGroup::from_value(&filter, max_limit, cfg.relay.max_tags_per_filter) - .map_err(|e| e.to_string())?; + let fg = NostrFilterGroup::from_value( + &filter, + max_limit, + cfg.relay.max_tags_per_filter, + cfg.relay.max_and_entries, + ) + .map_err(|e| e.to_string())?; let query_cost = fg.estimated_cost(false); if cfg.relay.abuse.enabled && cfg.relay.abuse.max_query_cost != 0 @@ -3148,7 +3158,8 @@ fn run_cron(env: Env, cfg: Arc>, shutdown: Arc 0 || vanished_deleted > 0 { diff --git a/crates/wok-ws/src/admin.rs b/crates/wok-ws/src/admin.rs index d83d99d..59634d7 100644 --- a/crates/wok-ws/src/admin.rs +++ b/crates/wok-ws/src/admin.rs @@ -67,6 +67,7 @@ struct LimitsPatch { query_timeslice_budget_us: Option, max_filter_limit: Option, max_tags_per_filter: Option, + max_and_entries: Option, max_filter_limit_count: Option, max_total_events_per_req: Option, max_subs_per_connection: Option, @@ -366,6 +367,7 @@ fn overview(handle: &RelayHandle) -> Response> { "query_timeslice_budget_us": cfg.relay.query_timeslice_budget_us, "max_filter_limit": cfg.relay.max_filter_limit, "max_tags_per_filter": cfg.relay.max_tags_per_filter, + "max_and_entries": cfg.relay.max_and_entries, "max_filter_limit_count": cfg.relay.max_filter_limit_count, "max_total_events_per_req": cfg.relay.max_total_events_per_req, "max_subs_per_connection": cfg.relay.max_subs_per_connection, @@ -525,6 +527,9 @@ fn apply_patch(config: &mut Config, patch: ConfigPatch) { if let Some(value) = limits.max_tags_per_filter { config.relay.max_tags_per_filter = value; } + if let Some(value) = limits.max_and_entries { + config.relay.max_and_entries = value; + } if let Some(value) = limits.max_filter_limit_count { config.relay.max_filter_limit_count = value; } @@ -779,7 +784,7 @@ const GROUPS=[ ['max_event_size','Maximum event bytes','Largest serialized event accepted by the relay.','number'],['max_num_tags','Maximum tags','Largest number of tags accepted on one event.','number'],['max_tag_val_size','Maximum tag value bytes','Largest individual tag value accepted.','number'],['reject_newer_than_secs','Future timestamp tolerance','Reject events this many seconds newer than the relay clock.','number'],['reject_older_than_secs','Maximum event age','Reject non-ephemeral events older than this many seconds.','number'],['reject_ephemeral_older_than_secs','Maximum ephemeral age','Reject ephemeral events older than this many seconds.','number'],['ephemeral_lifetime_secs','Ephemeral TTL','How long TTL-persisted ephemeral events remain available.','number'],['ephemeral_persistence','Ephemeral persistence','Keep ephemeral events live-only or persist them until their TTL.','select',null,[['live_only','Live only'],['ttl','TTL persistence']]] ]}, {key:'limits',title:'Queries and protocol limits',help:'Bounds for subscriptions, filters, result sets, queues, plugins, and Negentropy.',fields:[ - ['max_req_filter_size','Maximum REQ filter bytes','Combined compact-JSON bytes allowed across all filter objects in one REQ or COUNT.','number'],['max_filters_per_req','Maximum filters per request','Unconditional ceiling for filter objects in one REQ or COUNT.','number'],['query_timeslice_budget_us','Query time slice (microseconds)','CPU time a query may use before yielding to other work.','number'],['max_filter_limit','REQ result ceiling','Maximum event limit accepted for a normal subscription filter.','number'],['max_filter_limit_count','COUNT result ceiling','Maximum event limit used while answering COUNT.','number'],['max_tags_per_filter','Tag constraints per filter','Maximum number of tag query keys allowed in one filter.','number'],['max_total_events_per_req','Events per REQ','Maximum total historical events emitted for one REQ; zero is unlimited.','number'],['max_subs_per_connection','Subscriptions per connection','Maximum simultaneous subscriptions on one WebSocket.','number'],['max_pending_outbound_bytes','Outbound queue bytes','Disconnect a slow client after its pending output exceeds this bound; zero is unlimited (never disconnects).','number'],['write_policy_timeout_secs','Write-policy timeout','Seconds to wait for the configured write-policy plugin.','number'],['negentropy_enabled','Negentropy enabled','Advertise and accept NIP-77 synchronization requests.','checkbox'],['max_sync_events','Negentropy event ceiling','Maximum events reconciled by one synchronization session.','number'] + ['max_req_filter_size','Maximum REQ filter bytes','Combined compact-JSON bytes allowed across all filter objects in one REQ or COUNT.','number'],['max_filters_per_req','Maximum filters per request','Unconditional ceiling for filter objects in one REQ or COUNT.','number'],['query_timeslice_budget_us','Query time slice (microseconds)','CPU time a query may use before yielding to other work.','number'],['max_filter_limit','REQ result ceiling','Maximum event limit accepted for a normal subscription filter.','number'],['max_filter_limit_count','COUNT result ceiling','Maximum event limit used while answering COUNT.','number'],['max_tags_per_filter','Tag constraints per filter','Maximum number of distinct tag names allowed across # and & filters.','number'],['max_and_entries','AND entries per filter','Maximum distinct values allowed across all NIP-91 & filters.','number'],['max_total_events_per_req','Events per REQ','Maximum total historical events emitted for one REQ; zero is unlimited.','number'],['max_subs_per_connection','Subscriptions per connection','Maximum simultaneous subscriptions on one WebSocket.','number'],['max_pending_outbound_bytes','Outbound queue bytes','Disconnect a slow client after its pending output exceeds this bound; zero is unlimited (never disconnects).','number'],['write_policy_timeout_secs','Write-policy timeout','Seconds to wait for the configured write-policy plugin.','number'],['negentropy_enabled','Negentropy enabled','Advertise and accept NIP-77 synchronization requests.','checkbox'],['max_sync_events','Negentropy event ceiling','Maximum events reconciled by one synchronization session.','number'] ]}, {key:'abuse',title:'Abuse protection',help:'Token-bucket rates, bursts, query budgets, quotas, and proof-of-work requirements.',fields:[ ['enabled','Abuse protection enabled','Apply all configured connection, message, query, quota, and PoW guards.','checkbox'],['connection_rate_per_second','Connection rate','New connections allowed per second before burst capacity is consumed; zero disables this bucket (unlimited).','number'],['connection_burst','Connection burst','Maximum accumulated burst capacity for new connections; zero disables this bucket (unlimited).','number'],['event_rate_per_second','Connection EVENT rate','EVENT messages allowed per second on one connection; zero disables this bucket (unlimited).','number'],['event_burst','Connection EVENT burst','Maximum accumulated EVENT burst per connection; zero disables this bucket (unlimited).','number'],['pubkey_event_rate_per_second','Pubkey EVENT rate','Accepted events per second across connections for one author; zero disables this bucket (unlimited).','number'],['pubkey_event_burst','Pubkey EVENT burst','Maximum accumulated EVENT burst for one author; zero disables this bucket (unlimited).','number'],['req_rate_per_second','REQ rate','REQ messages allowed per second on one connection; zero disables this bucket (unlimited).','number'],['req_burst','REQ burst','Maximum accumulated REQ burst per connection; zero disables this bucket (unlimited).','number'],['count_rate_per_second','COUNT rate','COUNT messages allowed per second on one connection; zero disables this bucket (unlimited).','number'],['count_burst','COUNT burst','Maximum accumulated COUNT burst per connection; zero disables this bucket (unlimited).','number'],['max_concurrent_historical_queries','Concurrent historical queries','Maximum historical scans running at once per connection; zero rejects new scans.','number'],['max_query_cost','Query cost ceiling','Reject filters whose estimated scan cost exceeds this value; zero is unlimited.','number'],['max_stored_events','Stored events globally','Total durable event ceiling across every author; zero means unlimited.','number'],['max_stored_events_per_pubkey','Stored events per pubkey','Per-author durable event ceiling; zero means unlimited.','number'],['min_pow_difficulty','Minimum proof of work','Required NIP-13 difficulty; zero disables the requirement.','number'] diff --git a/crates/wok-ws/src/lib.rs b/crates/wok-ws/src/lib.rs index 12409d8..8f77032 100644 --- a/crates/wok-ws/src/lib.rs +++ b/crates/wok-ws/src/lib.rs @@ -253,6 +253,8 @@ fn nip11(cfg: &Config, handle: &RelayHandle) -> serde_json::Value { "max_limit": cfg.relay.max_filter_limit, "max_total_events_per_req": cfg.relay.max_total_events_per_req, "max_event_tags": cfg.events.max_num_tags, + "max_tags_per_filter": cfg.relay.max_tags_per_filter, + "max_and_entries": cfg.relay.max_and_entries, "created_at_lower_limit": cfg.events.reject_older_than_secs, "created_at_upper_limit": cfg.events.reject_newer_than_secs, "default_limit": cfg.relay.max_filter_limit, @@ -309,6 +311,7 @@ fn nip_description(nip: u64) -> Option<&'static str> { 62 => Some("Lets a key request complete, durable deletion of its relay-hosted footprint."), 70 => Some("Restricts publication of protected events to their authenticated author."), 77 => Some("Synchronizes event sets efficiently with the Negentropy reconciliation protocol."), + 91 => Some("Requires every value in an &-prefixed tag filter while preserving compatible #-filter fallback."), _ => None, } } diff --git a/docs/config.md b/docs/config.md index 7e9bbca..b8fbff0 100644 --- a/docs/config.md +++ b/docs/config.md @@ -104,7 +104,8 @@ interface or protect the path at the reverse proxy. | `relay.frame_read_timeout_secs` | `30` | Live (new connections) | No | Maximum idle gap between socket reads while a partial WebSocket frame or unfinished fragmented message is buffered (slow-trickle guard); zero disables it. | | `relay.query_timeslice_budget_us` | `10000` | Live | Yes | Query CPU budget before cooperative yielding. | | `relay.max_filter_limit` | `500` | Live | Yes | Maximum normal REQ filter limit. | -| `relay.max_tags_per_filter` | `3` | Live | Yes | Maximum tag query keys in one filter. | +| `relay.max_tags_per_filter` | `3` | Live | Yes | Maximum distinct tag names across `#` and NIP-91 `&` query keys in one filter. | +| `relay.max_and_entries` | `16` | Live | Yes | Maximum distinct values across all NIP-91 `&` query keys in one filter; zero disables NIP-91 and removes it from NIP-11. | | `relay.max_filter_limit_count` | `1000000` | Live | Yes | Maximum COUNT filter limit. | | `relay.max_total_events_per_req` | `2000` | Live | Yes | Deduplicated historical events across a REQ; zero is unlimited. | | `relay.max_subs_per_connection` | `200` | Live | Yes | Simultaneous subscriptions per connection. | diff --git a/docs/nips.md b/docs/nips.md index 7317cd4..52f370d 100644 --- a/docs/nips.md +++ b/docs/nips.md @@ -22,6 +22,7 @@ arbitrary list. | 70 | Protected events | `-` tag + AUTH | `nip_conformance.rs` | always | | 77 | Negentropy | `wok-negentropy` | protocol unit tests | `negentropy.enabled` | | 86 | Relay management API | `wok-ws` RPC + `wok-db` moderation tables + `wok-relay` enforcement | `e2e_transports.rs` | `admin.enabled` with operator pubkeys | +| 91 | AND operator in filters | `wok-query` compiled AND tags + full-event verification | matcher property tests, historical/live/COUNT conformance | both tag-filter limits are nonzero | NIP-02, NIP-04, and NIP-28 event kinds are accepted and stored, but those client/application semantics are deliberately not advertised as relay @@ -41,6 +42,24 @@ implements all specified target forms: raw event/pubkey hex, an address's pubkey, or SHA-256 of any other string. Multi-filter, multi-target, and limited counts omit HLL because their sketches would be ambiguous or incomplete. +NIP-91 support follows the draft in +[nostr-protocol/nips#2252](https://github.com/nostr-protocol/nips/pull/2252) +at head `b93bda29d45998866e81c65e0693616294a78672`. An `&x` filter requires every +listed value to occur in an `x` tag. Values duplicated in the corresponding +`#x` compatibility filter are removed from its OR alternatives; if no OR +alternatives remain, that OR clause is omitted. Historical queries use one +required value as an index seed and verify the complete packed event. COUNT is +exact but omits the optional NIP-45 HLL sketch for filters containing `&`. + +Outbound filters are adapted to the upstream. Router `REQ`s and `wok sync` +`NEG-OPEN`s forward `&` clauses only to relays whose NIP-11 document advertises +NIP-91; the document is cached per relay and is not consulted for anything +else. For every other upstream — including one that serves no usable NIP-11 +document — each `&x` value set is folded into the `#x` compatibility clause the +draft requires clients to send, so the remote answers a superset that is then +narrowed by the exact AND filter locally. The `&` keys are never simply +dropped, which would turn an `&`-only filter into a match-all request. + NIP-62 accepts a signed kind 62 request containing either a matching `["relay", ""]` tag or `["relay", "ALL_RELAYS"]`. The relay immediately suppresses qualifying authored events and gift wraps for the diff --git a/docs/wok.toml b/docs/wok.toml index f4a5da3..8bcee7c 100644 --- a/docs/wok.toml +++ b/docs/wok.toml @@ -65,6 +65,8 @@ frame_read_timeout_secs = 30 query_timeslice_budget_us = 10000 max_filter_limit = 500 max_tags_per_filter = 3 +# Maximum distinct values across all NIP-91 & filters in one filter object. +max_and_entries = 16 max_filter_limit_count = 1000000 # Zero means unlimited. max_total_events_per_req = 2000 diff --git a/fuzz/Cargo.lock b/fuzz/Cargo.lock index 9b6819b..ea319b3 100644 --- a/fuzz/Cargo.lock +++ b/fuzz/Cargo.lock @@ -1471,6 +1471,7 @@ dependencies = [ "wok-db", "wok-event", "wok-negentropy", + "wok-query", "wok-relay", "wok-ws", ] diff --git a/fuzz/Cargo.toml b/fuzz/Cargo.toml index 192c9c0..bec0e95 100644 --- a/fuzz/Cargo.toml +++ b/fuzz/Cargo.toml @@ -12,6 +12,7 @@ libfuzzer-sys = "0.4" wok-db = { path = "../crates/wok-db" } wok-event = { path = "../crates/wok-event" } wok-negentropy = { path = "../crates/wok-negentropy" } +wok-query = { path = "../crates/wok-query" } wok-relay = { path = "../crates/wok-relay" } wok-ws = { path = "../crates/wok-ws" } diff --git a/fuzz/fuzz_targets/ingress.rs b/fuzz/fuzz_targets/ingress.rs index 1276899..bbc8392 100644 --- a/fuzz/fuzz_targets/ingress.rs +++ b/fuzz/fuzz_targets/ingress.rs @@ -60,7 +60,9 @@ impl BTreeBackend for FuzzBackend { fuzz_target!(|data: &[u8]| { if let Ok(text) = std::str::from_utf8(data) { - let _ = wok_event::json::parse_strict(text); + if let Ok(value) = wok_event::json::parse_strict(text) { + let _ = wok_query::NostrFilterGroup::from_value(&value, 500, 3, 16); + } let _ = ClientCommand::parse(text); } let _ = wok_event::PackedEventView::new(data);