From 0b6681de29561a1000e0e6d44dd9945cca7015cc Mon Sep 17 00:00:00 2001 From: Pinky <5f5ab050ec58ae208332edd544ebf705221e24c1b86d82a6ca07038a7a8f6ac9@buzz.block.builderlab.xyz> Date: Mon, 7 Sep 2026 15:17:07 -0600 Subject: [PATCH 1/2] fix(desktop): scope quota backoff to its transport Keep WS quota and concurrency CLOSED retries from gating HTTP, and keep HTTP 429 errors out of the renderer WS gate. Preserve damping for an explicitly unavailable shared admission service. Reuse existing channel live streams for Home mention invalidation instead of subscribing to each channel twice. Keep kind, recipient, membership, identity, and replay guards. Add production-path HTTP/WS and mounted-hook regressions. No pacing, reconnect-readiness, replay, routing, or relay policy changes. Signed-off-by: Pinky <5f5ab050ec58ae208332edd544ebf705221e24c1b86d82a6ca07038a7a8f6ac9@buzz.block.builderlab.xyz> --- desktop/src-tauri/src/egress_guard_tests.rs | 1 + desktop/src-tauri/src/native_relay_client.rs | 13 +- .../src/native_relay_client_tests.rs | 6 +- .../native_relay_client_transport_tests.rs | 169 ++++++++ desktop/src-tauri/src/relay.rs | 8 +- desktop/src-tauri/src/relay/tests.rs | 2 +- desktop/src-tauri/src/relay_admission.rs | 7 +- .../channels/useLiveChannelUpdates.test.mjs | 380 ++++++++++++++++++ .../channels/useLiveChannelUpdates.ts | 142 +------ desktop/src/shared/api/relayChannelFilters.ts | 15 - .../api/relayClientPublishRejection.test.mjs | 58 +++ desktop/src/shared/api/relayClientSession.ts | 11 - .../shared/api/relayRateLimitGate.test.mjs | 23 -- desktop/src/shared/api/relayRateLimitGate.ts | 16 +- desktop/src/shared/api/tauri.test.mjs | 118 ------ desktop/src/shared/api/tauri.ts | 24 +- 16 files changed, 657 insertions(+), 336 deletions(-) create mode 100644 desktop/src-tauri/src/native_relay_client_transport_tests.rs create mode 100644 desktop/src/features/channels/useLiveChannelUpdates.test.mjs diff --git a/desktop/src-tauri/src/egress_guard_tests.rs b/desktop/src-tauri/src/egress_guard_tests.rs index 29e74cfb506..fe053602f9e 100644 --- a/desktop/src-tauri/src/egress_guard_tests.rs +++ b/desktop/src-tauri/src/egress_guard_tests.rs @@ -271,6 +271,7 @@ const EVENTS_INVENTORY: &[(&str, usize, usize)] = &[ ("src/native_websocket.rs", 0, 2), // boundary 8 (WS frames; no events URL) // Test-only fixtures — no production egress, no guard: ("src/relay_admission.rs", 1, 0), + ("src/native_relay_client_transport_tests.rs", 1, 0), ("src/archive/mod_tests.rs", 1, 0), ("src/managed_agents/persona_events/tests.rs", 1, 0), ("src/commands/team_snapshot/tests.rs", 1, 0), diff --git a/desktop/src-tauri/src/native_relay_client.rs b/desktop/src-tauri/src/native_relay_client.rs index 2237076a926..19740dd0197 100644 --- a/desktop/src-tauri/src/native_relay_client.rs +++ b/desktop/src-tauri/src/native_relay_client.rs @@ -771,11 +771,16 @@ impl ClosedRetry { self.due_at = None; } ClosedClass::RateLimited => { - // Arm the process-wide gate so the HTTP bridge backs off too, - // rather than keeping a second private notion of the same - // relay's back-pressure. + // WS quota/concurrency limits do not consume HTTP's ApiCalls + // budget. Only an explicit failure of the shared admission + // service warrants damping the other transport too. let hint = parse_retry_in_seconds(message); - crate::relay_admission::activate_rate_limit(hint); + if message + .trim() + .eq_ignore_ascii_case("rate-limited: shared admission unavailable") + { + crate::relay_admission::activate_rate_limit(None); + } let hinted = hint .map(Duration::from_secs) .unwrap_or(CLOSED_RATE_LIMIT_DEFAULT); diff --git a/desktop/src-tauri/src/native_relay_client_tests.rs b/desktop/src-tauri/src/native_relay_client_tests.rs index 96ec39a4bf0..ca2324327f3 100644 --- a/desktop/src-tauri/src/native_relay_client_tests.rs +++ b/desktop/src-tauri/src/native_relay_client_tests.rs @@ -487,7 +487,6 @@ async fn a_reconcile_preserves_the_backoff_of_a_still_desired_subscription() { reconcile, got {reopened:?}" ); - crate::relay_admission::reset_rate_limit_gate(); session.shutdown(); } @@ -726,7 +725,6 @@ fn a_rate_limited_closed_waits_at_least_the_relay_hint() { due >= Instant::now() + Duration::from_secs(11), "a 12s hint must not be undercut by the base backoff" ); - crate::relay_admission::reset_rate_limit_gate(); } #[test] @@ -739,7 +737,6 @@ fn a_hintless_rate_limited_closed_uses_the_shared_default() { due >= Instant::now() + CLOSED_RATE_LIMIT_DEFAULT - Duration::from_secs(1), "a hintless rate-limit must fall back to the shared default window" ); - crate::relay_admission::reset_rate_limit_gate(); } #[test] @@ -894,3 +891,6 @@ async fn the_first_lease_installs_a_session_the_archive_then_reuses() { replacing an identically scoped one" ); } + +#[path = "native_relay_client_transport_tests.rs"] +mod transport_tests; diff --git a/desktop/src-tauri/src/native_relay_client_transport_tests.rs b/desktop/src-tauri/src/native_relay_client_transport_tests.rs new file mode 100644 index 00000000000..d064b521e52 --- /dev/null +++ b/desktop/src-tauri/src/native_relay_client_transport_tests.rs @@ -0,0 +1,169 @@ +//! Real persistent-WS receive loop -> real HTTP submit admission regressions. +use super::*; +use crate::relay_admission::{reset_rate_limit_gate, TEST_SERIAL}; +use axum::{routing::post, Json, Router}; + +async fn http_relay() -> ( + String, + mpsc::Receiver<(nostr::Event, std::time::Instant)>, + tokio::task::JoinHandle<()>, +) { + let (sent, received) = mpsc::channel(4); + let router = Router::new() + .route( + "/query", + post(|| async { + ( + axum::http::StatusCode::TOO_MANY_REQUESTS, + Json(serde_json::json!({"error": "rate-limited: quota exceeded; retry in 1s"})), + ) + }), + ) + .route( + "/events", + post(move |Json(event): Json| { + let sent = sent.clone(); + async move { + let received_at = std::time::Instant::now(); + let id = event.id.to_hex(); + assert!(event.verify().is_ok()); + sent.send((event, received_at)).await.unwrap(); + Json(serde_json::json!({"event_id": id, "accepted": true, "message": ""})) + } + }), + ); + let listener = tokio::net::TcpListener::bind("127.0.0.1:0").await.unwrap(); + let address = listener.local_addr().unwrap(); + let server = tokio::spawn(async move { + axum::serve(listener, router).await.unwrap(); + }); + (format!("http://{address}"), received, server) +} + +fn reply(keys: &Keys) -> nostr::Event { + EventBuilder::new(nostr::Kind::Custom(9), "startup reply") + .tags([ + nostr::Tag::parse(["h", "5b130804-d759-40ad-a564-d64cc907fa8e"]).unwrap(), + nostr::Tag::parse(["e", &"a".repeat(64), "", "reply"]).unwrap(), + ]) + .sign_with_keys(keys) + .unwrap() +} + +async fn closed_then_http_submit(message: &str, shared_unavailable: bool) { + let _serial = TEST_SERIAL.lock().await; + reset_rate_limit_gate(); + let (ws_url, mut frames, commands) = stub_relay().await; + let keys = Keys::generate(); + let (session, mut events) = start(ws_url, keys.clone(), None).await; + session + .set_subscriptions(vec![ + probe_subscription(), + Subscription { + id: "barrier".into(), + filter: serde_json::json!({"kinds": [1], "limit": 0}), + }, + ]) + .await; + assert_eq!(next_req(&mut frames, "probe REQ").await, PROBE_ID); + assert_eq!(next_req(&mut frames, "barrier REQ").await, "barrier"); + commands + .send(StubCommand::Closed(PROBE_ID.into(), message.into())) + .await + .unwrap(); + // An ordered frame on an unaffected persistent subscription proves CLOSED + // went through the receive loop. No test-side gate activation or sleeps. + let barrier = EventBuilder::text_note("barrier") + .sign_with_keys(&keys) + .unwrap(); + commands + .send(StubCommand::Event( + "barrier".into(), + serde_json::to_value(barrier).unwrap(), + )) + .await + .unwrap(); + assert_eq!( + tokio::time::timeout(Duration::from_secs(3), events.recv()) + .await + .unwrap() + .unwrap() + .subscription_id, + "barrier" + ); + + let (http_url, mut submitted, server) = http_relay().await; + let state = crate::app_state::build_app_state(); + let event = reply(&keys); + let submit = crate::relay::submit_signed_event_at_with_keys(&event, &state, &http_url, &keys); + let outcome = tokio::time::timeout(Duration::from_secs(1), submit).await; + session.shutdown(); + server.abort(); + reset_rate_limit_gate(); + if shared_unavailable { + assert!( + outcome.is_err(), + "shared admission outage must still damp HTTP" + ); + assert!( + submitted.try_recv().is_err(), + "HTTP must not dispatch during the shared outage" + ); + } else { + let response = outcome + .expect("WS quota must not withhold HTTP submission") + .unwrap(); + assert!(response.accepted); + assert_eq!(response.event_id, event.id.to_hex()); + assert_eq!(submitted.recv().await.unwrap().0, event); + } + assert!( + frames.try_recv().is_err(), + "the limited WS subscription must not reopen early" + ); +} + +#[tokio::test] +async fn persistent_ws_quota_does_not_withhold_http_reply() { + closed_then_http_submit("rate-limited: quota exceeded; retry in 50s", false).await; +} + +#[tokio::test] +async fn persistent_ws_concurrency_does_not_withhold_http_reply() { + closed_then_http_submit("rate-limited: too many concurrent requests", false).await; +} + +#[tokio::test] +async fn persistent_ws_shared_unavailable_still_withholds_http_reply() { + closed_then_http_submit("rate-limited: shared admission unavailable", true).await; +} + +#[tokio::test] +async fn http_429_still_withholds_http_reply_then_accepts_it() { + let _serial = TEST_SERIAL.lock().await; + reset_rate_limit_gate(); + let (http_url, mut submitted, server) = http_relay().await; + let state = crate::app_state::build_app_state(); + let keys = Keys::generate(); + let event = reply(&keys); + let error = crate::relay::query_relay_at(&state, &http_url, &[serde_json::json!({"limit": 1})]) + .await + .unwrap_err(); + assert_eq!(error, "relay rate-limited: retry in 1s"); + let before = std::time::Instant::now(); + let outcome = tokio::time::timeout( + Duration::from_secs(3), + crate::relay::submit_signed_event_at_with_keys(&event, &state, &http_url, &keys), + ) + .await; + server.abort(); + reset_rate_limit_gate(); + let response = outcome.unwrap().unwrap(); + let (received, received_at) = submitted.recv().await.unwrap(); + assert!( + received_at.duration_since(before) >= Duration::from_millis(900), + "HTTP submit must honour its own cooldown" + ); + assert!(response.accepted); + assert_eq!(received, event); +} diff --git a/desktop/src-tauri/src/relay.rs b/desktop/src-tauri/src/relay.rs index 676b9656ff2..2484fc7d14c 100644 --- a/desktop/src-tauri/src/relay.rs +++ b/desktop/src-tauri/src/relay.rs @@ -316,7 +316,7 @@ pub async fn relay_error_message(response: reqwest::Response) -> String { }; // 429 Too Many Requests → typed `relay rate-limited:` prefix so the TS - // client can activate the rate-limit gate without confusing it with a + // client can report back-pressure without confusing it with a // connectivity failure (`relay unreachable:`). Also arm the Rust-side // admission gate here — the one place every relay HTTP error funnels // through — so the next relay-backed command waits out the quota window @@ -324,10 +324,8 @@ pub async fn relay_error_message(response: reqwest::Response) -> String { if status == reqwest::StatusCode::TOO_MANY_REQUESTS { let hint = extract_retry_in_hint(&body); // Clamp the hint to MAX_HINT_SECONDS before arming the Rust gate AND - // before embedding it in the returned string. Every consumer (Rust gate - // via `activate_rate_limit` and TS gate via `applyTauriRateLimitIfNeeded`) - // must see the same capped value — a single policy point prevents the TS - // gate from receiving an uncapped hint from an untrusted relay. + // before embedding it in the returned string, so the caller sees the + // same bounded hint the native HTTP gate actually honours. let capped_hint = hint.map(|s| s.min(crate::relay_admission::MAX_HINT_SECONDS)); crate::relay_admission::activate_rate_limit(capped_hint); if let Some(secs) = capped_hint { diff --git a/desktop/src-tauri/src/relay/tests.rs b/desktop/src-tauri/src/relay/tests.rs index 0fcbc891b79..f2928cbf612 100644 --- a/desktop/src-tauri/src/relay/tests.rs +++ b/desktop/src-tauri/src/relay/tests.rs @@ -45,7 +45,7 @@ fn overlong_digit_string_returns_none() { // // Verify that an oversized relay hint is capped in the returned message // string, not just inside `activate_rate_limit()`. This guarantees every -// consumer — including the TS gate via `applyTauriRateLimitIfNeeded` — +// consumer of the returned error — // receives the capped value rather than the raw untrusted relay value. #[tokio::test] diff --git a/desktop/src-tauri/src/relay_admission.rs b/desktop/src-tauri/src/relay_admission.rs index 4b0dd1f3696..29d6cdb309b 100644 --- a/desktop/src-tauri/src/relay_admission.rs +++ b/desktop/src-tauri/src/relay_admission.rs @@ -17,6 +17,10 @@ //! are driven by user-initiated file transfers rather than bridge event flow, //! and they have independent retry logic. //! +//! The native WS client also arms this gate for the explicit +//! `shared admission unavailable` signal. Quota/concurrency CLOSEDs stay on +//! WebSocket; they do not consume the HTTP bridge's separate ApiCalls budget. +//! //! **Community scope:** the gate is reset on every `apply_workspace` call, //! mirroring the TS gate's `resetRateLimitGate()` on community switch in //! `useCommunityInit.ts`. A 429 from community A cannot stall community B. @@ -36,8 +40,7 @@ const DEFAULT_RATE_LIMIT_SECONDS: u64 = 10; /// Prevents an untrusted relay from pinning traffic for an unreasonable window /// or overflowing `Instant` arithmetic. /// Exposed `pub` so `relay.rs` can clamp the hint before embedding it in the -/// returned error string — ensuring every consumer (Rust gate and TS gate via -/// `applyTauriRateLimitIfNeeded`) sees the same capped value. +/// returned error string, matching the window the native HTTP gate honours. pub const MAX_HINT_SECONDS: u64 = 300; static GATE_EXPIRY: Mutex> = Mutex::new(None); diff --git a/desktop/src/features/channels/useLiveChannelUpdates.test.mjs b/desktop/src/features/channels/useLiveChannelUpdates.test.mjs new file mode 100644 index 00000000000..5b80bd91f6d --- /dev/null +++ b/desktop/src/features/channels/useLiveChannelUpdates.test.mjs @@ -0,0 +1,380 @@ +import assert from "node:assert/strict"; +import { after, before, test } from "node:test"; +import { JSDOM } from "jsdom"; +import { HOME_MENTION_EVENT_KINDS } from "@/shared/constants/kinds"; + +const dom = new JSDOM("", { + url: "http://localhost", +}); +before(() => { + Object.assign(globalThis, { + document: dom.window.document, + HTMLElement: dom.window.HTMLElement, + IS_REACT_ACT_ENVIRONMENT: true, + window: dom.window, + localStorage: dom.window.localStorage, + }); +}); +after(() => dom.window.close()); + +const VIEWER = "a".repeat(64); +const PEER = "b".repeat(64); +function channels(count) { + return Array.from({ length: count }, (_, i) => ({ + id: `channel-${i}`, + name: `channel-${i}`, + channelType: i === 0 ? "dm" : "stream", + })); +} +function message(id, overrides = {}) { + return { + id, + kind: 9, + pubkey: PEER, + content: "hello", + created_at: Math.floor(Date.now() / 1000), + tags: [ + ["h", "channel-0"], + ["p", VIEWER], + ], + sig: "", + ...overrides, + }; +} + +async function mount(initialChannels, options = {}, subscribeImpl) { + const { act, cleanup, renderHook } = await import("@testing-library/react"); + const React = await import("react"); + const { QueryClient, QueryClientProvider } = await import( + "@tanstack/react-query" + ); + const { relayClient } = await import("@/shared/api/relayClient"); + const { useLiveChannelUpdates } = await import("./useLiveChannelUpdates.ts"); + const { channelMessagesKey } = await import( + "@/features/messages/lib/messageQueryKeys" + ); + const originalLive = relayClient.subscribeLive; + // Keep a tripwire for the removed API so restoring the second family fails. + const originalMention = Object.getOwnPropertyDescriptor( + relayClient, + "subscribeToChannelMentionEvents", + ); + const subscriptions = []; + const mentionSubscriptions = []; + relayClient.subscribeLive = async (filter, onEvent) => { + const sub = { filter, onEvent, disposed: false }; + subscriptions.push(sub); + await subscribeImpl?.(sub); + return async () => { + sub.disposed = true; + }; + }; + relayClient.subscribeToChannelMentionEvents = async (...args) => { + mentionSubscriptions.push(args); + return async () => {}; + }; + const queryClient = new QueryClient({ + defaultOptions: { queries: { retry: false, gcTime: Infinity } }, + }); + const wrapper = ({ children }) => + React.createElement(QueryClientProvider, { client: queryClient }, children); + const hook = renderHook( + ({ members, opts }) => useLiveChannelUpdates(members, null, opts), + { + wrapper, + initialProps: { + members: initialChannels, + opts: { currentPubkey: VIEWER, ...options }, + }, + }, + ); + const settle = () => + act(async () => { + // Drain subscription setup and React Query notifications, not relay timing. + await new Promise((resolve) => setTimeout(resolve, 0)); + }); + await settle(); + return { + act, + settle, + subscriptions, + mentionSubscriptions, + queryClient, + channelMessagesKey, + rerender(members, opts = options) { + hook.rerender({ members, opts: { currentPubkey: VIEWER, ...opts } }); + }, + async deliver(sub, event) { + await act(async () => { + sub.onEvent(event); + }); + }, + unmount: hook.unmount, + restore() { + hook.unmount(); + cleanup(); + queryClient.clear(); + relayClient.subscribeLive = originalLive; + if (originalMention) { + Object.defineProperty( + relayClient, + "subscribeToChannelMentionEvents", + originalMention, + ); + } else { + delete relayClient.subscribeToChannelMentionEvents; + } + }, + }; +} + +for (const count of [20, 50, 129]) { + test(`${count} member channels use one live stream each, not a second mention family`, async () => { + const h = await mount(channels(count), { onLiveMention() {} }); + try { + assert.equal(h.subscriptions.length, count); + assert.equal( + h.mentionSubscriptions.length, + 0, + "channel streams already include mention kinds", + ); + assert.deepEqual( + h.subscriptions.map((s) => s.filter["#h"]), + channels(count) + .map((c) => [c.id]) + .sort(), + ); + assert.ok(h.subscriptions.every((s) => s.filter.since > 0)); + assert.ok( + h.subscriptions.every((s) => + HOME_MENTION_EVENT_KINDS.every((kind) => + s.filter.kinds.includes(kind), + ), + ), + "the remaining wire filter must cover every Home mention kind", + ); + } finally { + h.restore(); + } + }); +} + +test("live channel stream drives mention, unread and DM callbacks once across replay", async () => { + const mentions = []; + const unreads = []; + const dms = []; + const h = await mount(channels(1), { + onLiveMention: () => mentions.push("mention"), + onChannelMessage: (id, event) => unreads.push([id, event.id]), + onDmMessage: (event, channel) => dms.push([channel.id, event.id]), + }); + try { + const event = message("mention"); + await h.deliver(h.subscriptions[0], event); + await h.deliver(h.subscriptions[0], event); + assert.deepEqual(mentions, ["mention"]); + assert.deepEqual(unreads, [["channel-0", "mention"]]); + assert.deepEqual(dms, [["channel-0", "mention"]]); + } finally { + h.restore(); + } +}); + +test("mention signal retains kind, recipient, self and member-channel boundaries", async () => { + let mentions = 0; + const h = await mount(channels(1), { onLiveMention: () => mentions++ }); + try { + const sub = h.subscriptions[0]; + await h.deliver( + sub, + message("not-mentioned", { tags: [["h", "channel-0"]] }), + ); + await h.deliver(sub, message("self", { pubkey: VIEWER.toUpperCase() })); + for (const kind of [5, 7, 9005, 40001, 40003, 40008, 40099, 48100, 48101]) { + await h.deliver(sub, message(`aux-${kind}`, { kind })); + } + await h.deliver( + sub, + message("outside", { + tags: [ + ["h", "outside"], + ["p", VIEWER], + ], + }), + ); + assert.equal(mentions, 0); + for (const kind of [9, 40002, 45001, 45003]) { + await h.deliver( + sub, + message(`mention-${kind}`, { + kind, + tags: [ + ["h", "channel-0"], + ["p", VIEWER.toUpperCase()], + ], + }), + ); + } + assert.equal(mentions, 4); + } finally { + h.restore(); + } +}); + +test("latest mention callback and identity are used without subscription churn", async () => { + const notifications = []; + const h = await mount(channels(1)); + try { + const sub = h.subscriptions[0]; + await h.deliver(sub, message("before-callback")); + h.rerender(channels(1), { + onLiveMention: () => notifications.push("first"), + }); + await h.deliver(sub, message("first")); + h.rerender(channels(1), { + onLiveMention: () => notifications.push("latest"), + }); + await h.deliver(sub, message("second")); + h.rerender(channels(1), { + currentPubkey: "c".repeat(64), + onLiveMention: () => notifications.push("new-identity"), + }); + await h.deliver(sub, message("old-viewer")); + await h.deliver( + sub, + message("new-viewer", { + tags: [ + ["h", "channel-0"], + ["p", "c".repeat(64)], + ], + }), + ); + h.rerender(channels(1)); + await h.deliver(sub, message("disabled")); + assert.deepEqual(notifications, ["first", "latest", "new-identity"]); + assert.equal(h.subscriptions.length, 1); + assert.equal(h.mentionSubscriptions.length, 0); + } finally { + h.restore(); + } +}); + +test("membership diff disposes removed channels and keeps unchanged channel delivery", async () => { + let mentions = 0; + const opts = { onLiveMention: () => mentions++ }; + const h = await mount(channels(2), opts); + try { + const [removed, retained] = h.subscriptions; + h.rerender(channels(3).slice(1), opts); + await h.settle(); + assert.equal( + h.subscriptions.length, + 3, + "only the newly joined channel subscribes", + ); + assert.equal(removed.disposed, true); + assert.equal(retained.disposed, false); + await h.deliver(removed, message("removed")); + await h.deliver( + retained, + message("retained", { + tags: [ + ["h", "channel-1"], + ["p", VIEWER], + ], + }), + ); + assert.equal( + mentions, + 1, + "late event from a removed channel must not notify", + ); + } finally { + h.restore(); + } +}); + +test("untagged auxiliary event keeps its single-channel context in the timeline cache", async () => { + let mentions = 0; + const h = await mount(channels(2), { onLiveMention: () => mentions++ }); + try { + const key = h.channelMessagesKey("channel-1"); + h.queryClient.setQueryData(key, [ + message("parent", { tags: [["h", "channel-1"]] }), + ]); + await h.deliver( + h.subscriptions[1], + message("reaction", { + kind: 7, + content: "👀", + tags: [ + ["e", "parent"], + ["p", VIEWER], + ], + }), + ); + const reaction = h.queryClient + .getQueryData(key) + .find((event) => event.id === "reaction"); + assert.ok(reaction); + assert.deepEqual(reaction.tags.at(-1), ["h", "channel-1"]); + assert.equal(mentions, 0); + } finally { + h.restore(); + } +}); + +test("one failed setup does not abort other channel streams and retries only that channel", async () => { + const originalSetTimeout = window.setTimeout; + const originalClearTimeout = window.clearTimeout; + const retries = new Map(); + let timer = 0; + window.setTimeout = (fn) => { + retries.set(++timer, fn); + return timer; + }; + window.clearTimeout = (id) => retries.delete(id); + let failOnce = true; + const h = await mount(channels(2), {}, async (sub) => { + if (sub.filter["#h"][0] === "channel-0" && failOnce) { + failOnce = false; + throw new Error("fixture subscription setup failure"); + } + }); + try { + assert.equal(h.subscriptions.length, 2); + assert.equal(retries.size, 1); + await h.act(async () => { + retries.values().next().value(); + }); + await h.settle(); + assert.equal(h.subscriptions.length, 3); + assert.deepEqual(h.subscriptions[2].filter["#h"], ["channel-0"]); + } finally { + h.restore(); + window.setTimeout = originalSetTimeout; + window.clearTimeout = originalClearTimeout; + } +}); + +test("unmount disposes both established and pending channel streams", async () => { + let release; + const h = await mount(channels(2), { onLiveMention() {} }, async (sub) => { + if (sub.filter["#h"][0] === "channel-1") { + await new Promise((resolve) => { + release = resolve; + }); + } + }); + try { + h.unmount(); + assert.equal(h.subscriptions[0].disposed, true); + await h.act(async () => { + release(); + }); + assert.ok(h.subscriptions.every((sub) => sub.disposed)); + assert.equal(h.mentionSubscriptions.length, 0); + } finally { + h.restore(); + } +}); diff --git a/desktop/src/features/channels/useLiveChannelUpdates.ts b/desktop/src/features/channels/useLiveChannelUpdates.ts index 7598e8db3e5..b25df3acf92 100644 --- a/desktop/src/features/channels/useLiveChannelUpdates.ts +++ b/desktop/src/features/channels/useLiveChannelUpdates.ts @@ -9,11 +9,15 @@ import { getChannelIdFromTags, isThreadReply, } from "@/features/messages/lib/threading"; -import { shouldNotifyForEvent } from "@/features/notifications/lib/shouldNotify"; +import { + hasMentionForEvent, + shouldNotifyForEvent, +} from "@/features/notifications/lib/shouldNotify"; import { relayClient } from "@/shared/api/relayClient"; import { CHANNEL_EVENT_KINDS, CHANNEL_MESSAGE_EVENT_KINDS, + HOME_MENTION_EVENT_KINDS, } from "@/shared/constants/kinds"; import type { Channel, RelayEvent } from "@/shared/api/types"; import { @@ -142,10 +146,8 @@ export function useLiveChannelUpdates( const normalizedCurrentPubkey = options.currentPubkey?.trim().toLowerCase() ?? ""; const seenMentionEventIdsRef = React.useRef(new Set()); - // Reconnect replay overlaps each live filter by five seconds so no message is - // lost at the boundary. Keep one shared guard for every notification side - // effect: the same event can be replayed repeatedly while a relay flaps, and - // mention events also arrive through both the channel and mention filters. + // Reconnect replay overlaps live filters so no message is lost at the + // boundary. Guard notification side effects against repeated delivery. const seenNotificationEventIdsRef = React.useRef(new Set()); const channelsInvalidateRef = React.useRef(null); if (channelsInvalidateRef.current === null) { @@ -243,6 +245,17 @@ export function useLiveChannelUpdates( } const isDmChannel = dmChannelMap.has(channelId); + // Mention kinds are already in the channel stream. Keep Home's narrower + // kind/recipient policy without opening a second REQ for every channel. + if ( + options.onLiveMention && + HOME_MENTION_EVENT_KINDS.some((kind) => kind === event.kind) && + isExternalMentionEvent(event, normalizedCurrentPubkey) && + hasMentionForEvent(event, normalizedCurrentPubkey) && + trackSeenEvent(seenMentionEventIdsRef.current, event.id) + ) { + options.onLiveMention(); + } const isUnreadTriggerKind = isChannelUnreadTriggerKind( event.kind, isDmChannel, @@ -338,19 +351,6 @@ export function useLiveChannelUpdates( ); }); - const handleMentionEvent = React.useEffectEvent((event: RelayEvent) => { - if (!isExternalMentionEvent(event, normalizedCurrentPubkey)) { - return; - } - - if (!trackSeenEvent(seenMentionEventIdsRef.current, event.id)) { - return; - } - - handleIncomingMessage(event); - options.onLiveMention?.(); - }); - React.useEffect(() => { return relayClient.subscribeToReconnects(() => { void queryClient.invalidateQueries({ queryKey: channelsQueryKey }); @@ -446,105 +446,6 @@ export function useLiveChannelUpdates( }; }, [channelIdsKey]); - // Subscribe to mention events per channel with a diff-based manager: only - // subscribe newly-added channels and unsubscribe removed ones on each sync. - // The ref survives re-renders so churn-with-identical-IDs does zero work. - const mentionSubsRef = React.useRef(new Map Promise>()); - const mentionSubsPubkeyRef = React.useRef(null); - - React.useEffect(() => { - if (!options.onLiveMention || normalizedCurrentPubkey.length === 0) { - return; - } - - let isCancelled = false; - let retryTimeout: number | undefined; - let retryAttempt = 0; - - const syncSubs = async (): Promise => { - const activeSubs = mentionSubsRef.current; - - if ( - mentionSubsPubkeyRef.current !== null && - mentionSubsPubkeyRef.current !== normalizedCurrentPubkey - ) { - const stale = Array.from(activeSubs.values()); - activeSubs.clear(); - await Promise.allSettled(stale.map((dispose) => dispose())); - if (isCancelled) return true; - } - mentionSubsPubkeyRef.current = normalizedCurrentPubkey; - - const targetIds = new Set(channelIdsKey ? channelIdsKey.split(",") : []); - - for (const [channelId, dispose] of activeSubs) { - if (!targetIds.has(channelId)) { - activeSubs.delete(channelId); - void dispose().catch(() => {}); - } - } - - let anyFailed = false; - // Pass handleMentionEvent directly — it's a stable useEffectEvent - // callback. Do NOT wrap in an isCancelled check here: subs persist - // across effect runs (that's the point of the diff manager), so a - // stale isCancelled flag from a prior run would silently drop events - // on long-lived subs. - const additions = Array.from(targetIds) - .filter((channelId) => !activeSubs.has(channelId)) - .map(async (channelId) => { - try { - const dispose = await relayClient.subscribeToChannelMentionEvents( - channelId, - normalizedCurrentPubkey, - handleMentionEvent, - ); - if (isCancelled) { - void dispose().catch(() => {}); - return; - } - activeSubs.set(channelId, dispose); - } catch (err) { - anyFailed = true; - console.error( - "Failed to subscribe to mention events", - channelId, - err, - ); - } - }); - await Promise.allSettled(additions); - return !anyFailed; - }; - - const runSync = async () => { - const ok = await syncSubs(); - if (isCancelled) return; - if (ok) { - retryAttempt = 0; - return; - } - const delayMs = Math.min( - LIVE_SUBSCRIPTION_RETRY_BASE_MS * 2 ** retryAttempt, - LIVE_SUBSCRIPTION_RETRY_MAX_MS, - ); - retryAttempt += 1; - retryTimeout = window.setTimeout(() => { - retryTimeout = undefined; - void runSync(); - }, delayMs); - }; - - void runSync(); - - return () => { - isCancelled = true; - if (retryTimeout !== undefined) { - window.clearTimeout(retryTimeout); - } - }; - }, [channelIdsKey, normalizedCurrentPubkey, options.onLiveMention]); - React.useEffect(() => { return () => { channelsInvalidateRef.current?.cancel(); @@ -553,13 +454,6 @@ export function useLiveChannelUpdates( void dispose().catch(() => {}); } liveSubsRef.current.clear(); - - const subs = mentionSubsRef.current; - for (const dispose of subs.values()) { - void dispose().catch(() => {}); - } - subs.clear(); - mentionSubsPubkeyRef.current = null; }; }, []); } diff --git a/desktop/src/shared/api/relayChannelFilters.ts b/desktop/src/shared/api/relayChannelFilters.ts index b9f623a75f8..b9a6cb11508 100644 --- a/desktop/src/shared/api/relayChannelFilters.ts +++ b/desktop/src/shared/api/relayChannelFilters.ts @@ -2,7 +2,6 @@ import { CHANNEL_AUX_EVENT_KINDS, CHANNEL_EVENT_KINDS, CHANNEL_TIMELINE_CONTENT_KINDS, - HOME_MENTION_EVENT_KINDS, KIND_DELETION, KIND_NIP29_DELETE_EVENT, KIND_REACTION, @@ -161,17 +160,3 @@ export function buildGlobalStreamFilter( limit, }; } - -export function buildChannelMentionFilter( - channelId: string, - pubkey: string, - limit: number, -): RelaySubscriptionFilter { - return { - kinds: [...HOME_MENTION_EVENT_KINDS], - "#h": [channelId], - "#p": [pubkey], - limit, - since: Math.floor(Date.now() / 1_000), - }; -} diff --git a/desktop/src/shared/api/relayClientPublishRejection.test.mjs b/desktop/src/shared/api/relayClientPublishRejection.test.mjs index 7875b8679ab..503ac329afa 100644 --- a/desktop/src/shared/api/relayClientPublishRejection.test.mjs +++ b/desktop/src/shared/api/relayClientPublishRejection.test.mjs @@ -16,6 +16,7 @@ const pendingTimers = new Map(); let nextTimerId = 1; const sendAttempts = []; const deliveredFrames = []; +let invokeError = null; let sendTransport = async (args) => { deliveredFrames.push(args); }; @@ -29,6 +30,7 @@ globalThis.window = { clearTimeout: (id) => pendingTimers.delete(id), __TAURI_INTERNALS__: { invoke: async (command, args) => { + if (command === "get_channels" && invokeError) throw invokeError; if (command === "plugin:websocket|send") { sendAttempts.push(args); return sendTransport(args); @@ -39,12 +41,14 @@ globalThis.window = { Date.now = () => fakeNow; const { RelayClient } = await import("./relayClientSession.ts"); +const { invokeTauri } = await import("./tauri.ts"); const { activateRateLimit, isRateLimited, resetRateLimitGate } = await import( "./relayRateLimitGate.ts" ); function reset() { resetRateLimitGate(); + invokeError = null; pendingTimers.clear(); nextTimerId = 1; sendAttempts.length = 0; @@ -293,3 +297,57 @@ test("a community switch after send failure cannot retry through its replacement ); assert.equal(eventFrames().length, 0); }); + +// Drive the actual invoke rejection and publisher, not a classifier helper. +// Restoring HTTP -> WS gate propagation must fail before any timer advances. +for (const message of [ + "relay rate-limited: retry in 50s", + "relay rate-limited: quota exceeded", + "relay rate-limited: retry in 1000000s", +]) { + test(`HTTP backoff does not withhold a WS publish: ${message}`, async () => { + reset(); + invokeError = message; + await assert.rejects(invokeTauri("get_channels"), { message }); + const client = connectedClient(); + const event = { id: "9".repeat(64), kind: 9 }; + const published = client.publishEvent(event, "timed out", "send failed"); + try { + await flushUntil(() => eventFrames().length === 1); + assert.equal(isRateLimited(), false); + await deliver(client, ["OK", event.id, true, ""]); + assert.equal(await published, event); + assert.equal(client.pendingEvents.size, 0); + } finally { + // Also drain safely when the pre-fix gate coupling is restored. + resetRateLimitGate(); + await flushUntil(() => eventFrames().length === 1); + await deliver(client, ["OK", event.id, true, ""]); + await published; + } + }); +} + +test("HTTP failure does not clear an existing WS backoff", async () => { + reset(); + const client = connectedClient(); + await deliver(client, [ + "NOTICE", + "rate-limited: quota exceeded; retry in 4s", + ]); + invokeError = "relay rate-limited: retry in 50s"; + await assert.rejects(invokeTauri("get_channels"), { message: invokeError }); + const event = { id: "8".repeat(64), kind: 9 }; + const published = client.publishEvent(event, "timed out", "send failed"); + await Promise.resolve(); + assert.equal(eventFrames().length, 0); + assert.equal(isRateLimited(), true); + assert.ok([...pendingTimers.values()].some(({ fireAt }) => fireAt === 4_000)); + assert.ok( + ![...pendingTimers.values()].some(({ fireAt }) => fireAt === 50_000), + ); + resetRateLimitGate(); + await flushUntil(() => eventFrames().length === 1); + await deliver(client, ["OK", event.id, true, ""]); + assert.equal(await published, event); +}); diff --git a/desktop/src/shared/api/relayClientSession.ts b/desktop/src/shared/api/relayClientSession.ts index 95bcec79700..d9be0e6c9bc 100644 --- a/desktop/src/shared/api/relayClientSession.ts +++ b/desktop/src/shared/api/relayClientSession.ts @@ -26,7 +26,6 @@ import { buildChannelAuxDeletionFilter, buildChannelFilter, buildChannelHistoryFilter, - buildChannelMentionFilter, buildGlobalStreamFilter, } from "@/shared/api/relayChannelFilters"; import { @@ -413,16 +412,6 @@ export class RelayClient { ) { return this.subscribe(filter, onEvent, onReady, readinessTimeoutMs); } - async subscribeToChannelMentionEvents( - channelId: string, - pubkey: string, - onEvent: (event: RelayEvent) => void, - ) { - return this.subscribe( - buildChannelMentionFilter(channelId, pubkey, 50), - onEvent, - ); - } async preconnect() { // Explicit re-engagement (reconnect card / community switch): clears the // terminal latch and AUTH rejection streak, and bypasses backoff once. diff --git a/desktop/src/shared/api/relayRateLimitGate.test.mjs b/desktop/src/shared/api/relayRateLimitGate.test.mjs index a458e87b487..dd2346911fe 100644 --- a/desktop/src/shared/api/relayRateLimitGate.test.mjs +++ b/desktop/src/shared/api/relayRateLimitGate.test.mjs @@ -177,29 +177,6 @@ test("hint exactly at MAX_HINT_SECONDS is honoured without clamping", () => { assert.equal(isRateLimited(), false); }); -test("applyTauriRateLimitIfNeeded with oversized hint clamps to MAX_HINT_SECONDS", async () => { - // This test imports applyTauriRateLimitIfNeeded separately and verifies that - // the TS cap applies even when the message string contains a large hint value - // (the Rust layer clamps in practice, but TS must be independently safe). - reset(0); - const { applyTauriRateLimitIfNeeded } = await import("./tauri.ts"); - // Simulate a message that somehow escaped the Rust cap (defence-in-depth). - applyTauriRateLimitIfNeeded("relay rate-limited: retry in 1000000s"); - // Gate should cap at MAX_HINT_SECONDS * 1000 ms. - tickTo(MAX_HINT_SECONDS * 1_000 - 1); - assert.equal( - isRateLimited(), - true, - "gate must still be active just before cap", - ); - tickTo(MAX_HINT_SECONDS * 1_000 + 1); - assert.equal( - isRateLimited(), - false, - "gate must expire at MAX_HINT_SECONDS, not 1 000 000s", - ); -}); - // ── waitForRateLimit ────────────────────────────────────────────────────────── test("waitForRateLimit resolves immediately when not rate-limited", async () => { diff --git a/desktop/src/shared/api/relayRateLimitGate.ts b/desktop/src/shared/api/relayRateLimitGate.ts index 040bedae780..827713eb069 100644 --- a/desktop/src/shared/api/relayRateLimitGate.ts +++ b/desktop/src/shared/api/relayRateLimitGate.ts @@ -1,9 +1,9 @@ /** - * Module-level rate-limit gate for the relay WebSocket and HTTP bridge. + * Module-level rate-limit gate for relay WebSocket operations. * - * When the relay signals back-pressure via a CLOSED `rate-limited:` message or - * an HTTP 429 response, callers activate the gate. Operations that must not run - * while rate-limited call `isRateLimited()` or await `waitForRateLimit()`. + * WebSocket back-pressure activates this gate. HTTP 429s are handled by the + * native HTTP gate: ApiCalls and WsEvents are separate relay quota budgets. + * Callers that must not run while WS-limited await `waitForRateLimit()`. * * The gate is a singleton: one shared expiry covers all concurrent callers so * overlapping hints (multiple CLOSED frames) extend to the latest expiry without @@ -15,12 +15,10 @@ const DEFAULT_RATE_LIMIT_SECONDS = 10; /** - * Maximum hint the TS gate will honour from a relay 429 response. + * Maximum hint the TS gate will honour from a WebSocket rejection. * - * Mirrors `MAX_HINT_SECONDS` in `relay_admission.rs` (Rust). The Rust relay - * layer clamps the hint before embedding it in the error string, so in practice - * this TS cap is a defence-in-depth guard against any future Rust path that - * forgets to clamp, keeping both gates on the same documented bound. + * Uses the same cap as the native HTTP gate (`MAX_HINT_SECONDS` in + * `relay_admission.rs`), but each transport consumes its own retry hints. */ export const MAX_HINT_SECONDS = 300; diff --git a/desktop/src/shared/api/tauri.test.mjs b/desktop/src/shared/api/tauri.test.mjs index 2273b55fca3..203c0fca2e9 100644 --- a/desktop/src/shared/api/tauri.test.mjs +++ b/desktop/src/shared/api/tauri.test.mjs @@ -1,117 +1,6 @@ -/** - * Unit tests for tauri.ts — focused on `applyTauriRateLimitIfNeeded`, the - * extracted `relay rate-limited:` classifier that activates the shared - * rate-limit gate when Rust emits an HTTP 429 error prefix. - * - * Testing the exported production function (not a local copy) ensures any - * change to the classifier logic is immediately covered here. - */ import assert from "node:assert/strict"; import test from "node:test"; -// ── Fake-timer + gate setup ─────────────────────────────────────────────────── - -let fakeNow = 0; -const pendingTimers = new Map(); -let nextTimerId = 1; - -function fakeSetTimeout(fn, ms) { - const id = nextTimerId++; - pendingTimers.set(id, { fn, fireAt: fakeNow + ms }); - return id; -} - -function fakeClearTimeout(id) { - pendingTimers.delete(id); -} - -function tickTo(ms) { - fakeNow = ms; - for (const [id, { fn, fireAt }] of Array.from(pendingTimers.entries())) { - if (fireAt <= fakeNow) { - pendingTimers.delete(id); - fn(); - } - } -} - -const origDateNow = Date.now; -function setFakeNow(ms) { - fakeNow = ms; - Date.now = () => fakeNow; -} - -globalThis.window = { - setTimeout: fakeSetTimeout, - clearTimeout: fakeClearTimeout, -}; - -setFakeNow(0); - -const { isRateLimited, resetRateLimitGate } = await import( - "./relayRateLimitGate.ts" -); - -// Import the production classifier from tauri.ts — tests must exercise the -// real function, not a local copy, so a logic change is always caught here. -const { applyTauriRateLimitIfNeeded } = await import("./tauri.ts"); - -function resetGate(startMs = 0) { - pendingTimers.clear(); - nextTimerId = 1; - setFakeNow(startMs); - resetRateLimitGate(); -} - -// ── applyTauriRateLimitIfNeeded: relay rate-limited: prefix ─────────────────── - -test("relay rate-limited: prefix activates the rate-limit gate", () => { - resetGate(0); - applyTauriRateLimitIfNeeded("relay rate-limited: retry in 10s"); - assert.equal(isRateLimited(), true, "gate must be active after 429 error"); -}); - -test("relay rate-limited: prefix parses the retry hint and arms the gate duration", () => { - resetGate(0); - applyTauriRateLimitIfNeeded("relay rate-limited: retry in 7s"); - // Gate should be active at 6s. - setFakeNow(6_000); - assert.equal(isRateLimited(), true); - // Gate should expire after 7s. - tickTo(7_001); - assert.equal(isRateLimited(), false); -}); - -test("relay rate-limited: with no hint uses the 10s default", () => { - resetGate(0); - applyTauriRateLimitIfNeeded("relay rate-limited: quota exceeded"); - tickTo(9_999); - assert.equal(isRateLimited(), true); - tickTo(10_001); - assert.equal(isRateLimited(), false); -}); - -test("non-rate-limited error does not activate the gate", () => { - resetGate(0); - applyTauriRateLimitIfNeeded("relay returned 404 Not Found"); - assert.equal( - isRateLimited(), - false, - "gate must remain inactive for unrelated errors", - ); -}); - -test("relay rate-limited: prefix check is case-sensitive (Rust always emits lowercase)", () => { - resetGate(0); - // The prefix from Rust is always lowercase; mixed-case must not trigger it. - applyTauriRateLimitIfNeeded("Relay rate-limited: retry in 5s"); - assert.equal( - isRateLimited(), - false, - "uppercase prefix must not activate gate (relay emits lowercase only)", - ); -}); - // ── fromRawAcpRuntimeCatalogEntry: custom row API-boundary (B-2) ───────────── // // These tests feed real raw custom catalog rows through fromRawAcpRuntimeCatalogEntry @@ -256,10 +145,3 @@ test("fromRawAcpRuntimeCatalogEntry omits maxParallelism when max_parallelism is "uncapped harness must have maxParallelism: undefined", ); }); - -// ── Teardown ────────────────────────────────────────────────────────────────── - -test("teardown — restore Date.now", () => { - Date.now = origDateNow; - assert.ok(true); -}); diff --git a/desktop/src/shared/api/tauri.ts b/desktop/src/shared/api/tauri.ts index 984b9d176df..76dadebe8b7 100644 --- a/desktop/src/shared/api/tauri.ts +++ b/desktop/src/shared/api/tauri.ts @@ -1,8 +1,4 @@ import { invoke as tauriInvoke } from "@tauri-apps/api/core"; -import { - activateRateLimit, - parseRateLimitHint, -} from "@/shared/api/relayRateLimitGate"; import { fromRawInstallRuntimeResult, type RawInstallRuntimeResult, @@ -280,18 +276,6 @@ function toTauriError(error: unknown): Error { } } -/** - * Inspect a Tauri error message and activate the shared rate-limit gate when - * the Rust relay layer emitted an HTTP 429 response (`relay rate-limited:` prefix). - * - * Extracted so it can be unit-tested without mocking the Tauri invoke bridge. - */ -export function applyTauriRateLimitIfNeeded(message: string): void { - if (message.startsWith("relay rate-limited:")) { - activateRateLimit(parseRateLimitHint(message)); - } -} - export async function invokeTauri( command: string, args?: Record, @@ -299,11 +283,9 @@ export async function invokeTauri( try { return await tauriInvoke(command, args); } catch (error) { - const err = toTauriError(error); - // Rust emits `relay rate-limited:` for HTTP 429 responses. Activate the - // shared gate so the TS relay client backs off for the same window. - applyTauriRateLimitIfNeeded(err.message); - throw err; + // HTTP backoff lives in Rust. Do not apply its separate ApiCalls quota + // to the WebSocket gate, but preserve the failure for the caller. + throw toTauriError(error); } } From 83ac81ff6ae439b0e68f4a3dafd4b2af6dbadcdb Mon Sep 17 00:00:00 2001 From: Pinky <5f5ab050ec58ae208332edd544ebf705221e24c1b86d82a6ca07038a7a8f6ac9@buzz.block.builderlab.xyz> Date: Mon, 7 Sep 2026 22:01:53 -0600 Subject: [PATCH 2/2] perf(desktop): reuse discovery rosters within channel fetches Use the complete membership events already fetched for discovery rather than reading the same rosters again. Query only uncovered directory or pending-owner channels, preserving covered rosters when the tolerant fallback fails. Keep reuse request-local; do not change refresh triggers, channel scope, activity batching, or introduce a cache. Cover real native HTTP requests, pagination, fallback and identity/community changes in regression tests. Signed-off-by: Pinky <5f5ab050ec58ae208332edd544ebf705221e24c1b86d82a6ca07038a7a8f6ac9@buzz.block.builderlab.xyz> --- .../src-tauri/src/commands/channels/fetch.rs | 45 ++- .../src/commands/channels/fetch_tests.rs | 378 ++++++++++++++++++ 2 files changed, 409 insertions(+), 14 deletions(-) create mode 100644 desktop/src-tauri/src/commands/channels/fetch_tests.rs diff --git a/desktop/src-tauri/src/commands/channels/fetch.rs b/desktop/src-tauri/src/commands/channels/fetch.rs index 36c24a35b7d..993d40807e2 100644 --- a/desktop/src-tauri/src/commands/channels/fetch.rs +++ b/desktop/src-tauri/src/commands/channels/fetch.rs @@ -186,9 +186,9 @@ pub(super) enum DirectoryScope { /// - Phase 1 (parallel): member-chain (kind:39002→kind:39000), the non-member /// metadata source (pending-owned ids when member-only, else the all-open /// kind:39000 scan), and the hidden-DM snapshot (kind:30622). -/// - Phase 2 (parallel): member counts (kind:39002 batch) and last-message -/// timestamps (bounded per-channel human-visible activity batches), fanned -/// out over the merged set. Member-count failures degrade to zero; timestamp +/// - Phase 2 (parallel): missing member counts (kind:39002 batch) and last-message +/// timestamps (bounded per-channel human-visible activity batches). Reuse the +/// member-chain rosters; missing-count failures degrade to zero. Timestamp /// failures abort so cached recency is never replaced by a false /// authoritative empty result. pub(super) async fn fetch_channels( @@ -263,7 +263,7 @@ pub(super) async fn fetch_channels( Vec::new() }; - Ok::<_, String>(meta_events) + Ok::<_, String>((meta_events, collect_members_by_channel(&member_events))) }, // Step 3: non-member channel metadata (kind:39000). // - IncludeOpenDirectory: scan ALL open channels so the discovery @@ -321,7 +321,7 @@ pub(super) async fn fetch_channels( #[cfg(debug_assertions)] let t_phase1 = _profile_start.elapsed(); - let meta_events = member_chain_result?; + let (meta_events, mut membership) = member_chain_result?; let open_meta_events = open_meta_result?; // hidden_dms is already a resolved HashSet (tolerant path above) @@ -371,8 +371,8 @@ pub(super) async fn fetch_channels( } } - // Phase 2 — concurrent: member counts (step 4) and last-message timestamps - // (step 5). Member-count failures degrade to zero. Timestamp failures + // Phase 2 — concurrent: missing member counts (step 4) and last-message + // timestamps (step 5). Missing-count failures degrade to zero. Timestamp failures // abort this refresh so the frontend keeps its previous Recent ordering. let all_channel_ids: Vec = channels.iter().map(|c| c.id.clone()).collect(); if !all_channel_ids.is_empty() { @@ -381,16 +381,27 @@ pub(super) async fn fetch_channels( .map(|id| last_message_filter(id)) .collect(); - // Bind both filter arrays before the join so their lifetimes cover - // both branches of the concurrent pair. + // Step 1 already returned complete rosters, not just the matching p-tag. + // Only directory-only or still-pending channels need another read. Keep + // reuse local to this fetch so the next refresh sees membership changes. + let missing_member_ids: Vec<&String> = all_channel_ids + .iter() + .filter(|id| !membership.contains_key(*id)) + .collect(); let member_count_filters = [serde_json::json!({ "kinds": [39002], - "#d": &all_channel_ids, - "limit": all_channel_ids.len(), + "#d": &missing_member_ids, + "limit": missing_member_ids.len(), })]; let (members_result, message_result) = tokio::join!( - // Step 4: batch-fetch kind:39002 for member counts. - query_relay(state, &member_count_filters), + // Step 4: do not send an empty #d filter (an unscoped roster query). + async { + if missing_member_ids.is_empty() { + Ok(Vec::new()) + } else { + query_relay(state, &member_count_filters).await + } + }, // Step 5: preserve one indexed filter per channel while keeping // every relay request within its aggregate explicit-channel cap. query_last_messages(state, &last_msg_filters), @@ -400,7 +411,9 @@ pub(super) async fn fetch_channels( // empty result and clear every cached timestamp in the frontend. let messages = message_result?; - let membership = collect_members_by_channel(&members_result.unwrap_or_default()); + membership.extend(collect_members_by_channel( + &members_result.unwrap_or_default(), + )); for channel in &mut channels { if let Some(info) = membership.get(&channel.id) { channel.member_count = info.count; @@ -488,3 +501,7 @@ pub(super) fn collect_members_by_channel( } map } + +#[cfg(test)] +#[path = "fetch_tests.rs"] +mod tests; diff --git a/desktop/src-tauri/src/commands/channels/fetch_tests.rs b/desktop/src-tauri/src/commands/channels/fetch_tests.rs new file mode 100644 index 00000000000..439dfbd6c56 --- /dev/null +++ b/desktop/src-tauri/src/commands/channels/fetch_tests.rs @@ -0,0 +1,378 @@ +//! Exercise the production channel-list fetch over the native HTTP bridge. +use super::*; +use axum::{extract::State, http::StatusCode, routing::post, Json, Router}; +use nostr::{Event, EventBuilder, Keys, Kind, Tag, Timestamp}; +use serde_json::{json, Value}; +use std::sync::{Arc, Mutex}; + +#[derive(Default)] +struct Fixture { + events: Vec, + requests: Vec>, + fail_discovery: bool, + fail_fallback: bool, + fail_messages: bool, +} + +struct Relay { + data: Arc>, + url: String, + task: tokio::task::JoinHandle<()>, +} + +impl Drop for Relay { + fn drop(&mut self) { + self.task.abort(); + } +} + +async fn query( + State(data): State>>, + Json(filters): Json>, +) -> (StatusCode, Json) { + let mut data = data.lock().unwrap(); + data.requests.push(filters.clone()); + let mut result = Vec::new(); + for filter in &filters { + let kind = filter["kinds"][0].as_u64().unwrap(); + if (kind == 39002 && filter.get("#p").is_some() && data.fail_discovery) + || (kind == 39002 && filter.get("#d").is_some() && data.fail_fallback) + || (kind == 9 && data.fail_messages) + { + return ( + StatusCode::SERVICE_UNAVAILABLE, + Json(json!({"error": "fixture unavailable"})), + ); + } + let mut page: Vec<_> = data + .events + .iter() + .filter(|event| { + if !filter["kinds"] + .as_array() + .unwrap() + .contains(&json!(event.kind.as_u16())) + { + return false; + } + for name in ["p", "d", "h"] { + if let Some(values) = filter.get(format!("#{name}")).and_then(Value::as_array) { + if !event.tags.iter().any(|tag| { + let s = tag.as_slice(); + s.len() >= 2 && s[0] == name && values.contains(&json!(s[1])) + }) { + return false; + } + } + } + if let Some(until) = filter["until"].as_u64() { + let ts = event.created_at.as_secs(); + if ts > until + || (ts == until + && filter["before_id"] + .as_str() + .is_some_and(|id| event.id.to_hex().as_str() <= id)) + { + return false; + } + } + true + }) + .cloned() + .collect(); + page.sort_by(|a, b| { + b.created_at + .cmp(&a.created_at) + .then_with(|| a.id.cmp(&b.id)) + }); + page.truncate(filter["limit"].as_u64().unwrap() as usize); + result.extend(page); + } + (StatusCode::OK, Json(json!(result))) +} + +impl Relay { + async fn new(events: Vec) -> Self { + let data = Arc::new(Mutex::new(Fixture { + events, + ..Default::default() + })); + let router = Router::new() + .route("/query", post(query)) + .with_state(data.clone()); + let listener = tokio::net::TcpListener::bind("127.0.0.1:0").await.unwrap(); + let url = format!("http://{}", listener.local_addr().unwrap()); + let task = tokio::spawn(async move { axum::serve(listener, router).await.unwrap() }); + Self { data, url, task } + } + + fn state(&self, keys: &Keys) -> AppState { + let state = crate::app_state::build_app_state(); + *state.keys.lock().unwrap() = keys.clone(); + *state.relay_url_override.lock().unwrap() = Some(self.url.clone()); + state + } + + fn roster_fallbacks(&self) -> Vec { + self.data + .lock() + .unwrap() + .requests + .iter() + .flatten() + .filter(|filter| filter["kinds"] == json!([39002]) && filter.get("#d").is_some()) + .cloned() + .collect() + } +} + +fn event(keys: &Keys, kind: u16, tags: Vec>) -> Event { + EventBuilder::new(Kind::from_u16(kind), "") + .allow_self_tagging() + .tags(tags.into_iter().map(|tag| Tag::parse(tag).unwrap())) + .custom_created_at(Timestamp::from(1_700_000_000)) + .sign_with_keys(keys) + .unwrap() +} + +fn metadata(keys: &Keys, id: &str) -> Event { + event(keys, 39000, vec![vec!["d", id], vec!["name", id]]) +} + +fn roster(keys: &Keys, id: &str, members: &[&str]) -> Event { + let mut tags = vec![vec!["d", id]]; + tags.extend(members.iter().map(|pk| vec!["p", *pk, "", "member"])); + event(keys, 39002, tags) +} + +async fn fetch(state: &AppState, scope: DirectoryScope) -> Result, String> { + tokio::time::timeout( + std::time::Duration::from_secs(10), + fetch_channels(state, scope), + ) + .await + .expect("bounded channel fetch") +} + +#[tokio::test] +async fn member_rosters_are_reused_including_every_discovery_page() { + let _serial = crate::relay_admission::TEST_SERIAL.lock().await; + let keys = Keys::generate(); + let me = keys.public_key().to_hex(); + let other = Keys::generate().public_key().to_hex(); + for count in [265, 501] { + let mut events = Vec::new(); + for index in 0..count { + let id = format!("channel-{index:04}"); + events.push(metadata(&keys, &id)); + // Repeated p-tag must not inflate the count; non-self members must survive. + events.push(roster(&keys, &id, &[&me, &other, &other])); + } + let relay = Relay::new(events).await; + let state = relay.state(&keys); + let started = std::time::Instant::now(); + let channels = fetch(&state, DirectoryScope::MemberOnly).await.unwrap(); + assert_eq!(channels.len(), count); + assert!(channels.iter().all(|c| c.is_member + && c.member_count == 2 + && c.member_pubkeys == vec![me.clone(), other.clone()])); + assert!( + relay.roster_fallbacks().is_empty(), + "covered rosters must not be fetched twice" + ); + let data = relay.data.lock().unwrap(); + let discovery: Vec<_> = data + .requests + .iter() + .flatten() + .filter(|f| f["kinds"] == json!([39002])) + .collect(); + assert_eq!(discovery.len(), count / DIRECTORY_PAGE_SIZE + 1); + if count > DIRECTORY_PAGE_SIZE { + assert_eq!(discovery[1]["until"], json!(1_700_000_000)); + assert!(discovery[1]["before_id"].is_string()); + } + // Membership pages + member metadata + hidden DMs + bounded activity batches. + let expected_reads = count / DIRECTORY_PAGE_SIZE + 1 + 2 + count.div_ceil(128); + assert_eq!(data.requests.len(), expected_reads); + eprintln!( + "roster-reuse fixture channels={count} reads={} elapsed={:?}", + data.requests.len(), + started.elapsed() + ); + } +} + +#[tokio::test] +async fn directory_fetches_only_uncovered_rosters_and_keeps_hidden_dm_behavior() { + let _serial = crate::relay_admission::TEST_SERIAL.lock().await; + let keys = Keys::generate(); + let me = keys.public_key().to_hex(); + let other = Keys::generate().public_key().to_hex(); + let relay = Relay::new(vec![ + metadata(&keys, "joined"), + roster(&keys, "joined", &[&me, &other]), + metadata(&keys, "open"), + roster(&keys, "open", &[&other]), + metadata(&keys, "empty"), + roster(&keys, "empty", &[]), + event(&keys, 39000, vec![vec!["d", "hidden-dm"], vec!["t", "dm"]]), + roster(&keys, "hidden-dm", &[&me, &other]), + event( + &keys, + buzz_core_pkg::kind::KIND_DM_VISIBILITY.try_into().unwrap(), + vec![vec!["p", &me], vec!["h", "hidden-dm"]], + ), + event(&keys, 9, vec![vec!["h", "joined"]]), + ]) + .await; + let channels = fetch(&relay.state(&keys), DirectoryScope::IncludeOpenDirectory) + .await + .unwrap(); + assert_eq!(channels.len(), 3); + let joined = channels.iter().find(|c| c.id == "joined").unwrap(); + assert!(joined.is_member); + assert_eq!(joined.member_pubkeys, vec![me, other.clone()]); + assert!(joined.last_message_at.is_some()); + let open = channels.iter().find(|c| c.id == "open").unwrap(); + assert!(!open.is_member); + assert_eq!(open.member_count, 1); + assert_eq!(open.member_pubkeys, vec![other]); + let empty = channels.iter().find(|c| c.id == "empty").unwrap(); + assert!(!empty.is_member); + assert_eq!(empty.member_count, 0); + assert!(empty.member_pubkeys.is_empty()); + let fallbacks = relay.roster_fallbacks(); + assert_eq!(fallbacks.len(), 1); + let mut ids = fallbacks[0]["#d"].as_array().unwrap().clone(); + ids.sort_by_key(Value::to_string); + assert_eq!(ids, vec![json!("empty"), json!("open")]); + assert_eq!(fallbacks[0]["limit"], json!(2)); +} + +#[tokio::test] +async fn pending_owner_fallback_failure_preserves_covered_rosters() { + let _serial = crate::relay_admission::TEST_SERIAL.lock().await; + let keys = Keys::generate(); + let me = keys.public_key().to_hex(); + let other = Keys::generate().public_key().to_hex(); + let relay = Relay::new(vec![ + metadata(&keys, "joined"), + roster(&keys, "joined", &[&me, &other]), + metadata(&keys, "pending"), + roster(&keys, "pending", &[&other]), + metadata(&keys, "unpropagated"), + metadata(&keys, "someone-elses-pending"), + ]) + .await; + let state = relay.state(&keys); + state.mark_pending_owned_channel(&me, "joined"); + state.mark_pending_owned_channel(&me, "pending"); + state.mark_pending_owned_channel(&me, "unpropagated"); + state.mark_pending_owned_channel(&other, "someone-elses-pending"); + for fail in [false, true] { + relay.data.lock().unwrap().fail_fallback = fail; + let channels = fetch(&state, DirectoryScope::MemberOnly).await.unwrap(); + assert_eq!(channels.len(), 3); + let unpropagated = channels.iter().find(|c| c.id == "unpropagated").unwrap(); + assert!(unpropagated.is_member); + assert_eq!(unpropagated.member_count, 0); + assert!(unpropagated.member_pubkeys.is_empty()); + let pending = channels.iter().find(|c| c.id == "pending").unwrap(); + assert!(pending.is_member); + assert_eq!(pending.member_count, if fail { 0 } else { 1 }); + assert_eq!( + channels + .iter() + .find(|c| c.id == "joined") + .unwrap() + .member_count, + 2 + ); + assert!(!state.is_pending_owned_channel(&me, "joined")); + assert!(state.is_pending_owned_channel(&me, "pending")); + } + for fallback in relay.roster_fallbacks() { + let mut ids = fallback["#d"].as_array().unwrap().clone(); + ids.sort_by_key(Value::to_string); + assert_eq!(ids, vec![json!("pending"), json!("unpropagated")]); + assert_eq!(fallback["limit"], json!(2)); + } +} + +#[tokio::test] +async fn each_fetch_observes_roster_changes_and_the_current_identity_and_relay() { + let _serial = crate::relay_admission::TEST_SERIAL.lock().await; + let keys = Keys::generate(); + let me = keys.public_key().to_hex(); + let next_keys = Keys::generate(); + let next = next_keys.public_key().to_hex(); + let relay = Relay::new(vec![ + metadata(&keys, "same-id"), + roster(&keys, "same-id", &[&me]), + ]) + .await; + let state = relay.state(&keys); + assert_eq!( + fetch(&state, DirectoryScope::MemberOnly).await.unwrap()[0].member_pubkeys, + vec![me.clone()] + ); + relay.data.lock().unwrap().events[1] = roster(&keys, "same-id", &[&me, &next]); + assert_eq!( + fetch(&state, DirectoryScope::MemberOnly).await.unwrap()[0].member_count, + 2 + ); + relay.data.lock().unwrap().events[1] = roster(&keys, "same-id", &[&next]); + assert!(fetch(&state, DirectoryScope::MemberOnly) + .await + .unwrap() + .is_empty()); + *state.keys.lock().unwrap() = next_keys.clone(); + assert_eq!( + fetch(&state, DirectoryScope::MemberOnly).await.unwrap()[0].member_pubkeys, + vec![next.clone()] + ); + let second = Relay::new(vec![ + metadata(&keys, "same-id"), + roster(&keys, "same-id", &[&next, &me]), + ]) + .await; + *state.relay_url_override.lock().unwrap() = Some(second.url.clone()); + assert_eq!( + fetch(&state, DirectoryScope::MemberOnly).await.unwrap()[0].member_pubkeys, + vec![next, me] + ); + assert!(relay.roster_fallbacks().is_empty()); + assert!(second.roster_fallbacks().is_empty()); +} + +#[tokio::test] +async fn empty_membership_does_not_issue_metadata_roster_or_activity_queries() { + let _serial = crate::relay_admission::TEST_SERIAL.lock().await; + let keys = Keys::generate(); + let relay = Relay::new(vec![]).await; + assert!(fetch(&relay.state(&keys), DirectoryScope::MemberOnly) + .await + .unwrap() + .is_empty()); + assert_eq!(relay.data.lock().unwrap().requests.len(), 2); + assert!(relay.roster_fallbacks().is_empty()); +} + +#[tokio::test] +async fn discovery_and_activity_failures_still_abort_the_refresh() { + let _serial = crate::relay_admission::TEST_SERIAL.lock().await; + let keys = Keys::generate(); + let me = keys.public_key().to_hex(); + let relay = Relay::new(vec![ + metadata(&keys, "joined"), + roster(&keys, "joined", &[&me]), + ]) + .await; + let state = relay.state(&keys); + relay.data.lock().unwrap().fail_discovery = true; + assert!(fetch(&state, DirectoryScope::MemberOnly).await.is_err()); + relay.data.lock().unwrap().fail_discovery = false; + relay.data.lock().unwrap().fail_messages = true; + assert!(fetch(&state, DirectoryScope::MemberOnly).await.is_err()); +}