1use nostr_sdk::prelude::*;
9use crate::ClientRelayExt;
10
11#[derive(Clone, Copy, Debug, Default, PartialEq, Eq)]
22pub enum Evidence {
23 Fast,
27 #[default]
34 Quorum,
35 Full,
40}
41
42#[derive(Clone, Debug, Default)]
46pub struct Query {
47 pub kinds: Vec<u16>,
48 pub z_tags: Vec<String>,
50 pub d_tags: Vec<String>,
53 pub p_tags: Vec<String>,
56 pub k_tags: Vec<String>,
59 pub authors: Vec<String>,
61 pub ids: Vec<String>,
64 pub since: Option<u64>,
66 pub until: Option<u64>,
69 pub limit: Option<usize>,
71 pub evidence: Evidence,
75}
76
77impl Query {
78 pub fn matches(&self, event: &Event) -> bool {
80 if !self.kinds.is_empty() && !self.kinds.iter().any(|k| Kind::Custom(*k) == event.kind) {
84 return false;
85 }
86 if let Some(since) = self.since {
87 if event.created_at.as_secs() < since {
88 return false;
89 }
90 }
91 if let Some(until) = self.until {
92 if event.created_at.as_secs() > until {
93 return false;
94 }
95 }
96 if !self.authors.is_empty() && !self.authors.iter().any(|a| *a == event.pubkey.to_hex()) {
97 return false;
98 }
99 if !self.ids.is_empty() && !self.ids.iter().any(|i| *i == event.id.to_hex()) {
100 return false;
101 }
102 if !self.z_tags.is_empty() && !self.matches_single_letter("z", &self.z_tags, event) {
103 return false;
104 }
105 if !self.d_tags.is_empty() && !self.matches_single_letter("d", &self.d_tags, event) {
106 return false;
107 }
108 if !self.p_tags.is_empty() && !self.matches_single_letter("p", &self.p_tags, event) {
109 return false;
110 }
111 if !self.k_tags.is_empty() && !self.matches_single_letter("k", &self.k_tags, event) {
112 return false;
113 }
114 true
115 }
116
117 fn matches_single_letter(&self, name: &str, wanted: &[String], event: &Event) -> bool {
118 event.tags.iter().any(|t| {
121 let s = t.as_slice();
122 s.len() >= 2 && s[0] == name && wanted.iter().any(|w| *w == s[1])
123 })
124 }
125
126 pub fn to_filter(&self) -> Filter {
128 let mut filter = Filter::new();
129 if !self.kinds.is_empty() {
130 filter = filter.kinds(self.kinds.iter().map(|k| Kind::Custom(*k)));
131 }
132 if !self.z_tags.is_empty() {
133 filter = filter
134 .custom_tags(SingleLetterTag::LOWERCASE_Z, self.z_tags.clone());
135 }
136 if !self.d_tags.is_empty() {
137 filter = filter.identifiers(self.d_tags.clone());
138 }
139 if !self.p_tags.is_empty() {
140 filter = filter
141 .custom_tags(SingleLetterTag::LOWERCASE_P, self.p_tags.clone());
142 }
143 if !self.k_tags.is_empty() {
144 filter = filter
145 .custom_tags(SingleLetterTag::LOWERCASE_K, self.k_tags.clone());
146 }
147 if !self.authors.is_empty() {
148 let authors: Vec<PublicKey> =
149 self.authors.iter().filter_map(|a| PublicKey::from_hex(a).ok()).collect();
150 if !authors.is_empty() {
151 filter = filter.authors(authors);
152 }
153 }
154 if !self.ids.is_empty() {
155 let ids: Vec<EventId> =
156 self.ids.iter().filter_map(|i| EventId::from_hex(i).ok()).collect();
157 if !ids.is_empty() {
158 filter = filter.ids(ids);
159 }
160 }
161 if let Some(since) = self.since {
162 filter = filter.since(Timestamp::from_secs(since));
163 }
164 if let Some(until) = self.until {
165 filter = filter.until(Timestamp::from_secs(until));
166 }
167 if let Some(limit) = self.limit {
168 filter = filter.limit(limit);
169 }
170 filter
171 }
172}
173
174#[async_trait::async_trait]
178pub trait Transport {
179 async fn publish(&self, event: &Event, relays: &[String]) -> Result<(), String>;
181 async fn fetch(&self, query: &Query, relays: &[String]) -> Result<Vec<Event>, String>;
183
184 async fn fetch_plane(&self, plane: &Keys, query: &Query, relays: &[String]) -> Result<Vec<Event>, String>;
193
194 async fn publish_durable(&self, event: &Event, relays: &[String]) -> Result<(), String>;
202}
203
204pub const MAX_PUBLISH_ATTEMPTS: usize = 30;
207
208pub const RESIDUAL_GRACE_MS: u64 = 400;
212
213pub const QUORUM_GRACE_MS: u64 = 2000;
219
220const BREAKER_TRIP_THRESHOLD: u8 = 2;
222
223const BREAKER_COOLDOWN: std::time::Duration = std::time::Duration::from_secs(30);
226
227const TRIPPED_TIMEOUT: std::time::Duration = std::time::Duration::from_secs(6);
231
232pub const CONFIRM_WINDOW: std::time::Duration = std::time::Duration::from_secs(30);
236
237pub async fn durable_broadcast<'a, F>(
244 relays: &[String],
245 max_attempts: usize,
246 backoff: std::time::Duration,
247 mut send_round: F,
248) -> Result<(), String>
249where
250 F: FnMut(Vec<String>) -> std::pin::Pin<Box<dyn std::future::Future<Output = Vec<String>> + Send + 'a>>,
251{
252 let mut pending: Vec<String> = Vec::new();
253 for r in relays {
254 if !pending.contains(r) {
255 pending.push(r.clone());
256 }
257 }
258 let total = pending.len();
259 if total == 0 {
260 return Err("no relays to broadcast to".to_string());
261 }
262 for attempt in 0..max_attempts {
263 if pending.is_empty() {
264 break;
265 }
266 let acked = send_round(pending.clone()).await;
267 pending.retain(|r| !acked.contains(r));
268 if pending.is_empty() || attempt + 1 == max_attempts {
269 break;
270 }
271 if !backoff.is_zero() {
272 tokio::time::sleep(backoff).await;
273 }
274 }
275 if pending.len() < total {
276 Ok(()) } else {
278 Err(format!("no relay accepted the event after {max_attempts} attempts each"))
279 }
280}
281
282pub trait CommunityIngestSink: Send + Sync + 'static {
291 fn ingest_stragglers(&self, events: Vec<Event>);
292}
293
294static INGEST_SINK: std::sync::OnceLock<Box<dyn CommunityIngestSink>> = std::sync::OnceLock::new();
295
296pub fn set_community_ingest_sink(sink: Box<dyn CommunityIngestSink>) {
298 let _ = INGEST_SINK.set(sink);
299}
300
301fn submit_stragglers(events: Vec<Event>) {
302 if events.is_empty() {
303 return;
304 }
305 if let Some(sink) = INGEST_SINK.get() {
306 sink.ingest_stragglers(events);
307 }
308}
309
310static WARMED_RELAYS: std::sync::LazyLock<std::sync::Mutex<(u64, std::collections::HashSet<String>)>> =
316 std::sync::LazyLock::new(|| std::sync::Mutex::new((0, std::collections::HashSet::new())));
317
318pub fn forget_warmed_relay(url: &str) {
323 WARMED_RELAYS.lock().unwrap_or_else(|e| e.into_inner()).1.remove(url);
324}
325
326#[derive(Default)]
333struct BreakerEntry {
334 consecutive_failures: u8,
335 tripped_until: Option<std::time::Instant>,
336}
337
338static RELAY_BREAKER: std::sync::LazyLock<
339 std::sync::Mutex<(u64, std::collections::HashMap<String, BreakerEntry>)>,
340> = std::sync::LazyLock::new(|| std::sync::Mutex::new((0, std::collections::HashMap::new())));
341
342fn with_breaker_at<R>(
347 generation: u64,
348 f: impl FnOnce(&mut std::collections::HashMap<String, BreakerEntry>) -> R,
349) -> R {
350 let mut guard = RELAY_BREAKER.lock().unwrap_or_else(|e| e.into_inner());
351 if guard.0 != generation {
352 guard.0 = generation;
353 guard.1.clear();
354 }
355 f(&mut guard.1)
356}
357
358fn drop_unrevivable(targets: Vec<String>, is_unrevivable: impl Fn(&str) -> bool) -> Vec<String> {
372 let alive: Vec<String> = targets.iter().filter(|r| !is_unrevivable(r)).cloned().collect();
373 if alive.is_empty() {
374 targets
375 } else {
376 alive
377 }
378}
379
380async fn drop_unrevivable_targets(client: &Client, targets: Vec<String>) -> Vec<String> {
382 use nostr_sdk::prelude::RelayStatus;
383 let pool = client.relays().all().await;
384 drop_unrevivable(targets, |r| {
385 RelayUrl::parse(r)
386 .ok()
387 .and_then(|u| pool.get(&u).map(|rl| matches!(rl.status(), RelayStatus::Terminated | RelayStatus::Shutdown | RelayStatus::Banned)))
388 .unwrap_or(false)
389 })
390}
391
392fn breaker_tripped(url: &str) -> bool {
393 breaker_tripped_at(crate::state::current_session_generation(), url)
394}
395
396fn breaker_tripped_at(generation: u64, url: &str) -> bool {
397 with_breaker_at(generation, |map| {
398 map.get(url)
399 .and_then(|e| e.tripped_until)
400 .map_or(false, |t| std::time::Instant::now() < t)
401 })
402}
403
404fn breaker_record(url: &str, success: bool, full_budget: bool) {
407 breaker_record_at(crate::state::current_session_generation(), url, success, full_budget)
408}
409
410fn breaker_record_at(generation: u64, url: &str, success: bool, full_budget: bool) {
411 with_breaker_at(generation, |map| {
412 if success {
413 map.remove(url);
414 return;
415 }
416 if !full_budget {
417 return;
418 }
419 let e = map.entry(url.to_string()).or_default();
420 e.consecutive_failures = e.consecutive_failures.saturating_add(1);
421 if e.consecutive_failures >= BREAKER_TRIP_THRESHOLD {
422 e.tripped_until = Some(std::time::Instant::now() + BREAKER_COOLDOWN);
423 }
424 })
425}
426
427fn effective_evidence(query: &Query) -> Evidence {
433 query.evidence
434}
435
436struct PooledPlane {
444 client: Client,
445 last_used: std::time::Instant,
446}
447
448static PLANE_POOL: std::sync::LazyLock<std::sync::Mutex<(u64, std::collections::HashMap<String, PooledPlane>)>> =
449 std::sync::LazyLock::new(|| std::sync::Mutex::new((0, std::collections::HashMap::new())));
450
451const PLANE_POOL_IDLE_TTL: std::time::Duration = std::time::Duration::from_secs(90);
454const PLANE_POOL_MAX: usize = 24;
456
457fn plane_pool_key(plane_pk: &str, relays: &[String]) -> String {
458 let mut rs: Vec<&str> = relays.iter().map(|s| s.as_str()).collect();
459 rs.sort_unstable();
460 let mut k = String::with_capacity(plane_pk.len() + 1 + rs.iter().map(|r| r.len() + 1).sum::<usize>());
461 k.push_str(plane_pk);
462 k.push('|');
463 k.push_str(&rs.join(","));
464 k
465}
466
467fn disconnect_clients(clients: Vec<Client>) {
472 if clients.is_empty() {
473 return;
474 }
475 match tokio::runtime::Handle::try_current() {
476 Ok(handle) => {
477 for c in clients {
478 handle.spawn(async move {
479 let _ = c.disconnect();
480 });
481 }
482 }
483 Err(_) => drop(clients),
484 }
485}
486
487pub fn clear_plane_pool() {
490 let drained: Vec<Client> = {
491 let mut g = PLANE_POOL.lock().unwrap_or_else(|e| e.into_inner());
492 g.0 = crate::state::current_session_generation();
497 g.1.drain().map(|(_, p)| p.client).collect()
498 };
499 disconnect_clients(drained);
500}
501
502fn plane_pool_take(generation: u64, key: &str) -> (Option<Client>, Vec<Client>) {
506 let mut g = PLANE_POOL.lock().unwrap_or_else(|e| e.into_inner());
507 let mut evicted: Vec<Client> = Vec::new();
508 if g.0 != generation {
509 evicted.extend(g.1.drain().map(|(_, p)| p.client));
510 g.0 = generation;
511 }
512 let now = std::time::Instant::now();
514 let expired: Vec<String> = g.1.iter()
515 .filter(|(_, p)| now.duration_since(p.last_used) >= PLANE_POOL_IDLE_TTL)
516 .map(|(k, _)| k.clone())
517 .collect();
518 for k in expired {
519 if let Some(p) = g.1.remove(&k) {
520 evicted.push(p.client);
521 }
522 }
523 let hit = g.1.get_mut(key).map(|p| {
524 p.last_used = now;
525 p.client.clone()
526 });
527 (hit, evicted)
528}
529
530fn plane_pool_insert(generation: u64, key: String, client: Client) -> Vec<Client> {
537 let mut g = PLANE_POOL.lock().unwrap_or_else(|e| e.into_inner());
538 if g.0 != generation || g.1.contains_key(&key) {
541 return Vec::new();
542 }
543 let mut evicted: Vec<Client> = Vec::new();
544 if g.1.len() >= PLANE_POOL_MAX {
545 if let Some(lru_key) = g.1.iter().min_by_key(|(_, p)| p.last_used).map(|(k, _)| k.clone()) {
546 if let Some(p) = g.1.remove(&lru_key) {
547 evicted.push(p.client);
548 }
549 }
550 }
551 g.1.insert(key, PooledPlane { client, last_used: std::time::Instant::now() });
552 evicted
553}
554
555fn demotion_allowed() -> bool {
559 #[cfg(feature = "tor")]
560 {
561 matches!(crate::tor::transport_state(), crate::tor::TorTransportState::Disabled)
562 }
563 #[cfg(not(feature = "tor"))]
564 {
565 true
566 }
567}
568
569pub async fn prewarm_held_communities(session: crate::state::SessionGuard) {
575 let mut relays: Vec<String> = Vec::new();
576 for id in crate::db::community::list_community_ids().unwrap_or_default() {
577 match crate::db::community::community_protocol(&id).ok().flatten() {
578 Some(crate::community::ConcordProtocol::V2) => {
579 if let Ok(Some(c)) = crate::db::community::load_community_v2(&id) {
580 relays.extend(c.relays.iter().cloned());
581 }
582 }
583 _ => {
584 if let Ok(Some(c)) = crate::db::community::load_community(&id) {
585 relays.extend(c.relays.iter().cloned());
586 }
587 }
588 }
589 }
590 relays.sort();
591 relays.dedup();
592 if relays.is_empty() || !session.is_valid() {
593 return;
594 }
595 let Ok(client) = LiveTransport::warm_client(&relays, std::time::Duration::from_secs(4)).await
596 else {
597 return;
598 };
599 crate::community::v2::streamauth::ensure_responder(&client);
605 let probe = Keys::generate();
606 let filter = Query {
607 kinds: vec![crate::community::v2::stream::KIND_WRAP],
608 authors: vec![probe.public_key().to_hex()],
609 limit: Some(1),
610 ..Default::default()
611 }
612 .to_filter();
613 let mut elicits = futures_util::stream::FuturesUnordered::new();
614 for r in &relays {
615 let c = client.clone();
616 let f = filter.clone();
617 let r = r.clone();
618 elicits.push(async move {
619 let _ = fetch_relay_eose_filters(&c, &r, vec![f], std::time::Duration::from_secs(3)).await;
620 });
621 }
622 use futures_util::StreamExt;
623 while elicits.next().await.is_some() {}
624}
625
626#[derive(Clone, Copy, PartialEq, Eq, Debug)]
630pub enum EoseFail {
631 Closed,
633 Deadline,
635 Gone,
637}
638
639pub async fn fetch_relay_eose(
650 client: &Client,
651 url: &str,
652 filter: Filter,
653 timeout: std::time::Duration,
654) -> Result<Vec<Event>, ()> {
655 fetch_relay_eose_filters(client, url, vec![filter], timeout).await.map_err(|_| ())
656}
657
658pub async fn fetch_relay_eose_filters(
663 client: &Client,
664 url: &str,
665 filters: Vec<Filter>,
666 timeout: std::time::Duration,
667) -> Result<Vec<Event>, EoseFail> {
668 let relay = client.relay(url).await.map_err(|_| EoseFail::Gone)?.ok_or(EoseFail::Gone)?;
669 let mut notifications = relay.notifications();
671 let sub_id = SubscriptionId::generate();
672 struct CloseGuard(Relay, SubscriptionId);
679 impl Drop for CloseGuard {
680 fn drop(&mut self) {
681 let relay = self.0.clone();
682 let id = self.1.clone();
683 tokio::spawn(async move {
684 let _ = relay
685 .send_msg(nostr_sdk::prelude::ClientMessage::Close(std::borrow::Cow::Owned(id)))
686 .await;
687 });
688 }
689 }
690 let _close = CloseGuard(relay.clone(), sub_id.clone());
691 relay
692 .send_msg(nostr_sdk::prelude::ClientMessage::Req {
693 subscription_id: std::borrow::Cow::Borrowed(&sub_id),
694 filters: filters.into_iter().map(std::borrow::Cow::Owned).collect(),
695 })
696 .await
697 .map_err(|_| EoseFail::Gone)?;
698 let deadline = tokio::time::Instant::now() + timeout;
699 let mut events: Vec<Event> = Vec::new();
700 let mut seen: std::collections::HashSet<EventId> = std::collections::HashSet::new();
701 loop {
702 let notification = match tokio::time::timeout_at(deadline, notifications.next()).await {
706 Ok(Some(n)) => n,
707 Ok(None) => return Err(EoseFail::Gone),
708 Err(_) => return Err(EoseFail::Deadline), };
710 match notification {
711 RelayNotification::Event { subscription_id, event } if subscription_id == sub_id => {
712 if seen.insert(event.id) {
713 events.push(*event);
714 }
715 }
716 RelayNotification::Message { message } => match *message {
717 RelayMessage::Event { subscription_id, event } if *subscription_id == sub_id => {
718 if seen.insert(event.id) {
719 events.push(event.into_owned());
720 }
721 }
722 RelayMessage::EndOfStoredEvents(id) if *id == sub_id => return Ok(events),
723 RelayMessage::Closed { subscription_id, .. } if *subscription_id == sub_id => {
724 return Err(EoseFail::Closed); }
726 _ => {}
727 },
728 RelayNotification::RelayStatus { status }
729 if status == nostr_sdk::prelude::RelayStatus::Shutdown =>
730 {
731 return Err(EoseFail::Gone);
732 }
733 _ => {}
734 }
735 }
736}
737
738pub(crate) struct UnionPlan {
742 attempted: usize,
743 successes: usize,
744 resolved: usize,
745 evidence: Evidence,
746}
747
748impl UnionPlan {
749 pub(crate) fn new(evidence: Evidence, attempted: usize) -> Self {
750 Self { attempted, successes: 0, resolved: 0, evidence }
751 }
752
753 pub(crate) fn record(&mut self, success: bool) {
754 self.resolved += 1;
755 if success {
756 self.successes += 1;
757 }
758 }
759
760 pub(crate) fn satisfied(&self) -> bool {
765 match self.evidence {
766 Evidence::Fast => self.successes >= 1,
767 Evidence::Quorum => self.successes >= (self.attempted / 2) + 1,
768 Evidence::Full => self.resolved >= self.attempted,
769 }
770 }
771
772 pub(crate) fn exhausted(&self) -> bool {
774 self.resolved >= self.attempted
775 }
776
777 pub(crate) fn successes(&self) -> usize {
778 self.successes
779 }
780
781 pub(crate) fn attempted(&self) -> usize {
782 self.attempted
783 }
784}
785
786#[cfg(feature = "tor")]
790const TOR_READY_WAIT: std::time::Duration = std::time::Duration::from_secs(30);
791
792#[allow(dead_code)]
802async fn wait_until_tor_ready<F: Fn() -> bool>(
803 is_blocked: F,
804 max_wait: std::time::Duration,
805) -> Result<(), String> {
806 if !is_blocked() {
807 return Ok(());
808 }
809 let deadline = std::time::Instant::now() + max_wait;
810 while is_blocked() {
811 if std::time::Instant::now() >= deadline {
812 return Err("Tor is still connecting. Wait a moment and try again.".to_string());
813 }
814 tokio::time::sleep(std::time::Duration::from_millis(250)).await;
815 }
816 Ok(())
817}
818
819pub async fn prune_unneeded_community_relays(candidates: &[String]) {
827 if candidates.is_empty() {
828 return;
829 }
830 let Some(client) = crate::state::nostr_client() else { return };
831
832 let mut still_needed: std::collections::HashSet<String> = std::collections::HashSet::new();
833 if let Ok(ids) = crate::db::community::list_community_ids() {
834 for id in ids {
835 if let Ok(Some(c)) = crate::db::community::load_community(&id) {
836 for r in &c.relays {
837 still_needed.insert(r.clone());
838 }
839 }
840 }
841 }
842
843 let pool = client;
844 let pooled = pool.relays().all().await;
846 for url in candidates {
847 if still_needed.contains(url) {
848 continue;
849 }
850 if let Ok(parsed) = nostr_sdk::prelude::RelayUrl::parse(url) {
851 if let Some(relay) = pooled.get(&parsed) {
852 if relay.capabilities().load().can_read() || relay.capabilities().load().can_write() {
853 continue; }
855 }
856 let _ = pool.remove_relay(parsed).force().await; forget_warmed_relay(url);
858 }
859 }
860}
861
862pub struct LiveTransport {
869 timeout: std::time::Duration,
870}
871
872impl Default for LiveTransport {
873 fn default() -> Self {
874 Self::with_timeout(std::time::Duration::from_secs(10))
875 }
876}
877
878impl LiveTransport {
879 pub fn new() -> Self {
880 Self::default()
881 }
882
883 pub fn with_timeout(timeout: std::time::Duration) -> Self {
889 Self { timeout: crate::relay_request_timeout(timeout) }
890 }
891
892 pub(crate) async fn warm_client(relays: &[String], connect_timeout: std::time::Duration) -> Result<Client, String> {
898 if relays.is_empty() {
899 return Err("community has no relays configured".to_string());
900 }
901 #[cfg(feature = "tor")]
907 wait_until_tor_ready(
908 || matches!(crate::tor::transport_state(), crate::tor::TorTransportState::RequiredButInactive),
909 TOR_READY_WAIT,
910 ).await?;
911 let client = crate::state::nostr_client().ok_or_else(|| "nostr client not initialized".to_string())?;
912
913 let generation = crate::state::current_session_generation();
917 {
918 let warmed = WARMED_RELAYS.lock().unwrap_or_else(|e| e.into_inner());
919 if warmed.0 == generation && relays.iter().all(|r| warmed.1.contains(r)) {
920 return Ok(client);
921 }
922 }
923
924 let mut added_new = false;
930 let mut succeeded: Vec<&String> = Vec::new();
931 for url in relays {
932 match client
933 .add_managed_relay(url.as_str())
934 .capabilities(crate::community_relay_capabilities())
935 .await
936 {
937 Ok(true) => { added_new = true; succeeded.push(url); }
938 Ok(false) => { succeeded.push(url); }
939 Err(_) => {}
940 }
941 }
942 if succeeded.is_empty() {
943 return Err("no valid community relays could be added".to_string());
944 }
945 if added_new {
946 let _ = client.try_connect().timeout(connect_timeout).await;
951 } else {
952 client.connect().await;
954 }
955
956 {
959 let mut warmed = WARMED_RELAYS.lock().unwrap_or_else(|e| e.into_inner());
960 if warmed.0 != generation {
961 warmed.0 = generation;
962 warmed.1.clear();
963 }
964 for url in succeeded {
965 warmed.1.insert(url.clone());
966 }
967 }
968 Ok(client)
969 }
970
971 pub async fn fetch_counted(&self, query: &Query, relays: &[String]) -> Result<(Vec<Event>, usize, usize), String> {
976 let client = Self::warm_client(relays, self.timeout).await?;
977 let base_timeout = self.timeout;
978 let filter = query.to_filter();
979
980 let evidence = effective_evidence(query);
981
982 let mut targets: Vec<String> = Vec::new();
983 for r in relays {
984 if !targets.contains(r) {
985 targets.push(r.clone());
986 }
987 }
988 targets = drop_unrevivable_targets(&client, targets).await;
993
994 if evidence == Evidence::Fast && targets.len() >= 2 {
999 let alive: Vec<String> =
1000 targets.iter().filter(|r| !breaker_tripped(r)).cloned().collect();
1001 if !alive.is_empty() {
1002 targets = alive;
1003 }
1004 }
1005
1006 if targets.len() <= 1 {
1009 let Some(url) = targets.first() else {
1010 return Err("no valid relay to fetch from".to_string());
1011 };
1012 let res = fetch_relay_eose(&client, url, filter, base_timeout)
1013 .await
1014 .map_err(|_| format!("relay did not answer the fetch: {url}"));
1015 breaker_record(url, res.is_ok(), true);
1016 return res.map(|evs| (evs, 1, 1));
1017 }
1018
1019 fn merge_events(
1020 evs: Vec<Event>,
1021 result: &mut Vec<Event>,
1022 seen: &mut std::collections::HashSet<EventId>,
1023 ) {
1024 for e in evs {
1025 if seen.insert(e.id) {
1026 result.push(e);
1027 }
1028 }
1029 }
1030
1031 let demote = evidence == Evidence::Full && demotion_allowed();
1038 use futures_util::stream::{FuturesUnordered, StreamExt};
1039 let mut fetches: FuturesUnordered<_> = targets
1040 .iter()
1041 .map(|r| {
1042 let client = client.clone();
1043 let filter = filter.clone();
1044 let r = r.clone();
1045 let timeout = if demote && breaker_tripped(&r) {
1046 TRIPPED_TIMEOUT.min(base_timeout)
1047 } else {
1048 base_timeout
1049 };
1050 let full_budget = timeout >= base_timeout;
1051 tokio::spawn(async move {
1052 let out = fetch_relay_eose(&client, &r, filter, timeout).await;
1053 (r, full_budget, out)
1054 })
1055 })
1056 .collect();
1057
1058 let mut plan = UnionPlan::new(evidence, targets.len());
1059 let mut result: Vec<Event> = Vec::new();
1060 let mut union_ids: std::collections::HashSet<EventId> = std::collections::HashSet::new();
1061
1062 let mut quorum_deadline: Option<tokio::time::Instant> = None;
1067 let mut quorum_window_closed = false;
1068 while !plan.satisfied() && !plan.exhausted() {
1069 let next = match quorum_deadline {
1070 Some(deadline) => match tokio::time::timeout_at(deadline, fetches.next()).await {
1071 Ok(n) => n,
1072 Err(_) => {
1073 quorum_window_closed = true;
1074 break; }
1076 },
1077 None => fetches.next().await,
1078 };
1079 let Some(joined) = next else { break };
1080 match joined {
1081 Ok((url, full_budget, Ok(evs))) => {
1082 breaker_record(&url, true, full_budget);
1083 merge_events(evs, &mut result, &mut union_ids);
1084 plan.record(true);
1085 if evidence == Evidence::Quorum && quorum_deadline.is_none() {
1086 quorum_deadline = Some(
1087 tokio::time::Instant::now()
1088 + std::time::Duration::from_millis(QUORUM_GRACE_MS),
1089 );
1090 }
1091 }
1092 Ok((url, full_budget, Err(()))) => {
1093 breaker_record(&url, false, full_budget);
1094 plan.record(false);
1095 }
1096 Err(_) => plan.record(false), }
1098 }
1099
1100 if plan.successes() == 0 {
1101 return Err(format!(
1102 "no relay answered the fetch (0/{} attempted)",
1103 plan.attempted()
1104 ));
1105 }
1106
1107 if !fetches.is_empty() && !quorum_window_closed {
1111 let grace = tokio::time::sleep(std::time::Duration::from_millis(RESIDUAL_GRACE_MS));
1112 tokio::pin!(grace);
1113 loop {
1114 tokio::select! {
1115 _ = &mut grace => break,
1116 next = fetches.next() => match next {
1117 Some(Ok((url, full_budget, Ok(evs)))) => {
1118 breaker_record(&url, true, full_budget);
1119 merge_events(evs, &mut result, &mut union_ids);
1120 }
1121 Some(Ok((url, full_budget, Err(())))) => {
1122 breaker_record(&url, false, full_budget);
1123 }
1124 Some(Err(_)) => continue,
1125 None => break,
1126 }
1127 }
1128 }
1129 }
1130
1131 if !fetches.is_empty() {
1135 let seen: std::collections::HashSet<EventId> = result.iter().map(|e| e.id).collect();
1136 let session = crate::state::SessionGuard::capture();
1139 tokio::spawn(async move {
1140 let mut extra: Vec<Event> = Vec::new();
1141 let mut extra_ids: std::collections::HashSet<EventId> = std::collections::HashSet::new();
1142 while let Some(joined) = fetches.next().await {
1143 if let Ok((url, full_budget, out)) = joined {
1144 breaker_record(&url, out.is_ok(), full_budget);
1148 if let Ok(evs) = out {
1149 for e in evs {
1150 if !seen.contains(&e.id) && extra_ids.insert(e.id) {
1151 extra.push(e);
1152 }
1153 }
1154 }
1155 }
1156 }
1157 if !session.is_valid() {
1158 return;
1159 }
1160 submit_stragglers(extra);
1161 });
1162 }
1163
1164 Ok((result, plan.successes(), plan.attempted()))
1165 }
1166}
1167
1168#[async_trait::async_trait]
1169impl Transport for LiveTransport {
1170 async fn publish(&self, event: &Event, relays: &[String]) -> Result<(), String> {
1171 let client = Self::warm_client(relays, self.timeout).await?;
1172 let timeout = self.timeout;
1173 let mut targets: Vec<String> = Vec::new();
1174 for r in relays { if !targets.contains(r) { targets.push(r.clone()); } }
1175 targets = drop_unrevivable_targets(&client, targets).await;
1177 use futures_util::stream::{FuturesUnordered, StreamExt};
1183 let mut sends: FuturesUnordered<_> = targets
1184 .into_iter()
1185 .map(|r| {
1186 let client = client.clone();
1187 let event = event.clone();
1188 tokio::spawn(async move {
1189 matches!(
1190 tokio::time::timeout(timeout, client.send_event(&event).to(vec![r.clone()])).await,
1191 Ok(Ok(out)) if RelayUrl::parse(&r).map(|u| out.success.contains_key(&u)).unwrap_or(false)
1192 )
1193 })
1194 })
1195 .collect();
1196 while let Some(joined) = sends.next().await {
1197 if matches!(joined, Ok(true)) {
1198 return Ok(()); }
1200 }
1201 Err("no relay accepted the event".to_string())
1202 }
1203
1204 async fn fetch(&self, query: &Query, relays: &[String]) -> Result<Vec<Event>, String> {
1205 self.fetch_counted(query, relays).await.map(|(events, _successes, _attempted)| events)
1206 }
1207
1208 async fn fetch_plane(&self, plane: &Keys, query: &Query, relays: &[String]) -> Result<Vec<Event>, String> {
1209 if relays.is_empty() {
1210 return Ok(Vec::new());
1211 }
1212 #[cfg(feature = "tor")]
1213 wait_until_tor_ready(
1214 || matches!(crate::tor::transport_state(), crate::tor::TorTransportState::RequiredButInactive),
1215 TOR_READY_WAIT,
1216 ).await?;
1217 let mut targets: Vec<String> = relays.iter().filter(|r| !breaker_tripped(r)).cloned().collect();
1222 if targets.is_empty() {
1223 targets = relays.to_vec();
1224 }
1225 let filter = query.to_filter();
1226 let generation = crate::state::current_session_generation();
1227 let key = plane_pool_key(&plane.public_key().to_hex(), &targets);
1228
1229 let (hit, evicted) = plane_pool_take(generation, &key);
1232 disconnect_clients(evicted);
1233 let client = if let Some(c) = hit {
1234 c
1235 } else {
1236 let client = crate::apply_tor_proxy(
1242 nostr_sdk::prelude::Client::builder().authenticator(
1243 nostr_sdk::prelude::SignerAuthenticator::new(plane.clone()),
1244 ),
1245 )
1246 .build();
1247 for r in &targets {
1252 let _ = client.add_managed_relay(r.clone()).capabilities(crate::community_relay_capabilities()).await;
1253 }
1254 client.connect().await;
1255 for r in &targets {
1258 let _ = client
1259 .fetch_events(nostr_sdk::prelude::ReqTarget::single(r.clone(), [filter.clone()]))
1260 .timeout(crate::relay_request_timeout(std::time::Duration::from_secs(5)))
1264 .await;
1265 }
1266 let ev = plane_pool_insert(generation, key, client.clone());
1267 disconnect_clients(ev);
1268 client
1269 };
1270
1271 let targets = drop_unrevivable_targets(&client, targets).await;
1275
1276 let mut result: Vec<Event> = Vec::new();
1277 let mut seen: std::collections::HashSet<EventId> = std::collections::HashSet::new();
1278 let mut successes = 0usize;
1279 for r in &targets {
1280 let res = fetch_relay_eose(&client, r, filter.clone(), self.timeout).await;
1281 breaker_record(r, res.is_ok(), true);
1284 if let Ok(events) = res {
1285 successes += 1;
1286 for e in events {
1287 if seen.insert(e.id) {
1288 result.push(e);
1289 }
1290 }
1291 }
1292 }
1293 if successes == 0 {
1298 return Err(format!("no relay answered the plane fetch (0/{} attempted)", targets.len()));
1299 }
1300 Ok(result)
1302 }
1303
1304 async fn publish_durable(&self, event: &Event, relays: &[String]) -> Result<(), String> {
1305 let client = Self::warm_client(relays, self.timeout).await?;
1310 let timeout = self.timeout;
1311 let event = event.clone();
1312 let backoff = std::time::Duration::from_millis(750);
1313 let mut pending: Vec<String> = Vec::new();
1314 for r in relays { if !pending.contains(r) { pending.push(r.clone()); } }
1315 if pending.is_empty() {
1316 return Err("no relays to broadcast to".to_string());
1317 }
1318
1319 let mut acked_any = false;
1324 let _ = tokio::time::timeout(CONFIRM_WINDOW, async {
1325 loop {
1326 let round = drop_unrevivable_targets(&client, pending.clone()).await;
1330 let sends = round.iter().cloned().map(|r| {
1331 let client = &client;
1332 let event = &event;
1333 Box::pin(async move {
1334 match tokio::time::timeout(timeout, client.send_event(event).to(vec![r.clone()])).await {
1335 Ok(Ok(out)) if RelayUrl::parse(&r).map(|u| out.success.contains_key(&u)).unwrap_or(false) => Ok(r),
1336 _ => Err(()),
1337 }
1338 })
1339 });
1340 if let Ok((winner, _losers)) = futures_util::future::select_ok(sends).await {
1341 acked_any = true;
1342 pending.retain(|r| r != &winner);
1343 break;
1344 }
1345 tokio::time::sleep(backoff).await;
1346 }
1347 })
1348 .await;
1349
1350 if !acked_any {
1351 return Err(format!("no relay accepted the event within {}s", CONFIRM_WINDOW.as_secs()));
1352 }
1353 if pending.is_empty() {
1354 return Ok(()); }
1356
1357 tokio::spawn(async move {
1362 let client_ref = &client;
1363 let event_ref = &event;
1364 let _ = durable_broadcast(&pending, MAX_PUBLISH_ATTEMPTS, backoff, move |round| {
1365 Box::pin(async move {
1366 match tokio::time::timeout(timeout, client_ref.send_event(event_ref).to(round.clone())).await {
1367 Ok(Ok(output)) => round.into_iter().filter(|p| RelayUrl::parse(p).map(|u| output.success.contains_key(&u)).unwrap_or(false)).collect(),
1368 _ => Vec::new(),
1369 }
1370 })
1371 })
1372 .await;
1373 });
1374 Ok(())
1375 }
1376}
1377
1378#[cfg(test)]
1379pub(crate) mod memory {
1380 use super::*;
1381 use std::collections::{HashMap, HashSet};
1382 use std::sync::Mutex;
1383
1384 pub struct MemoryRelay {
1389 per_relay: Mutex<HashMap<String, Vec<Event>>>,
1390 subscribers: Mutex<Vec<(Query, tokio::sync::mpsc::UnboundedSender<Event>)>>,
1391 }
1392
1393 fn is_ephemeral(kind: u16) -> bool {
1396 (20000..30000).contains(&kind)
1397 }
1398
1399 impl MemoryRelay {
1400 pub fn new() -> Self {
1401 MemoryRelay {
1402 per_relay: Mutex::new(HashMap::new()),
1403 subscribers: Mutex::new(Vec::new()),
1404 }
1405 }
1406
1407 pub fn stored_count(&self) -> usize {
1411 self.per_relay.lock().unwrap().values().map(|v| v.len()).sum()
1412 }
1413
1414 pub fn subscribe(&self, query: Query) -> tokio::sync::mpsc::UnboundedReceiver<Event> {
1417 let (tx, rx) = tokio::sync::mpsc::unbounded_channel();
1418 self.subscribers.lock().unwrap().push((query, tx));
1419 rx
1420 }
1421
1422 fn deliver(&self, event: &Event) {
1424 self.subscribers.lock().unwrap().retain(|(q, tx)| {
1425 if q.matches(event) {
1426 tx.send(event.clone()).is_ok()
1427 } else {
1428 !tx.is_closed()
1429 }
1430 });
1431 }
1432
1433 pub fn inject(&self, event: &Event, relays: &[String]) {
1436 self.deliver(event);
1437 if is_ephemeral(event.kind.as_u16()) {
1438 return; }
1440 let d_tag = |e: &Event| e.tags.iter().find_map(|t| {
1446 let s = t.as_slice();
1447 (s.len() >= 2 && s[0] == "d").then(|| s[1].clone())
1448 }).unwrap_or_default();
1449 let k = event.kind.as_u16();
1450 let replaceable = (30000..40000).contains(&k) || (10000..20000).contains(&k) || k == 0 || k == 3;
1451 let coord = (event.kind.as_u16(), event.pubkey, d_tag(event));
1452 let mut map = self.per_relay.lock().unwrap();
1453 for r in relays {
1454 let v = map.entry(r.clone()).or_default();
1455 if replaceable {
1456 v.retain(|e| (e.kind.as_u16(), e.pubkey, d_tag(e)) != coord);
1457 }
1458 v.push(event.clone());
1459 }
1460 }
1461
1462 pub fn count_on(&self, relay: &str) -> usize {
1464 self.per_relay.lock().unwrap().get(relay).map_or(0, |v| v.len())
1465 }
1466
1467 fn apply_deletion(&self, deletion: &Event, relays: &[String]) {
1472 let mut id_targets: HashSet<String> = HashSet::new();
1473 let mut coord_targets: HashSet<String> = HashSet::new();
1474 for t in deletion.tags.iter() {
1475 let s = t.as_slice();
1476 if s.len() >= 2 && s[0] == "e" {
1477 id_targets.insert(s[1].clone());
1478 } else if s.len() >= 2 && s[0] == "a" {
1479 coord_targets.insert(s[1].clone());
1480 }
1481 }
1482 let mut map = self.per_relay.lock().unwrap();
1483 for r in relays {
1484 if let Some(events) = map.get_mut(r) {
1485 events.retain(|e| {
1486 if e.pubkey != deletion.pubkey {
1487 return true;
1488 }
1489 if id_targets.contains(&e.id.to_hex()) {
1490 return false;
1491 }
1492 let d = e.tags.iter().find_map(|t| {
1494 let s = t.as_slice();
1495 (s.len() >= 2 && s[0] == "d").then(|| s[1].clone())
1496 }).unwrap_or_default();
1497 let coord = format!("{}:{}:{}", e.kind.as_u16(), e.pubkey.to_hex(), d);
1498 !coord_targets.contains(&coord)
1499 });
1500 }
1501 }
1502 }
1503 }
1504
1505 #[async_trait::async_trait]
1506 impl Transport for MemoryRelay {
1507 async fn publish(&self, event: &Event, relays: &[String]) -> Result<(), String> {
1508 if event.kind == Kind::EventDeletion {
1510 self.apply_deletion(event, relays);
1511 self.deliver(event);
1512 } else {
1513 self.inject(event, relays);
1514 }
1515 Ok(())
1516 }
1517
1518 async fn publish_durable(&self, event: &Event, relays: &[String]) -> Result<(), String> {
1521 self.publish(event, relays).await
1522 }
1523
1524 async fn fetch(&self, query: &Query, relays: &[String]) -> Result<Vec<Event>, String> {
1525 let map = self.per_relay.lock().unwrap();
1526 let mut seen = HashSet::new();
1527 let mut out = Vec::new();
1528 for r in relays {
1529 if let Some(events) = map.get(r) {
1530 for ev in events {
1531 if is_ephemeral(ev.kind.as_u16()) {
1534 continue;
1535 }
1536 if query.matches(ev) && seen.insert(ev.id) {
1537 out.push(ev.clone());
1538 }
1539 }
1540 }
1541 }
1542 if let Some(limit) = query.limit {
1545 out.sort_by(|a, b| b.created_at.cmp(&a.created_at).then_with(|| b.id.cmp(&a.id)));
1546 out.truncate(limit);
1547 }
1548 Ok(out)
1549 }
1550
1551 async fn fetch_plane(&self, _plane: &Keys, query: &Query, relays: &[String]) -> Result<Vec<Event>, String> {
1552 self.fetch(query, relays).await
1554 }
1555 }
1556}
1557
1558#[cfg(test)]
1559mod tests {
1560 use super::*;
1561
1562 #[test]
1565 fn a_dead_socket_is_skipped_but_a_connecting_or_unknown_one_is_not() {
1566 let targets: Vec<String> = ["wss://dead", "wss://connecting", "wss://unknown"].map(String::from).into();
1567 let out = drop_unrevivable(targets, |r| r == "wss://dead");
1568 assert_eq!(out, ["wss://connecting", "wss://unknown"].map(String::from).to_vec());
1569 }
1570
1571 #[test]
1572 fn an_all_dead_set_falls_back_to_the_full_list() {
1573 let targets: Vec<String> = ["wss://a", "wss://b"].map(String::from).into();
1576 let out = drop_unrevivable(targets.clone(), |_| true);
1577 assert_eq!(out, targets);
1578 }
1579
1580 #[tokio::test]
1585 async fn tor_gate_passes_immediately_when_not_blocked() {
1586 let start = std::time::Instant::now();
1587 let res = wait_until_tor_ready(|| false, std::time::Duration::from_secs(30)).await;
1588 assert!(res.is_ok());
1589 assert!(start.elapsed() < std::time::Duration::from_secs(1), "must not wait when Tor is ready");
1590 }
1591
1592 #[tokio::test]
1593 async fn tor_gate_errors_honestly_after_timeout_when_perpetually_blocked() {
1594 let res = wait_until_tor_ready(|| true, std::time::Duration::from_millis(300)).await;
1597 let err = res.expect_err("should error when Tor never activates");
1598 assert!(err.to_lowercase().contains("tor"), "error must name Tor, got: {err}");
1599 }
1600
1601 #[tokio::test]
1602 async fn tor_gate_passes_once_circuit_comes_up_mid_wait() {
1603 let calls = std::sync::atomic::AtomicUsize::new(0);
1604 let res = wait_until_tor_ready(
1606 || calls.fetch_add(1, std::sync::atomic::Ordering::SeqCst) < 3,
1607 std::time::Duration::from_secs(5),
1608 ).await;
1609 assert!(res.is_ok(), "should succeed once Tor activates within the window");
1610 }
1611
1612 fn evt(kind: u16, z: &str) -> Event {
1613 EventBuilder::new(Kind::Custom(kind), "x")
1614 .tags([Tag::custom(
1615 "z",
1616 [z.to_string()],
1617 )])
1618 .finalize(&Keys::generate())
1619 .unwrap()
1620 }
1621
1622 #[test]
1623 fn query_matches_kind_and_z() {
1624 let e = evt(3300, "abc");
1625 assert!(Query { kinds: vec![3300], z_tags: vec!["abc".into()], since: None, ..Default::default() }.matches(&e));
1626 assert!(!Query { kinds: vec![3301], ..Default::default() }.matches(&e));
1627 assert!(!Query { kinds: vec![], z_tags: vec!["xyz".into()], since: None, ..Default::default() }.matches(&e));
1628 assert!(Query::default().matches(&e), "empty query matches anything");
1629 }
1630
1631 fn evt_at(kind: u16, secs: u64) -> Event {
1632 EventBuilder::new(Kind::Custom(kind), "x")
1633 .custom_created_at(Timestamp::from(secs))
1634 .finalize(&Keys::generate())
1635 .unwrap()
1636 }
1637
1638 fn evt_z_at(kind: u16, z: &str, secs: u64) -> Event {
1641 EventBuilder::new(Kind::Custom(kind), "x")
1642 .custom_created_at(Timestamp::from(secs))
1643 .tags([Tag::custom(
1644 "z",
1645 [z.to_string()],
1646 )])
1647 .finalize(&Keys::generate())
1648 .unwrap()
1649 }
1650
1651 #[tokio::test]
1652 async fn fetch_pages_with_until_and_limit_newest_first() {
1653 let relay = super::memory::MemoryRelay::new();
1655 let relays = vec!["r1".to_string()];
1656 for s in 1..=5u64 {
1657 relay.inject(&evt_z_at(3300, "pg", s), &relays);
1658 }
1659 let secs = |evs: &[Event]| evs.iter().map(|e| e.created_at.as_secs()).collect::<Vec<_>>();
1660
1661 let latest = relay
1663 .fetch(&Query { kinds: vec![3300], z_tags: vec!["pg".into()], limit: Some(2), ..Default::default() }, &relays)
1664 .await
1665 .unwrap();
1666 assert_eq!(secs(&latest), vec![5, 4]);
1667
1668 let older = relay
1670 .fetch(&Query { kinds: vec![3300], z_tags: vec!["pg".into()], until: Some(3), limit: Some(2), ..Default::default() }, &relays)
1671 .await
1672 .unwrap();
1673 assert_eq!(secs(&older), vec![3, 2]);
1674
1675 let start = relay
1678 .fetch(&Query { kinds: vec![3300], z_tags: vec!["pg".into()], until: Some(1), limit: Some(2), ..Default::default() }, &relays)
1679 .await
1680 .unwrap();
1681 assert_eq!(secs(&start), vec![1]);
1682 }
1683
1684 #[test]
1685 fn to_filter_translates_kinds_z_and_since() {
1686 let q = Query { kinds: vec![3300], z_tags: vec!["abc".into()], since: Some(100), ..Default::default() };
1689 let filter = q.to_filter();
1690 let matching = EventBuilder::new(Kind::Custom(3300), "x")
1691 .custom_created_at(Timestamp::from(150))
1692 .tags([Tag::custom(
1693 "z",
1694 ["abc".to_string()],
1695 )])
1696 .finalize(&Keys::generate())
1697 .unwrap();
1698 assert!(filter.match_event(&matching, MatchEventOptions::new()), "to_filter must accept what matches() accepts");
1699 assert!(q.matches(&matching));
1700
1701 let wrong_kind = EventBuilder::new(Kind::Custom(3301), "x")
1703 .custom_created_at(Timestamp::from(150))
1704 .tags([Tag::custom(
1705 "z",
1706 ["abc".to_string()],
1707 )])
1708 .finalize(&Keys::generate())
1709 .unwrap();
1710 assert!(!filter.match_event(&wrong_kind, MatchEventOptions::new()));
1711 }
1712
1713 fn evt_sl(kind: u16, letter: SingleLetterTag, value: &str) -> Event {
1715 EventBuilder::new(Kind::Custom(kind), "x")
1716 .tags([Tag::custom(
1717 letter.as_char().to_string(),
1718 [value.to_string()],
1719 )])
1720 .finalize(&Keys::generate())
1721 .unwrap()
1722 }
1723
1724 #[test]
1725 fn to_filter_and_matches_agree_on_authors() {
1726 let keys = Keys::generate();
1727 let e = EventBuilder::new(Kind::Custom(1059), "x").finalize(&keys).unwrap();
1728 let q = Query { kinds: vec![1059], authors: vec![keys.public_key().to_hex()], ..Default::default() };
1729 assert!(q.matches(&e));
1730 assert!(q.to_filter().match_event(&e, MatchEventOptions::new()));
1731 let miss = Query { kinds: vec![1059], authors: vec![Keys::generate().public_key().to_hex()], ..Default::default() };
1732 assert!(!miss.matches(&e));
1733 assert!(!miss.to_filter().match_event(&e, MatchEventOptions::new()));
1734 }
1735
1736 #[test]
1737 fn to_filter_and_matches_agree_on_p_tags() {
1738 let recipient = Keys::generate().public_key().to_hex();
1739 let e = evt_sl(1059, SingleLetterTag::LOWERCASE_P, &recipient);
1740 let q = Query { kinds: vec![1059], p_tags: vec![recipient], ..Default::default() };
1741 assert!(q.matches(&e));
1742 assert!(q.to_filter().match_event(&e, MatchEventOptions::new()));
1743 let miss = Query { kinds: vec![1059], p_tags: vec![Keys::generate().public_key().to_hex()], ..Default::default() };
1744 assert!(!miss.matches(&e));
1745 assert!(!miss.to_filter().match_event(&e, MatchEventOptions::new()));
1746 }
1747
1748 #[test]
1749 fn to_filter_and_matches_agree_on_k_tags() {
1750 let e = evt_sl(1059, SingleLetterTag::LOWERCASE_K, "3311");
1751 let q = Query { kinds: vec![1059], k_tags: vec!["3311".into()], ..Default::default() };
1752 assert!(q.matches(&e));
1753 assert!(q.to_filter().match_event(&e, MatchEventOptions::new()));
1754 let miss = Query { kinds: vec![1059], k_tags: vec!["3300".into()], ..Default::default() };
1755 assert!(!miss.matches(&e));
1756 assert!(!miss.to_filter().match_event(&e, MatchEventOptions::new()));
1757 }
1758
1759 #[test]
1760 fn to_filter_empty_kinds_only_constrains_z_and_since() {
1761 let q = Query { kinds: vec![], z_tags: vec!["p".into()], since: None, ..Default::default() };
1763 let filter = q.to_filter();
1764 let e = evt(3300, "p");
1765 assert!(filter.match_event(&e, MatchEventOptions::new()));
1766 assert!(q.matches(&e));
1767 }
1768
1769 #[test]
1770 fn since_is_an_inclusive_lower_bound() {
1771 let below = evt_at(3300, 99);
1772 let exact = evt_at(3300, 100);
1773 let above = evt_at(3300, 101);
1774 let q = Query { kinds: vec![3300], z_tags: vec![], since: Some(100), ..Default::default() };
1775 assert!(!q.matches(&below), "below the floor is excluded");
1776 assert!(q.matches(&exact), "exactly the floor is included");
1777 assert!(q.matches(&above), "above the floor is included");
1778 }
1779
1780 #[test]
1781 fn z_tags_match_as_or_set() {
1782 let e = evt(3300, "p2");
1785 let q = Query { kinds: vec![3300], z_tags: vec!["p1".into(), "p2".into()], since: None, ..Default::default() };
1786 assert!(q.matches(&e));
1787 let miss = Query { kinds: vec![3300], z_tags: vec!["p1".into(), "p3".into()], since: None, ..Default::default() };
1788 assert!(!miss.matches(&e));
1789 }
1790
1791 #[tokio::test]
1792 async fn fetch_unions_and_dedups_across_relays() {
1793 use super::memory::MemoryRelay;
1794 let relay = MemoryRelay::new();
1795 let relays = vec!["r1".to_string(), "r2".to_string(), "r3".to_string()];
1796 let e = evt(3300, "p");
1797 relay.publish(&e, &relays).await.unwrap();
1798 let got = relay
1799 .fetch(&Query { kinds: vec![3300], z_tags: vec!["p".into()], since: None, ..Default::default() }, &relays)
1800 .await
1801 .unwrap();
1802 assert_eq!(got.len(), 1, "same event on 3 relays dedups to 1");
1803 }
1804
1805 #[tokio::test]
1806 async fn durable_broadcast_retries_only_the_failing_relays_until_they_ack() {
1807 use std::cell::Cell;
1810 let relays = vec!["r1".to_string(), "r2".to_string(), "r3".to_string()];
1811 let round = Cell::new(0usize);
1812 let r2_round_seen = Cell::new(0usize);
1813 let res = durable_broadcast(&relays, 30, std::time::Duration::ZERO, |pending| {
1814 let n = round.get();
1815 round.set(n + 1);
1816 if n >= 1 {
1818 assert_eq!(pending, vec!["r2".to_string()], "only the failing relay is retried");
1819 r2_round_seen.set(r2_round_seen.get() + 1);
1820 }
1821 Box::pin(async move {
1822 pending.into_iter().filter(|r| r != "r2" || n >= 4).collect()
1823 })
1824 })
1825 .await;
1826 assert!(res.is_ok(), "all relays eventually ACK → Ok");
1827 assert!(round.get() >= 5, "kept retrying r2 across rounds");
1828 }
1829
1830 #[tokio::test]
1831 async fn durable_broadcast_is_ok_if_some_ack_even_when_one_never_does() {
1832 let relays = vec!["r1".to_string(), "r2".to_string()];
1835 let res = durable_broadcast(&relays, 5, std::time::Duration::ZERO, |pending| {
1836 Box::pin(async move { pending.into_iter().filter(|r| r == "r1").collect() })
1837 })
1838 .await;
1839 assert!(res.is_ok(), "≥1 relay accepted → Ok despite a permanently-failing relay");
1840 }
1841
1842 #[tokio::test]
1843 async fn durable_broadcast_errs_only_if_zero_relays_ever_accept() {
1844 let relays = vec!["r1".to_string(), "r2".to_string()];
1845 let res = durable_broadcast(&relays, 5, std::time::Duration::ZERO, |_pending| {
1846 Box::pin(async move { Vec::new() }) })
1848 .await;
1849 assert!(res.is_err(), "zero acceptances after the retry cap → Err");
1850 }
1851
1852 #[tokio::test]
1853 async fn redundancy_self_heals_a_missing_relay() {
1854 use super::memory::MemoryRelay;
1855 let relay = MemoryRelay::new();
1856 let all = vec!["r1".to_string(), "r2".to_string(), "r3".to_string()];
1857 let e = evt(3300, "p");
1858 relay.inject(&e, &["r2".to_string()]);
1860 assert_eq!(relay.count_on("r1"), 0);
1861 assert_eq!(relay.count_on("r2"), 1);
1862 let got = relay.fetch(&Query { kinds: vec![3300], ..Default::default() }, &all).await.unwrap();
1864 assert_eq!(got.len(), 1);
1865 }
1866
1867 #[tokio::test]
1868 async fn ephemeral_kind_streams_live_but_is_never_stored_or_fetched() {
1869 use super::memory::MemoryRelay;
1870 let relay = MemoryRelay::new();
1871 let relays = vec!["r1".to_string()];
1872 let mut sub = relay.subscribe(Query { kinds: vec![21059], ..Default::default() });
1873 let e = evt(21059, "p");
1874 relay.publish(&e, &relays).await.unwrap();
1875 assert_eq!(relay.count_on("r1"), 0, "ephemeral is never stored");
1876 let got = relay
1877 .fetch(&Query { kinds: vec![21059], ..Default::default() }, &relays)
1878 .await
1879 .unwrap();
1880 assert!(got.is_empty(), "a real relay never serves an ephemeral from a fetch");
1881 assert_eq!(sub.try_recv().unwrap().id, e.id, "but a live subscriber receives it");
1882 }
1883
1884 #[tokio::test]
1885 async fn stored_kind_is_fetchable_and_delivered_live() {
1886 use super::memory::MemoryRelay;
1887 let relay = MemoryRelay::new();
1888 let relays = vec!["r1".to_string()];
1889 let mut sub = relay.subscribe(Query { kinds: vec![1059], ..Default::default() });
1890 let e = evt(1059, "p");
1891 relay.publish(&e, &relays).await.unwrap();
1892 let got = relay
1893 .fetch(&Query { kinds: vec![1059], ..Default::default() }, &relays)
1894 .await
1895 .unwrap();
1896 assert_eq!(got.len(), 1, "stored kind is fetchable");
1897 assert_eq!(sub.try_recv().unwrap().id, e.id, "and delivered to the live subscriber");
1898 }
1899
1900 #[tokio::test]
1901 async fn p_tags_route_a_giftwrap_to_the_matching_subscriber() {
1902 use super::memory::MemoryRelay;
1903 let relay = MemoryRelay::new();
1904 let relays = vec!["r1".to_string()];
1905 let alice = Keys::generate().public_key().to_hex();
1906 let bob = Keys::generate().public_key().to_hex();
1907 let mut sub_alice =
1908 relay.subscribe(Query { kinds: vec![1059], p_tags: vec![alice.clone()], ..Default::default() });
1909 let mut sub_bob =
1910 relay.subscribe(Query { kinds: vec![1059], p_tags: vec![bob.clone()], ..Default::default() });
1911 let wrap = evt_sl(1059, SingleLetterTag::LOWERCASE_P, &alice);
1912 relay.publish(&wrap, &relays).await.unwrap();
1913 assert_eq!(sub_alice.try_recv().unwrap().id, wrap.id, "addressed recipient gets it live");
1914 assert!(sub_bob.try_recv().is_err(), "a differently-addressed subscriber does not");
1915 let for_alice = relay
1917 .fetch(&Query { kinds: vec![1059], p_tags: vec![alice], ..Default::default() }, &relays)
1918 .await
1919 .unwrap();
1920 assert_eq!(for_alice.len(), 1);
1921 let for_bob = relay
1922 .fetch(&Query { kinds: vec![1059], p_tags: vec![bob], ..Default::default() }, &relays)
1923 .await
1924 .unwrap();
1925 assert!(for_bob.is_empty());
1926 }
1927
1928 #[test]
1931 fn union_plan_fast_satisfied_on_first_success() {
1932 let mut p = UnionPlan::new(Evidence::Fast, 4);
1933 p.record(false);
1934 assert!(!p.satisfied(), "a failure is not evidence");
1935 p.record(true);
1936 assert!(p.satisfied(), "one genuine EOSE satisfies Fast");
1937 assert!(!p.exhausted());
1938 }
1939
1940 #[test]
1941 fn union_plan_quorum_majority_math() {
1942 for (n, need) in [(2usize, 2usize), (3, 2), (4, 3), (5, 3)] {
1944 let mut p = UnionPlan::new(Evidence::Quorum, n);
1945 for _ in 0..need - 1 {
1946 p.record(true);
1947 }
1948 assert!(!p.satisfied(), "{}/{} must not satisfy quorum", need - 1, n);
1949 p.record(true);
1950 assert!(p.satisfied(), "{}/{} satisfies quorum", need, n);
1951 }
1952 }
1953
1954 #[test]
1955 fn union_plan_quorum_failures_never_substitute_for_successes() {
1956 let mut p = UnionPlan::new(Evidence::Quorum, 3);
1957 p.record(true);
1958 p.record(false);
1959 p.record(false);
1960 assert!(!p.satisfied(), "1 success + 2 failures is not a majority");
1961 assert!(p.exhausted(), "all resolved — the degraded path returns best-effort");
1962 assert_eq!(p.successes(), 1);
1963 }
1964
1965 #[test]
1966 fn union_plan_full_requires_every_relay_resolved() {
1967 let mut p = UnionPlan::new(Evidence::Full, 3);
1968 p.record(true);
1969 p.record(true);
1970 assert!(!p.satisfied(), "Full waits for the last relay even after 2 EOSEs");
1971 p.record(false);
1972 assert!(p.satisfied(), "a timeout is a resolution — Full is done");
1973 assert!(p.exhausted());
1974 }
1975
1976 #[test]
1977 fn union_plan_all_dead_is_reportable_not_a_confident_empty() {
1978 let mut p = UnionPlan::new(Evidence::Quorum, 2);
1979 p.record(false);
1980 p.record(false);
1981 assert!(p.exhausted());
1982 assert_eq!(p.successes(), 0, "the caller must map this to Err, never Ok(vec![])");
1983 }
1984
1985 const BREAKER_TEST_GEN: u64 = u64::MAX;
1991
1992 #[test]
1993 fn breaker_trips_only_after_consecutive_full_budget_failures() {
1994 let url = "wss://breaker-test-full-budget.example";
1995 breaker_record_at(BREAKER_TEST_GEN, url, false, true);
1996 assert!(!breaker_tripped_at(BREAKER_TEST_GEN, url), "one failure is below the threshold");
1997 breaker_record_at(BREAKER_TEST_GEN, url, false, true);
1998 assert!(breaker_tripped_at(BREAKER_TEST_GEN, url), "two consecutive full-budget failures trip");
1999 }
2000
2001 #[test]
2002 fn breaker_demoted_budget_failures_never_count() {
2003 let url = "wss://breaker-test-demoted.example";
2004 breaker_record_at(BREAKER_TEST_GEN, url, false, false);
2005 breaker_record_at(BREAKER_TEST_GEN, url, false, false);
2006 breaker_record_at(BREAKER_TEST_GEN, url, false, false);
2007 assert!(
2008 !breaker_tripped_at(BREAKER_TEST_GEN, url),
2009 "demoted-budget failures must not trip (anti-starvation: the post-cooldown probe must stay reachable)"
2010 );
2011 }
2012
2013 #[test]
2014 fn breaker_success_resets_the_entry() {
2015 let url = "wss://breaker-test-reset.example";
2016 breaker_record_at(BREAKER_TEST_GEN, url, false, true);
2017 breaker_record_at(BREAKER_TEST_GEN, url, false, true);
2018 assert!(breaker_tripped_at(BREAKER_TEST_GEN, url));
2019 breaker_record_at(BREAKER_TEST_GEN, url, true, false);
2021 assert!(!breaker_tripped_at(BREAKER_TEST_GEN, url), "any success unconditionally resets");
2022 breaker_record_at(BREAKER_TEST_GEN, url, false, true);
2023 assert!(!breaker_tripped_at(BREAKER_TEST_GEN, url), "and the failure count restarted from zero");
2024 }
2025
2026 #[test]
2029 fn declared_evidence_stands_and_default_is_quorum() {
2030 assert_eq!(Query::default().evidence, Evidence::Quorum, "unclassified sites get Quorum");
2031 assert_eq!(
2032 effective_evidence(&Query { until: Some(1), evidence: Evidence::Fast, ..Default::default() }),
2033 Evidence::Fast,
2034 "chat pagination rides its declared tier — absence verdicts request Full themselves"
2035 );
2036 assert_eq!(
2037 effective_evidence(&Query { evidence: Evidence::Fast, ..Default::default() }),
2038 Evidence::Fast,
2039 "without `until` the declared tier stands"
2040 );
2041 }
2042}