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()); +} 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); } }