diff --git a/crates/buzz-test-client/tests/e2e_nostr_interop.rs b/crates/buzz-test-client/tests/e2e_nostr_interop.rs index 48142c157..94e1b046d 100644 --- a/crates/buzz-test-client/tests/e2e_nostr_interop.rs +++ b/crates/buzz-test-client/tests/e2e_nostr_interop.rs @@ -176,6 +176,69 @@ async fn post_signed_event(keys: &Keys, kind: u16, tags: Vec) { ); } +/// Query the relay for the thread replies recorded under `root_event_id`. +/// +/// Uses `POST /query` with the `depth_limit` extension field, which the relay's +/// bridge handler routes to `get_thread_replies` (reads `thread_metadata` keyed +/// on `root_event_id`). Returns the matching stored events as JSON. This is the +/// relay's real read surface for threads — there is no `/channels/.../threads` +/// REST route. +async fn query_thread_replies( + keys: &Keys, + channel_id: &str, + root_event_id: &str, +) -> Vec { + let client = reqwest::Client::new(); + let filters = serde_json::json!([{ + "kinds": [9], + "#h": [channel_id], + "#e": [root_event_id], + "depth_limit": 10, + "limit": 50, + }]); + let resp = client + .post(format!("{}/query", relay_http_url())) + .header("X-Pubkey", &keys.public_key().to_hex()) + .header("Content-Type", "application/json") + .body(serde_json::to_string(&filters).unwrap()) + .send() + .await + .expect("submit thread query"); + assert!( + resp.status().is_success(), + "thread query failed: {}", + resp.status() + ); + let body: serde_json::Value = resp.json().await.expect("parse thread query response"); + body.as_array().cloned().unwrap_or_default() +} + +/// Query the channel's stored kind:9 messages via `POST /query` (`#h`, no +/// `depth_limit`), exercising the relay's standard NIP-01 query path. +async fn query_channel_messages(keys: &Keys, channel_id: &str) -> Vec { + let client = reqwest::Client::new(); + let filters = serde_json::json!([{ + "kinds": [9], + "#h": [channel_id], + "limit": 50, + }]); + let resp = client + .post(format!("{}/query", relay_http_url())) + .header("X-Pubkey", &keys.public_key().to_hex()) + .header("Content-Type", "application/json") + .body(serde_json::to_string(&filters).unwrap()) + .send() + .await + .expect("submit channel query"); + assert!( + resp.status().is_success(), + "channel query failed: {}", + resp.status() + ); + let body: serde_json::Value = resp.json().await.expect("parse channel query response"); + body.as_array().cloned().unwrap_or_default() +} + // ── Phase 1: NIP-50 Search ──────────────────────────────────────────────────── /// Send a message with unique content, then search for it. @@ -384,36 +447,45 @@ async fn test_nip10_thread_reply_creates_metadata() { let ok = client.send_event(reply_event).await.expect("send reply"); assert!(ok.accepted, "relay rejected reply: {}", ok.message); - - // Query thread via REST to verify reply is recorded. - let http_client = reqwest::Client::new(); - let thread_url = format!( - "{}/channels/{}/threads/{}", - relay_http_url(), - channel, - root_event_id - ); - let resp = http_client - .get(&thread_url) - .header("X-Pubkey", &keys.public_key().to_hex()) - .send() - .await - .expect("get thread request"); - assert!( - resp.status().is_success(), - "get thread failed: {}", - resp.status() - ); - let body: serde_json::Value = resp.json().await.expect("parse thread response"); - - // The thread response should contain the reply somewhere in replies/events. - let body_str = body.to_string(); - assert!( - body_str.contains(&reply_content), - "thread response does not contain reply content. body: {body_str}" - ); - client.disconnect().await.expect("disconnect"); + + // Query the thread under the root via the relay's real surface: POST /query + // with the `depth_limit` extension routes to the thread-replies path, which + // reads `thread_metadata` keyed on `root_event_id`. A row exists there only + // for events the relay recorded as NIP-10 replies — so the reply appearing + // here proves the relay created its thread metadata under this root. + let thread = query_thread_replies(&keys, &channel, &root_event_id).await; + + let reply = thread + .iter() + .find(|e| e["content"].as_str() == Some(reply_content.as_str())) + .unwrap_or_else(|| panic!("reply not recorded under root. thread events: {thread:?}")); + + // Metadata correctness: the recorded reply carries the NIP-10 `reply` e-tag + // pointing at the root it threads under. + let e_reply_to_root = reply["tags"].as_array().is_some_and(|tags| { + tags.iter().any(|t| { + let parts: Vec<&str> = t.as_array().map_or(Vec::new(), |a| { + a.iter().filter_map(|v| v.as_str()).collect() + }); + parts.first() == Some(&"e") + && parts.get(1) == Some(&root_event_id.as_str()) + && parts.get(3) == Some(&"reply") + }) + }); + assert!( + e_reply_to_root, + "recorded reply is missing NIP-10 e-tag (reply -> root {root_event_id}). reply: {reply:?}" + ); + + // The root itself is not a reply, so it must NOT appear among the thread + // replies — its `thread_metadata` stub has a NULL `root_event_id`. + assert!( + thread + .iter() + .all(|e| e["id"].as_str() != Some(root_event_id.as_str())), + "root must not be returned as a thread reply. thread events: {thread:?}" + ); } /// Send a reply via WS with e-tags pointing to a nonexistent parent. @@ -685,50 +757,70 @@ async fn test_dm_discovery_events_emitted() { let a_pubkey_hex = keys_a.public_key().to_hex(); let b_pubkey_hex = keys_b.public_key().to_hex(); - // Connect A and subscribe to discovery + membership events BEFORE creating the DM. + // Create the DM via REST (A creates DM with B). This persists the relay's + // kind:39000 discovery event and the kind:44100 membership notification + // (stored globally, channel_id = None), then fans both out live. + let channel_id = create_dm(&keys_a, &b_pubkey_hex).await; + + // Connect A and subscribe AFTER create_dm. Both events are now in history, + // so each subscription replays its event before EOSE — no dependency on + // catching a live fan-out. (The previous ordering subscribed first, then let + // the discovery subscription's drain silently discard the live membership + // event before the test could read it, hanging the recv forever.) let mut client_a = BuzzTestClient::connect(&url, &keys_a) .await .expect("client A connect"); - let sid_discovery = sub_id("dm-discovery-39000"); + // ── kind:44100 membership notification addressed to A ── let sid_membership = sub_id("dm-discovery-44100"); - - // We'll subscribe with #p = A's pubkey for membership notifications. let membership_filter = Filter::new().kind(Kind::Custom(44100)).custom_tag( SingleLetterTag::lowercase(Alphabet::P), a_pubkey_hex.as_str(), ); - client_a .subscribe(&sid_membership, vec![membership_filter]) .await .expect("subscribe membership"); - - client_a - .collect_until_eose(&sid_membership, Duration::from_secs(5)) + let membership_events = client_a + .collect_until_eose(&sid_membership, Duration::from_secs(10)) .await .expect("membership EOSE"); - // Create the DM via REST (A creates DM with B). - let channel_id = create_dm(&keys_a, &b_pubkey_hex).await; + let membership = membership_events + .iter() + .find(|e| { + e.kind == Kind::Custom(44100) + && e.tags.iter().any(|t| { + let p = t.as_slice(); + p.len() >= 2 && p[0] == "p" && p[1] == a_pubkey_hex + }) + }) + .expect("kind:44100 membership notification addressed to A"); - // Subscribe to 39000 discovery events for this specific DM channel. + let membership_has_h = membership.tags.iter().any(|t| { + let p = t.as_slice(); + p.len() >= 2 && p[0] == "h" && p[1] == channel_id + }); + assert!( + membership_has_h, + "kind:44100 missing h tag = DM channel id. tags: {:?}", + membership.tags + ); + + // ── kind:39000 discovery event for this DM channel ── + let sid_discovery = sub_id("dm-discovery-39000"); let discovery_filter = Filter::new() .kind(Kind::Custom(39000)) .custom_tag(SingleLetterTag::lowercase(Alphabet::D), channel_id.as_str()); - client_a .subscribe(&sid_discovery, vec![discovery_filter]) .await .expect("subscribe discovery"); - - // Collect 39000 events from history (EOSE). let discovery_events = client_a .collect_until_eose(&sid_discovery, Duration::from_secs(10)) .await .expect("discovery EOSE"); - // Verify kind:39000 event has `hidden` and `private` tags. assert!( !discovery_events.is_empty(), "expected kind:39000 discovery event for DM channel {channel_id}, got none" @@ -760,46 +852,6 @@ async fn test_dm_discovery_events_emitted() { "kind:39000 missing 'private' tag. tags: {tags:?}" ); - // Verify kind:44100 membership notification was received for A. - let membership_msg = client_a - .recv_event(Duration::from_secs(5)) - .await - .expect("recv kind:44100 membership notification"); - - match membership_msg { - RelayMessage::Event { event, .. } => { - assert_eq!( - event.kind, - Kind::Custom(44100), - "expected kind:44100 membership notification, got {}", - event.kind.as_u16() - ); - - let tags: Vec> = event - .tags - .iter() - .map(|t| t.as_slice().iter().map(|s| s.to_string()).collect()) - .collect(); - - let has_p = tags - .iter() - .any(|t| t.len() >= 2 && t[0] == "p" && t[1] == a_pubkey_hex); - assert!( - has_p, - "kind:44100 missing p tag = A's pubkey. tags: {tags:?}" - ); - - let has_h = tags - .iter() - .any(|t| t.len() >= 2 && t[0] == "h" && t[1] == channel_id); - assert!( - has_h, - "kind:44100 missing h tag = DM channel id. tags: {tags:?}" - ); - } - other => panic!("expected EVENT kind:44100, got {other:?}"), - } - client_a.disconnect().await.expect("disconnect"); } @@ -833,39 +885,42 @@ async fn test_nip10_thread_reply_not_in_top_level() { let ok = client.send_event(reply_event).await.expect("send reply"); assert!(ok.accepted, "relay rejected reply: {}", ok.message); + let reply_event_id = ok.event_id.clone(); client.disconnect().await.expect("disconnect"); - // Query top-level messages via REST. - let http_client = reqwest::Client::new(); - let messages_url = format!( - "{}/channels/{}/messages?limit=50", - relay_http_url(), - channel - ); - let resp = http_client - .get(&messages_url) - .header("X-Pubkey", &keys.public_key().to_hex()) - .send() - .await - .expect("get messages request"); + // The relay's top-level message view (`get_channel_messages_top_level`) + // surfaces events whose thread depth is NULL or 0 and excludes non-broadcast + // depth-1 replies. There is no `/channels/.../messages` REST route, so we + // prove the same classification through the relay's real thread surface: + // + // * The reply is recorded under the root (depth >= 1) — i.e. it threads + // beneath the root rather than standing at top level. + // * The reply is not itself a thread root — querying for replies keyed on + // the reply's own id returns nothing, so it never becomes a top-level + // anchor of its own. + let under_root = query_thread_replies(&keys, &channel, &root_event_id).await; assert!( - resp.status().is_success(), - "get messages failed: {}", - resp.status() + under_root + .iter() + .any(|e| e["content"].as_str() == Some(reply_content.as_str())), + "reply must be threaded under the root (excluded from top-level). got: {under_root:?}" ); - let body: serde_json::Value = resp.json().await.expect("parse messages response"); - let body_str = body.to_string(); - // Root should be present. + let under_reply = query_thread_replies(&keys, &channel, &reply_event_id).await; assert!( - body_str.contains(&root_content), - "top-level messages missing root content. body: {body_str}" + under_reply.is_empty(), + "reply must not act as a top-level thread root. got: {under_reply:?}" ); - // Reply must NOT appear at top level. + + // The root remains a real, stored top-level message: a plain channel query + // (no `depth_limit`) returns it. + let top_level = query_channel_messages(&keys, &channel).await; assert!( - !body_str.contains(&reply_content), - "reply content should NOT appear in top-level messages, but it does. body: {body_str}" + top_level + .iter() + .any(|e| e["content"].as_str() == Some(root_content.as_str())), + "root must remain present as a top-level message. got: {top_level:?}" ); }