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::db::current_session_id(), 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::db::current_session_id(), 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::db::current_session_id();
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() {
575 crate::db::scoped(async move {
576 let mut relays: Vec<String> = Vec::new();
577 for id in crate::db::community::list_community_ids().unwrap_or_default() {
578 match crate::db::community::community_protocol(&id).ok().flatten() {
579 Some(crate::community::ConcordProtocol::V2) => {
580 if let Ok(Some(c)) = crate::db::community::load_community_v2(&id) {
581 relays.extend(c.relays.iter().cloned());
582 }
583 }
584 _ => {
585 if let Ok(Some(c)) = crate::db::community::load_community(&id) {
586 relays.extend(c.relays.iter().cloned());
587 }
588 }
589 }
590 }
591 relays.sort();
592 relays.dedup();
593 if relays.is_empty() {
594 return;
595 }
596 let Ok(client) = LiveTransport::warm_client(&relays, std::time::Duration::from_secs(4)).await
597 else {
598 return;
599 };
600 crate::community::v2::streamauth::ensure_responder(&client);
606 let probe = Keys::generate();
607 let filter = Query {
608 kinds: vec![crate::community::v2::stream::KIND_WRAP],
609 authors: vec![probe.public_key().to_hex()],
610 limit: Some(1),
611 ..Default::default()
612 }
613 .to_filter();
614 let mut elicits = futures_util::stream::FuturesUnordered::new();
615 for r in &relays {
616 let c = client.clone();
617 let f = filter.clone();
618 let r = r.clone();
619 elicits.push(async move {
620 let _ = fetch_relay_eose_filters(&c, &r, vec![f], std::time::Duration::from_secs(3)).await;
621 });
622 }
623 use futures_util::StreamExt;
624 while elicits.next().await.is_some() {}
625 })
626 .await
627}
628
629#[derive(Clone, Copy, PartialEq, Eq, Debug)]
633pub enum EoseFail {
634 Closed,
636 Deadline,
638 Gone,
640}
641
642pub async fn fetch_relay_eose(
653 client: &Client,
654 url: &str,
655 filter: Filter,
656 timeout: std::time::Duration,
657) -> Result<Vec<Event>, ()> {
658 fetch_relay_eose_filters(client, url, vec![filter], timeout).await.map_err(|_| ())
659}
660
661pub async fn fetch_relay_eose_filters(
666 client: &Client,
667 url: &str,
668 filters: Vec<Filter>,
669 timeout: std::time::Duration,
670) -> Result<Vec<Event>, EoseFail> {
671 let relay = client.relay(url).await.map_err(|_| EoseFail::Gone)?.ok_or(EoseFail::Gone)?;
672 let mut notifications = relay.notifications();
674 let sub_id = SubscriptionId::generate();
675 struct CloseGuard(Relay, SubscriptionId);
682 impl Drop for CloseGuard {
683 fn drop(&mut self) {
684 let relay = self.0.clone();
685 let id = self.1.clone();
686 tokio::spawn(async move {
688 let _ = relay
689 .send_msg(nostr_sdk::prelude::ClientMessage::Close(std::borrow::Cow::Owned(id)))
690 .await;
691 });
692 }
693 }
694 let _close = CloseGuard(relay.clone(), sub_id.clone());
695 relay
696 .send_msg(nostr_sdk::prelude::ClientMessage::Req {
697 subscription_id: std::borrow::Cow::Borrowed(&sub_id),
698 filters: filters.into_iter().map(std::borrow::Cow::Owned).collect(),
699 })
700 .await
701 .map_err(|_| EoseFail::Gone)?;
702 let deadline = tokio::time::Instant::now() + timeout;
703 let mut events: Vec<Event> = Vec::new();
704 let mut seen: std::collections::HashSet<EventId> = std::collections::HashSet::new();
705 loop {
706 let notification = match tokio::time::timeout_at(deadline, notifications.next()).await {
710 Ok(Some(n)) => n,
711 Ok(None) => return Err(EoseFail::Gone),
712 Err(_) => return Err(EoseFail::Deadline), };
714 match notification {
715 RelayNotification::Event { subscription_id, event } if subscription_id == sub_id => {
716 if seen.insert(event.id) {
717 events.push(*event);
718 }
719 }
720 RelayNotification::Message { message } => match *message {
721 RelayMessage::Event { subscription_id, event } if *subscription_id == sub_id => {
722 if seen.insert(event.id) {
723 events.push(event.into_owned());
724 }
725 }
726 RelayMessage::EndOfStoredEvents(id) if *id == sub_id => return Ok(events),
727 RelayMessage::Closed { subscription_id, .. } if *subscription_id == sub_id => {
728 return Err(EoseFail::Closed); }
730 _ => {}
731 },
732 RelayNotification::RelayStatus { status }
733 if status == nostr_sdk::prelude::RelayStatus::Shutdown =>
734 {
735 return Err(EoseFail::Gone);
736 }
737 _ => {}
738 }
739 }
740}
741
742pub(crate) struct UnionPlan {
746 attempted: usize,
747 successes: usize,
748 resolved: usize,
749 evidence: Evidence,
750}
751
752impl UnionPlan {
753 pub(crate) fn new(evidence: Evidence, attempted: usize) -> Self {
754 Self { attempted, successes: 0, resolved: 0, evidence }
755 }
756
757 pub(crate) fn record(&mut self, success: bool) {
758 self.resolved += 1;
759 if success {
760 self.successes += 1;
761 }
762 }
763
764 pub(crate) fn satisfied(&self) -> bool {
769 match self.evidence {
770 Evidence::Fast => self.successes >= 1,
771 Evidence::Quorum => self.successes >= (self.attempted / 2) + 1,
772 Evidence::Full => self.resolved >= self.attempted,
773 }
774 }
775
776 pub(crate) fn exhausted(&self) -> bool {
778 self.resolved >= self.attempted
779 }
780
781 pub(crate) fn successes(&self) -> usize {
782 self.successes
783 }
784
785 pub(crate) fn attempted(&self) -> usize {
786 self.attempted
787 }
788}
789
790#[cfg(feature = "tor")]
794const TOR_READY_WAIT: std::time::Duration = std::time::Duration::from_secs(30);
795
796#[allow(dead_code)]
806async fn wait_until_tor_ready<F: Fn() -> bool>(
807 is_blocked: F,
808 max_wait: std::time::Duration,
809) -> Result<(), String> {
810 if !is_blocked() {
811 return Ok(());
812 }
813 let deadline = std::time::Instant::now() + max_wait;
814 while is_blocked() {
815 if std::time::Instant::now() >= deadline {
816 return Err("Tor is still connecting. Wait a moment and try again.".to_string());
817 }
818 tokio::time::sleep(std::time::Duration::from_millis(250)).await;
819 }
820 Ok(())
821}
822
823pub async fn prune_unneeded_community_relays(candidates: &[String]) {
831 if candidates.is_empty() {
832 return;
833 }
834 let Some(client) = crate::state::nostr_client() else { return };
835
836 let mut still_needed: std::collections::HashSet<String> = std::collections::HashSet::new();
837 if let Ok(ids) = crate::db::community::list_community_ids() {
838 for id in ids {
839 if let Ok(Some(c)) = crate::db::community::load_community(&id) {
840 for r in &c.relays {
841 still_needed.insert(r.clone());
842 }
843 }
844 }
845 }
846
847 let pool = client;
848 let pooled = pool.relays().all().await;
850 for url in candidates {
851 if still_needed.contains(url) {
852 continue;
853 }
854 if let Ok(parsed) = nostr_sdk::prelude::RelayUrl::parse(url) {
855 if let Some(relay) = pooled.get(&parsed) {
856 if relay.capabilities().load().can_read() || relay.capabilities().load().can_write() {
857 continue; }
859 }
860 let _ = pool.remove_relay(parsed).force().await; forget_warmed_relay(url);
862 }
863 }
864}
865
866pub struct LiveTransport {
873 timeout: std::time::Duration,
874}
875
876impl Default for LiveTransport {
877 fn default() -> Self {
878 Self::with_timeout(std::time::Duration::from_secs(10))
879 }
880}
881
882impl LiveTransport {
883 pub fn new() -> Self {
884 Self::default()
885 }
886
887 pub fn with_timeout(timeout: std::time::Duration) -> Self {
893 Self { timeout: crate::relay_request_timeout(timeout) }
894 }
895
896 pub(crate) async fn warm_client(relays: &[String], connect_timeout: std::time::Duration) -> Result<Client, String> {
902 if relays.is_empty() {
903 return Err("community has no relays configured".to_string());
904 }
905 #[cfg(feature = "tor")]
911 wait_until_tor_ready(
912 || matches!(crate::tor::transport_state(), crate::tor::TorTransportState::RequiredButInactive),
913 TOR_READY_WAIT,
914 ).await?;
915 let client = crate::state::nostr_client().ok_or_else(|| "nostr client not initialized".to_string())?;
916
917 let generation = crate::db::current_session_id();
921 {
922 let warmed = WARMED_RELAYS.lock().unwrap_or_else(|e| e.into_inner());
923 if warmed.0 == generation && relays.iter().all(|r| warmed.1.contains(r)) {
924 return Ok(client);
925 }
926 }
927
928 let mut added_new = false;
934 let mut succeeded: Vec<&String> = Vec::new();
935 for url in relays {
936 match client
937 .add_managed_relay(url.as_str())
938 .capabilities(crate::community_relay_capabilities())
939 .await
940 {
941 Ok(true) => { added_new = true; succeeded.push(url); }
942 Ok(false) => { succeeded.push(url); }
943 Err(_) => {}
944 }
945 }
946 if succeeded.is_empty() {
947 return Err("no valid community relays could be added".to_string());
948 }
949 if added_new {
950 let _ = client.try_connect().timeout(connect_timeout).await;
955 } else {
956 client.connect().await;
958 }
959
960 {
963 let mut warmed = WARMED_RELAYS.lock().unwrap_or_else(|e| e.into_inner());
964 if warmed.0 != generation {
965 warmed.0 = generation;
966 warmed.1.clear();
967 }
968 for url in succeeded {
969 warmed.1.insert(url.clone());
970 }
971 }
972 Ok(client)
973 }
974
975 pub async fn fetch_counted(&self, query: &Query, relays: &[String]) -> Result<(Vec<Event>, usize, usize), String> {
980 let client = Self::warm_client(relays, self.timeout).await?;
981 let base_timeout = self.timeout;
982 let filter = query.to_filter();
983
984 let evidence = effective_evidence(query);
985
986 let mut targets: Vec<String> = Vec::new();
987 for r in relays {
988 if !targets.contains(r) {
989 targets.push(r.clone());
990 }
991 }
992 targets = drop_unrevivable_targets(&client, targets).await;
997
998 if evidence == Evidence::Fast && targets.len() >= 2 {
1003 let alive: Vec<String> =
1004 targets.iter().filter(|r| !breaker_tripped(r)).cloned().collect();
1005 if !alive.is_empty() {
1006 targets = alive;
1007 }
1008 }
1009
1010 if targets.len() <= 1 {
1013 let Some(url) = targets.first() else {
1014 return Err("no valid relay to fetch from".to_string());
1015 };
1016 let res = fetch_relay_eose(&client, url, filter, base_timeout)
1017 .await
1018 .map_err(|_| format!("relay did not answer the fetch: {url}"));
1019 breaker_record(url, res.is_ok(), true);
1020 return res.map(|evs| (evs, 1, 1));
1021 }
1022
1023 fn merge_events(
1024 evs: Vec<Event>,
1025 result: &mut Vec<Event>,
1026 seen: &mut std::collections::HashSet<EventId>,
1027 ) {
1028 for e in evs {
1029 if seen.insert(e.id) {
1030 result.push(e);
1031 }
1032 }
1033 }
1034
1035 let demote = evidence == Evidence::Full && demotion_allowed();
1042 use futures_util::stream::{FuturesUnordered, StreamExt};
1043 let mut fetches: FuturesUnordered<_> = targets
1044 .iter()
1045 .map(|r| {
1046 let client = client.clone();
1047 let filter = filter.clone();
1048 let r = r.clone();
1049 let timeout = if demote && breaker_tripped(&r) {
1050 TRIPPED_TIMEOUT.min(base_timeout)
1051 } else {
1052 base_timeout
1053 };
1054 let full_budget = timeout >= base_timeout;
1055 tokio::spawn(async move {
1057 let out = fetch_relay_eose(&client, &r, filter, timeout).await;
1058 (r, full_budget, out)
1059 })
1060 })
1061 .collect();
1062
1063 let mut plan = UnionPlan::new(evidence, targets.len());
1064 let mut result: Vec<Event> = Vec::new();
1065 let mut union_ids: std::collections::HashSet<EventId> = std::collections::HashSet::new();
1066
1067 let mut quorum_deadline: Option<tokio::time::Instant> = None;
1072 let mut quorum_window_closed = false;
1073 while !plan.satisfied() && !plan.exhausted() {
1074 let next = match quorum_deadline {
1075 Some(deadline) => match tokio::time::timeout_at(deadline, fetches.next()).await {
1076 Ok(n) => n,
1077 Err(_) => {
1078 quorum_window_closed = true;
1079 break; }
1081 },
1082 None => fetches.next().await,
1083 };
1084 let Some(joined) = next else { break };
1085 match joined {
1086 Ok((url, full_budget, Ok(evs))) => {
1087 breaker_record(&url, true, full_budget);
1088 merge_events(evs, &mut result, &mut union_ids);
1089 plan.record(true);
1090 if evidence == Evidence::Quorum && quorum_deadline.is_none() {
1091 quorum_deadline = Some(
1092 tokio::time::Instant::now()
1093 + std::time::Duration::from_millis(QUORUM_GRACE_MS),
1094 );
1095 }
1096 }
1097 Ok((url, full_budget, Err(()))) => {
1098 breaker_record(&url, false, full_budget);
1099 plan.record(false);
1100 }
1101 Err(_) => plan.record(false), }
1103 }
1104
1105 if plan.successes() == 0 {
1106 return Err(format!(
1107 "no relay answered the fetch (0/{} attempted)",
1108 plan.attempted()
1109 ));
1110 }
1111
1112 if !fetches.is_empty() && !quorum_window_closed {
1116 let grace = tokio::time::sleep(std::time::Duration::from_millis(RESIDUAL_GRACE_MS));
1117 tokio::pin!(grace);
1118 loop {
1119 tokio::select! {
1120 _ = &mut grace => break,
1121 next = fetches.next() => match next {
1122 Some(Ok((url, full_budget, Ok(evs)))) => {
1123 breaker_record(&url, true, full_budget);
1124 merge_events(evs, &mut result, &mut union_ids);
1125 }
1126 Some(Ok((url, full_budget, Err(())))) => {
1127 breaker_record(&url, false, full_budget);
1128 }
1129 Some(Err(_)) => continue,
1130 None => break,
1131 }
1132 }
1133 }
1134 }
1135
1136 if !fetches.is_empty() {
1140 let seen: std::collections::HashSet<EventId> = result.iter().map(|e| e.id).collect();
1141 crate::db::spawn_bound(async move {
1144 let mut extra: Vec<Event> = Vec::new();
1145 let mut extra_ids: std::collections::HashSet<EventId> = std::collections::HashSet::new();
1146 while let Some(joined) = fetches.next().await {
1147 if let Ok((url, full_budget, out)) = joined {
1148 breaker_record(&url, out.is_ok(), full_budget);
1152 if let Ok(evs) = out {
1153 for e in evs {
1154 if !seen.contains(&e.id) && extra_ids.insert(e.id) {
1155 extra.push(e);
1156 }
1157 }
1158 }
1159 }
1160 }
1161 submit_stragglers(extra);
1162 });
1163 }
1164
1165 Ok((result, plan.successes(), plan.attempted()))
1166 }
1167}
1168
1169#[async_trait::async_trait]
1170impl Transport for LiveTransport {
1171 async fn publish(&self, event: &Event, relays: &[String]) -> Result<(), String> {
1172 let client = Self::warm_client(relays, self.timeout).await?;
1173 let timeout = self.timeout;
1174 let mut targets: Vec<String> = Vec::new();
1175 for r in relays { if !targets.contains(r) { targets.push(r.clone()); } }
1176 targets = drop_unrevivable_targets(&client, targets).await;
1178 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 {
1190 matches!(
1191 tokio::time::timeout(timeout, client.send_event(&event).to(vec![r.clone()])).await,
1192 Ok(Ok(out)) if RelayUrl::parse(&r).map(|u| out.success.contains_key(&u)).unwrap_or(false)
1193 )
1194 })
1195 })
1196 .collect();
1197 while let Some(joined) = sends.next().await {
1198 if matches!(joined, Ok(true)) {
1199 return Ok(()); }
1201 }
1202 Err("no relay accepted the event".to_string())
1203 }
1204
1205 async fn fetch(&self, query: &Query, relays: &[String]) -> Result<Vec<Event>, String> {
1206 self.fetch_counted(query, relays).await.map(|(events, _successes, _attempted)| events)
1207 }
1208
1209 async fn fetch_plane(&self, plane: &Keys, query: &Query, relays: &[String]) -> Result<Vec<Event>, String> {
1210 if relays.is_empty() {
1211 return Ok(Vec::new());
1212 }
1213 #[cfg(feature = "tor")]
1214 wait_until_tor_ready(
1215 || matches!(crate::tor::transport_state(), crate::tor::TorTransportState::RequiredButInactive),
1216 TOR_READY_WAIT,
1217 ).await?;
1218 let mut targets: Vec<String> = relays.iter().filter(|r| !breaker_tripped(r)).cloned().collect();
1223 if targets.is_empty() {
1224 targets = relays.to_vec();
1225 }
1226 let filter = query.to_filter();
1227 let generation = crate::db::current_session_id();
1228 let key = plane_pool_key(&plane.public_key().to_hex(), &targets);
1229
1230 let (hit, evicted) = plane_pool_take(generation, &key);
1233 disconnect_clients(evicted);
1234 let client = if let Some(c) = hit {
1235 c
1236 } else {
1237 let client = crate::apply_tor_proxy(
1243 nostr_sdk::prelude::Client::builder().authenticator(
1244 nostr_sdk::prelude::SignerAuthenticator::new(plane.clone()),
1245 ),
1246 )
1247 .build();
1248 for r in &targets {
1253 let _ = client.add_managed_relay(r.clone()).capabilities(crate::community_relay_capabilities()).await;
1254 }
1255 client.connect().await;
1256 for r in &targets {
1259 let _ = client
1260 .fetch_events(nostr_sdk::prelude::ReqTarget::single(r.clone(), [filter.clone()]))
1261 .timeout(crate::relay_request_timeout(std::time::Duration::from_secs(5)))
1265 .await;
1266 }
1267 let ev = plane_pool_insert(generation, key, client.clone());
1268 disconnect_clients(ev);
1269 client
1270 };
1271
1272 let targets = drop_unrevivable_targets(&client, targets).await;
1276
1277 let mut result: Vec<Event> = Vec::new();
1278 let mut seen: std::collections::HashSet<EventId> = std::collections::HashSet::new();
1279 let mut successes = 0usize;
1280 for r in &targets {
1281 let res = fetch_relay_eose(&client, r, filter.clone(), self.timeout).await;
1282 breaker_record(r, res.is_ok(), true);
1285 if let Ok(events) = res {
1286 successes += 1;
1287 for e in events {
1288 if seen.insert(e.id) {
1289 result.push(e);
1290 }
1291 }
1292 }
1293 }
1294 if successes == 0 {
1299 return Err(format!("no relay answered the plane fetch (0/{} attempted)", targets.len()));
1300 }
1301 Ok(result)
1303 }
1304
1305 async fn publish_durable(&self, event: &Event, relays: &[String]) -> Result<(), String> {
1306 let client = Self::warm_client(relays, self.timeout).await?;
1311 let timeout = self.timeout;
1312 let event = event.clone();
1313 let backoff = std::time::Duration::from_millis(750);
1314 let mut pending: Vec<String> = Vec::new();
1315 for r in relays { if !pending.contains(r) { pending.push(r.clone()); } }
1316 if pending.is_empty() {
1317 return Err("no relays to broadcast to".to_string());
1318 }
1319
1320 let mut acked_any = false;
1325 let _ = tokio::time::timeout(CONFIRM_WINDOW, async {
1326 loop {
1327 let round = drop_unrevivable_targets(&client, pending.clone()).await;
1331 let sends = round.iter().cloned().map(|r| {
1332 let client = &client;
1333 let event = &event;
1334 Box::pin(async move {
1335 match tokio::time::timeout(timeout, client.send_event(event).to(vec![r.clone()])).await {
1336 Ok(Ok(out)) if RelayUrl::parse(&r).map(|u| out.success.contains_key(&u)).unwrap_or(false) => Ok(r),
1337 _ => Err(()),
1338 }
1339 })
1340 });
1341 if let Ok((winner, _losers)) = futures_util::future::select_ok(sends).await {
1342 acked_any = true;
1343 pending.retain(|r| r != &winner);
1344 break;
1345 }
1346 tokio::time::sleep(backoff).await;
1347 }
1348 })
1349 .await;
1350
1351 if !acked_any {
1352 return Err(format!("no relay accepted the event within {}s", CONFIRM_WINDOW.as_secs()));
1353 }
1354 if pending.is_empty() {
1355 return Ok(()); }
1357
1358 tokio::spawn(async move {
1364 let client_ref = &client;
1365 let event_ref = &event;
1366 let _ = durable_broadcast(&pending, MAX_PUBLISH_ATTEMPTS, backoff, move |round| {
1367 Box::pin(async move {
1368 match tokio::time::timeout(timeout, client_ref.send_event(event_ref).to(round.clone())).await {
1369 Ok(Ok(output)) => round.into_iter().filter(|p| RelayUrl::parse(p).map(|u| output.success.contains_key(&u)).unwrap_or(false)).collect(),
1370 _ => Vec::new(),
1371 }
1372 })
1373 })
1374 .await;
1375 });
1376 Ok(())
1377 }
1378}
1379
1380#[cfg(test)]
1381pub(crate) mod memory {
1382 use super::*;
1383 use std::collections::{HashMap, HashSet};
1384 use std::sync::Mutex;
1385
1386 pub struct MemoryRelay {
1391 per_relay: Mutex<HashMap<String, Vec<Event>>>,
1392 subscribers: Mutex<Vec<(Query, tokio::sync::mpsc::UnboundedSender<Event>)>>,
1393 }
1394
1395 fn is_ephemeral(kind: u16) -> bool {
1398 (20000..30000).contains(&kind)
1399 }
1400
1401 impl MemoryRelay {
1402 pub fn new() -> Self {
1403 MemoryRelay {
1404 per_relay: Mutex::new(HashMap::new()),
1405 subscribers: Mutex::new(Vec::new()),
1406 }
1407 }
1408
1409 pub fn stored_count(&self) -> usize {
1413 self.per_relay.lock().unwrap().values().map(|v| v.len()).sum()
1414 }
1415
1416 pub fn subscribe(&self, query: Query) -> tokio::sync::mpsc::UnboundedReceiver<Event> {
1419 let (tx, rx) = tokio::sync::mpsc::unbounded_channel();
1420 self.subscribers.lock().unwrap().push((query, tx));
1421 rx
1422 }
1423
1424 fn deliver(&self, event: &Event) {
1426 self.subscribers.lock().unwrap().retain(|(q, tx)| {
1427 if q.matches(event) {
1428 tx.send(event.clone()).is_ok()
1429 } else {
1430 !tx.is_closed()
1431 }
1432 });
1433 }
1434
1435 pub fn inject(&self, event: &Event, relays: &[String]) {
1438 self.deliver(event);
1439 if is_ephemeral(event.kind.as_u16()) {
1440 return; }
1442 let d_tag = |e: &Event| e.tags.iter().find_map(|t| {
1448 let s = t.as_slice();
1449 (s.len() >= 2 && s[0] == "d").then(|| s[1].clone())
1450 }).unwrap_or_default();
1451 let k = event.kind.as_u16();
1452 let replaceable = (30000..40000).contains(&k) || (10000..20000).contains(&k) || k == 0 || k == 3;
1453 let coord = (event.kind.as_u16(), event.pubkey, d_tag(event));
1454 let mut map = self.per_relay.lock().unwrap();
1455 for r in relays {
1456 let v = map.entry(r.clone()).or_default();
1457 if replaceable {
1458 v.retain(|e| (e.kind.as_u16(), e.pubkey, d_tag(e)) != coord);
1459 }
1460 v.push(event.clone());
1461 }
1462 }
1463
1464 pub fn count_on(&self, relay: &str) -> usize {
1466 self.per_relay.lock().unwrap().get(relay).map_or(0, |v| v.len())
1467 }
1468
1469 fn apply_deletion(&self, deletion: &Event, relays: &[String]) {
1474 let mut id_targets: HashSet<String> = HashSet::new();
1475 let mut coord_targets: HashSet<String> = HashSet::new();
1476 for t in deletion.tags.iter() {
1477 let s = t.as_slice();
1478 if s.len() >= 2 && s[0] == "e" {
1479 id_targets.insert(s[1].clone());
1480 } else if s.len() >= 2 && s[0] == "a" {
1481 coord_targets.insert(s[1].clone());
1482 }
1483 }
1484 let mut map = self.per_relay.lock().unwrap();
1485 for r in relays {
1486 if let Some(events) = map.get_mut(r) {
1487 events.retain(|e| {
1488 if e.pubkey != deletion.pubkey {
1489 return true;
1490 }
1491 if id_targets.contains(&e.id.to_hex()) {
1492 return false;
1493 }
1494 let d = e.tags.iter().find_map(|t| {
1496 let s = t.as_slice();
1497 (s.len() >= 2 && s[0] == "d").then(|| s[1].clone())
1498 }).unwrap_or_default();
1499 let coord = format!("{}:{}:{}", e.kind.as_u16(), e.pubkey.to_hex(), d);
1500 !coord_targets.contains(&coord)
1501 });
1502 }
1503 }
1504 }
1505 }
1506
1507 #[async_trait::async_trait]
1508 impl Transport for MemoryRelay {
1509 async fn publish(&self, event: &Event, relays: &[String]) -> Result<(), String> {
1510 if event.kind == Kind::EventDeletion {
1512 self.apply_deletion(event, relays);
1513 self.deliver(event);
1514 } else {
1515 self.inject(event, relays);
1516 }
1517 Ok(())
1518 }
1519
1520 async fn publish_durable(&self, event: &Event, relays: &[String]) -> Result<(), String> {
1523 self.publish(event, relays).await
1524 }
1525
1526 async fn fetch(&self, query: &Query, relays: &[String]) -> Result<Vec<Event>, String> {
1527 let map = self.per_relay.lock().unwrap();
1528 let mut seen = HashSet::new();
1529 let mut out = Vec::new();
1530 for r in relays {
1531 if let Some(events) = map.get(r) {
1532 for ev in events {
1533 if is_ephemeral(ev.kind.as_u16()) {
1536 continue;
1537 }
1538 if query.matches(ev) && seen.insert(ev.id) {
1539 out.push(ev.clone());
1540 }
1541 }
1542 }
1543 }
1544 if let Some(limit) = query.limit {
1547 out.sort_by(|a, b| b.created_at.cmp(&a.created_at).then_with(|| b.id.cmp(&a.id)));
1548 out.truncate(limit);
1549 }
1550 Ok(out)
1551 }
1552
1553 async fn fetch_plane(&self, _plane: &Keys, query: &Query, relays: &[String]) -> Result<Vec<Event>, String> {
1554 self.fetch(query, relays).await
1556 }
1557 }
1558}
1559
1560#[cfg(test)]
1561mod tests {
1562 use super::*;
1563
1564 #[test]
1567 fn a_dead_socket_is_skipped_but_a_connecting_or_unknown_one_is_not() {
1568 let targets: Vec<String> = ["wss://dead", "wss://connecting", "wss://unknown"].map(String::from).into();
1569 let out = drop_unrevivable(targets, |r| r == "wss://dead");
1570 assert_eq!(out, ["wss://connecting", "wss://unknown"].map(String::from).to_vec());
1571 }
1572
1573 #[test]
1574 fn an_all_dead_set_falls_back_to_the_full_list() {
1575 let targets: Vec<String> = ["wss://a", "wss://b"].map(String::from).into();
1578 let out = drop_unrevivable(targets.clone(), |_| true);
1579 assert_eq!(out, targets);
1580 }
1581
1582 #[tokio::test]
1587 async fn tor_gate_passes_immediately_when_not_blocked() {
1588 let start = std::time::Instant::now();
1589 let res = wait_until_tor_ready(|| false, std::time::Duration::from_secs(30)).await;
1590 assert!(res.is_ok());
1591 assert!(start.elapsed() < std::time::Duration::from_secs(1), "must not wait when Tor is ready");
1592 }
1593
1594 #[tokio::test]
1595 async fn tor_gate_errors_honestly_after_timeout_when_perpetually_blocked() {
1596 let res = wait_until_tor_ready(|| true, std::time::Duration::from_millis(300)).await;
1599 let err = res.expect_err("should error when Tor never activates");
1600 assert!(err.to_lowercase().contains("tor"), "error must name Tor, got: {err}");
1601 }
1602
1603 #[tokio::test]
1604 async fn tor_gate_passes_once_circuit_comes_up_mid_wait() {
1605 let calls = std::sync::atomic::AtomicUsize::new(0);
1606 let res = wait_until_tor_ready(
1608 || calls.fetch_add(1, std::sync::atomic::Ordering::SeqCst) < 3,
1609 std::time::Duration::from_secs(5),
1610 ).await;
1611 assert!(res.is_ok(), "should succeed once Tor activates within the window");
1612 }
1613
1614 fn evt(kind: u16, z: &str) -> Event {
1615 EventBuilder::new(Kind::Custom(kind), "x")
1616 .tags([Tag::custom(
1617 "z",
1618 [z.to_string()],
1619 )])
1620 .finalize(&Keys::generate())
1621 .unwrap()
1622 }
1623
1624 #[test]
1625 fn query_matches_kind_and_z() {
1626 let e = evt(3300, "abc");
1627 assert!(Query { kinds: vec![3300], z_tags: vec!["abc".into()], since: None, ..Default::default() }.matches(&e));
1628 assert!(!Query { kinds: vec![3301], ..Default::default() }.matches(&e));
1629 assert!(!Query { kinds: vec![], z_tags: vec!["xyz".into()], since: None, ..Default::default() }.matches(&e));
1630 assert!(Query::default().matches(&e), "empty query matches anything");
1631 }
1632
1633 fn evt_at(kind: u16, secs: u64) -> Event {
1634 EventBuilder::new(Kind::Custom(kind), "x")
1635 .custom_created_at(Timestamp::from(secs))
1636 .finalize(&Keys::generate())
1637 .unwrap()
1638 }
1639
1640 fn evt_z_at(kind: u16, z: &str, secs: u64) -> Event {
1643 EventBuilder::new(Kind::Custom(kind), "x")
1644 .custom_created_at(Timestamp::from(secs))
1645 .tags([Tag::custom(
1646 "z",
1647 [z.to_string()],
1648 )])
1649 .finalize(&Keys::generate())
1650 .unwrap()
1651 }
1652
1653 #[tokio::test]
1654 async fn fetch_pages_with_until_and_limit_newest_first() {
1655 let relay = super::memory::MemoryRelay::new();
1657 let relays = vec!["r1".to_string()];
1658 for s in 1..=5u64 {
1659 relay.inject(&evt_z_at(3300, "pg", s), &relays);
1660 }
1661 let secs = |evs: &[Event]| evs.iter().map(|e| e.created_at.as_secs()).collect::<Vec<_>>();
1662
1663 let latest = relay
1665 .fetch(&Query { kinds: vec![3300], z_tags: vec!["pg".into()], limit: Some(2), ..Default::default() }, &relays)
1666 .await
1667 .unwrap();
1668 assert_eq!(secs(&latest), vec![5, 4]);
1669
1670 let older = relay
1672 .fetch(&Query { kinds: vec![3300], z_tags: vec!["pg".into()], until: Some(3), limit: Some(2), ..Default::default() }, &relays)
1673 .await
1674 .unwrap();
1675 assert_eq!(secs(&older), vec![3, 2]);
1676
1677 let start = relay
1680 .fetch(&Query { kinds: vec![3300], z_tags: vec!["pg".into()], until: Some(1), limit: Some(2), ..Default::default() }, &relays)
1681 .await
1682 .unwrap();
1683 assert_eq!(secs(&start), vec![1]);
1684 }
1685
1686 #[test]
1687 fn to_filter_translates_kinds_z_and_since() {
1688 let q = Query { kinds: vec![3300], z_tags: vec!["abc".into()], since: Some(100), ..Default::default() };
1691 let filter = q.to_filter();
1692 let matching = EventBuilder::new(Kind::Custom(3300), "x")
1693 .custom_created_at(Timestamp::from(150))
1694 .tags([Tag::custom(
1695 "z",
1696 ["abc".to_string()],
1697 )])
1698 .finalize(&Keys::generate())
1699 .unwrap();
1700 assert!(filter.match_event(&matching, MatchEventOptions::new()), "to_filter must accept what matches() accepts");
1701 assert!(q.matches(&matching));
1702
1703 let wrong_kind = EventBuilder::new(Kind::Custom(3301), "x")
1705 .custom_created_at(Timestamp::from(150))
1706 .tags([Tag::custom(
1707 "z",
1708 ["abc".to_string()],
1709 )])
1710 .finalize(&Keys::generate())
1711 .unwrap();
1712 assert!(!filter.match_event(&wrong_kind, MatchEventOptions::new()));
1713 }
1714
1715 fn evt_sl(kind: u16, letter: SingleLetterTag, value: &str) -> Event {
1717 EventBuilder::new(Kind::Custom(kind), "x")
1718 .tags([Tag::custom(
1719 letter.as_char().to_string(),
1720 [value.to_string()],
1721 )])
1722 .finalize(&Keys::generate())
1723 .unwrap()
1724 }
1725
1726 #[test]
1727 fn to_filter_and_matches_agree_on_authors() {
1728 let keys = Keys::generate();
1729 let e = EventBuilder::new(Kind::Custom(1059), "x").finalize(&keys).unwrap();
1730 let q = Query { kinds: vec![1059], authors: vec![keys.public_key().to_hex()], ..Default::default() };
1731 assert!(q.matches(&e));
1732 assert!(q.to_filter().match_event(&e, MatchEventOptions::new()));
1733 let miss = Query { kinds: vec![1059], authors: vec![Keys::generate().public_key().to_hex()], ..Default::default() };
1734 assert!(!miss.matches(&e));
1735 assert!(!miss.to_filter().match_event(&e, MatchEventOptions::new()));
1736 }
1737
1738 #[test]
1739 fn to_filter_and_matches_agree_on_p_tags() {
1740 let recipient = Keys::generate().public_key().to_hex();
1741 let e = evt_sl(1059, SingleLetterTag::LOWERCASE_P, &recipient);
1742 let q = Query { kinds: vec![1059], p_tags: vec![recipient], ..Default::default() };
1743 assert!(q.matches(&e));
1744 assert!(q.to_filter().match_event(&e, MatchEventOptions::new()));
1745 let miss = Query { kinds: vec![1059], p_tags: vec![Keys::generate().public_key().to_hex()], ..Default::default() };
1746 assert!(!miss.matches(&e));
1747 assert!(!miss.to_filter().match_event(&e, MatchEventOptions::new()));
1748 }
1749
1750 #[test]
1751 fn to_filter_and_matches_agree_on_k_tags() {
1752 let e = evt_sl(1059, SingleLetterTag::LOWERCASE_K, "3311");
1753 let q = Query { kinds: vec![1059], k_tags: vec!["3311".into()], ..Default::default() };
1754 assert!(q.matches(&e));
1755 assert!(q.to_filter().match_event(&e, MatchEventOptions::new()));
1756 let miss = Query { kinds: vec![1059], k_tags: vec!["3300".into()], ..Default::default() };
1757 assert!(!miss.matches(&e));
1758 assert!(!miss.to_filter().match_event(&e, MatchEventOptions::new()));
1759 }
1760
1761 #[test]
1762 fn to_filter_empty_kinds_only_constrains_z_and_since() {
1763 let q = Query { kinds: vec![], z_tags: vec!["p".into()], since: None, ..Default::default() };
1765 let filter = q.to_filter();
1766 let e = evt(3300, "p");
1767 assert!(filter.match_event(&e, MatchEventOptions::new()));
1768 assert!(q.matches(&e));
1769 }
1770
1771 #[test]
1772 fn since_is_an_inclusive_lower_bound() {
1773 let below = evt_at(3300, 99);
1774 let exact = evt_at(3300, 100);
1775 let above = evt_at(3300, 101);
1776 let q = Query { kinds: vec![3300], z_tags: vec![], since: Some(100), ..Default::default() };
1777 assert!(!q.matches(&below), "below the floor is excluded");
1778 assert!(q.matches(&exact), "exactly the floor is included");
1779 assert!(q.matches(&above), "above the floor is included");
1780 }
1781
1782 #[test]
1783 fn z_tags_match_as_or_set() {
1784 let e = evt(3300, "p2");
1787 let q = Query { kinds: vec![3300], z_tags: vec!["p1".into(), "p2".into()], since: None, ..Default::default() };
1788 assert!(q.matches(&e));
1789 let miss = Query { kinds: vec![3300], z_tags: vec!["p1".into(), "p3".into()], since: None, ..Default::default() };
1790 assert!(!miss.matches(&e));
1791 }
1792
1793 #[tokio::test]
1794 async fn fetch_unions_and_dedups_across_relays() {
1795 use super::memory::MemoryRelay;
1796 let relay = MemoryRelay::new();
1797 let relays = vec!["r1".to_string(), "r2".to_string(), "r3".to_string()];
1798 let e = evt(3300, "p");
1799 relay.publish(&e, &relays).await.unwrap();
1800 let got = relay
1801 .fetch(&Query { kinds: vec![3300], z_tags: vec!["p".into()], since: None, ..Default::default() }, &relays)
1802 .await
1803 .unwrap();
1804 assert_eq!(got.len(), 1, "same event on 3 relays dedups to 1");
1805 }
1806
1807 #[tokio::test]
1808 async fn durable_broadcast_retries_only_the_failing_relays_until_they_ack() {
1809 use std::cell::Cell;
1812 let relays = vec!["r1".to_string(), "r2".to_string(), "r3".to_string()];
1813 let round = Cell::new(0usize);
1814 let r2_round_seen = Cell::new(0usize);
1815 let res = durable_broadcast(&relays, 30, std::time::Duration::ZERO, |pending| {
1816 let n = round.get();
1817 round.set(n + 1);
1818 if n >= 1 {
1820 assert_eq!(pending, vec!["r2".to_string()], "only the failing relay is retried");
1821 r2_round_seen.set(r2_round_seen.get() + 1);
1822 }
1823 Box::pin(async move {
1824 pending.into_iter().filter(|r| r != "r2" || n >= 4).collect()
1825 })
1826 })
1827 .await;
1828 assert!(res.is_ok(), "all relays eventually ACK → Ok");
1829 assert!(round.get() >= 5, "kept retrying r2 across rounds");
1830 }
1831
1832 #[tokio::test]
1833 async fn durable_broadcast_is_ok_if_some_ack_even_when_one_never_does() {
1834 let relays = vec!["r1".to_string(), "r2".to_string()];
1837 let res = durable_broadcast(&relays, 5, std::time::Duration::ZERO, |pending| {
1838 Box::pin(async move { pending.into_iter().filter(|r| r == "r1").collect() })
1839 })
1840 .await;
1841 assert!(res.is_ok(), "≥1 relay accepted → Ok despite a permanently-failing relay");
1842 }
1843
1844 #[tokio::test]
1845 async fn durable_broadcast_errs_only_if_zero_relays_ever_accept() {
1846 let relays = vec!["r1".to_string(), "r2".to_string()];
1847 let res = durable_broadcast(&relays, 5, std::time::Duration::ZERO, |_pending| {
1848 Box::pin(async move { Vec::new() }) })
1850 .await;
1851 assert!(res.is_err(), "zero acceptances after the retry cap → Err");
1852 }
1853
1854 #[tokio::test]
1855 async fn redundancy_self_heals_a_missing_relay() {
1856 use super::memory::MemoryRelay;
1857 let relay = MemoryRelay::new();
1858 let all = vec!["r1".to_string(), "r2".to_string(), "r3".to_string()];
1859 let e = evt(3300, "p");
1860 relay.inject(&e, &["r2".to_string()]);
1862 assert_eq!(relay.count_on("r1"), 0);
1863 assert_eq!(relay.count_on("r2"), 1);
1864 let got = relay.fetch(&Query { kinds: vec![3300], ..Default::default() }, &all).await.unwrap();
1866 assert_eq!(got.len(), 1);
1867 }
1868
1869 #[tokio::test]
1870 async fn ephemeral_kind_streams_live_but_is_never_stored_or_fetched() {
1871 use super::memory::MemoryRelay;
1872 let relay = MemoryRelay::new();
1873 let relays = vec!["r1".to_string()];
1874 let mut sub = relay.subscribe(Query { kinds: vec![21059], ..Default::default() });
1875 let e = evt(21059, "p");
1876 relay.publish(&e, &relays).await.unwrap();
1877 assert_eq!(relay.count_on("r1"), 0, "ephemeral is never stored");
1878 let got = relay
1879 .fetch(&Query { kinds: vec![21059], ..Default::default() }, &relays)
1880 .await
1881 .unwrap();
1882 assert!(got.is_empty(), "a real relay never serves an ephemeral from a fetch");
1883 assert_eq!(sub.try_recv().unwrap().id, e.id, "but a live subscriber receives it");
1884 }
1885
1886 #[tokio::test]
1887 async fn stored_kind_is_fetchable_and_delivered_live() {
1888 use super::memory::MemoryRelay;
1889 let relay = MemoryRelay::new();
1890 let relays = vec!["r1".to_string()];
1891 let mut sub = relay.subscribe(Query { kinds: vec![1059], ..Default::default() });
1892 let e = evt(1059, "p");
1893 relay.publish(&e, &relays).await.unwrap();
1894 let got = relay
1895 .fetch(&Query { kinds: vec![1059], ..Default::default() }, &relays)
1896 .await
1897 .unwrap();
1898 assert_eq!(got.len(), 1, "stored kind is fetchable");
1899 assert_eq!(sub.try_recv().unwrap().id, e.id, "and delivered to the live subscriber");
1900 }
1901
1902 #[tokio::test]
1903 async fn p_tags_route_a_giftwrap_to_the_matching_subscriber() {
1904 use super::memory::MemoryRelay;
1905 let relay = MemoryRelay::new();
1906 let relays = vec!["r1".to_string()];
1907 let alice = Keys::generate().public_key().to_hex();
1908 let bob = Keys::generate().public_key().to_hex();
1909 let mut sub_alice =
1910 relay.subscribe(Query { kinds: vec![1059], p_tags: vec![alice.clone()], ..Default::default() });
1911 let mut sub_bob =
1912 relay.subscribe(Query { kinds: vec![1059], p_tags: vec![bob.clone()], ..Default::default() });
1913 let wrap = evt_sl(1059, SingleLetterTag::LOWERCASE_P, &alice);
1914 relay.publish(&wrap, &relays).await.unwrap();
1915 assert_eq!(sub_alice.try_recv().unwrap().id, wrap.id, "addressed recipient gets it live");
1916 assert!(sub_bob.try_recv().is_err(), "a differently-addressed subscriber does not");
1917 let for_alice = relay
1919 .fetch(&Query { kinds: vec![1059], p_tags: vec![alice], ..Default::default() }, &relays)
1920 .await
1921 .unwrap();
1922 assert_eq!(for_alice.len(), 1);
1923 let for_bob = relay
1924 .fetch(&Query { kinds: vec![1059], p_tags: vec![bob], ..Default::default() }, &relays)
1925 .await
1926 .unwrap();
1927 assert!(for_bob.is_empty());
1928 }
1929
1930 #[test]
1933 fn union_plan_fast_satisfied_on_first_success() {
1934 let mut p = UnionPlan::new(Evidence::Fast, 4);
1935 p.record(false);
1936 assert!(!p.satisfied(), "a failure is not evidence");
1937 p.record(true);
1938 assert!(p.satisfied(), "one genuine EOSE satisfies Fast");
1939 assert!(!p.exhausted());
1940 }
1941
1942 #[test]
1943 fn union_plan_quorum_majority_math() {
1944 for (n, need) in [(2usize, 2usize), (3, 2), (4, 3), (5, 3)] {
1946 let mut p = UnionPlan::new(Evidence::Quorum, n);
1947 for _ in 0..need - 1 {
1948 p.record(true);
1949 }
1950 assert!(!p.satisfied(), "{}/{} must not satisfy quorum", need - 1, n);
1951 p.record(true);
1952 assert!(p.satisfied(), "{}/{} satisfies quorum", need, n);
1953 }
1954 }
1955
1956 #[test]
1957 fn union_plan_quorum_failures_never_substitute_for_successes() {
1958 let mut p = UnionPlan::new(Evidence::Quorum, 3);
1959 p.record(true);
1960 p.record(false);
1961 p.record(false);
1962 assert!(!p.satisfied(), "1 success + 2 failures is not a majority");
1963 assert!(p.exhausted(), "all resolved — the degraded path returns best-effort");
1964 assert_eq!(p.successes(), 1);
1965 }
1966
1967 #[test]
1968 fn union_plan_full_requires_every_relay_resolved() {
1969 let mut p = UnionPlan::new(Evidence::Full, 3);
1970 p.record(true);
1971 p.record(true);
1972 assert!(!p.satisfied(), "Full waits for the last relay even after 2 EOSEs");
1973 p.record(false);
1974 assert!(p.satisfied(), "a timeout is a resolution — Full is done");
1975 assert!(p.exhausted());
1976 }
1977
1978 #[test]
1979 fn union_plan_all_dead_is_reportable_not_a_confident_empty() {
1980 let mut p = UnionPlan::new(Evidence::Quorum, 2);
1981 p.record(false);
1982 p.record(false);
1983 assert!(p.exhausted());
1984 assert_eq!(p.successes(), 0, "the caller must map this to Err, never Ok(vec![])");
1985 }
1986
1987 const BREAKER_TEST_GEN: u64 = u64::MAX;
1993
1994 #[test]
1995 fn breaker_trips_only_after_consecutive_full_budget_failures() {
1996 let url = "wss://breaker-test-full-budget.example";
1997 breaker_record_at(BREAKER_TEST_GEN, url, false, true);
1998 assert!(!breaker_tripped_at(BREAKER_TEST_GEN, url), "one failure is below the threshold");
1999 breaker_record_at(BREAKER_TEST_GEN, url, false, true);
2000 assert!(breaker_tripped_at(BREAKER_TEST_GEN, url), "two consecutive full-budget failures trip");
2001 }
2002
2003 #[test]
2004 fn breaker_demoted_budget_failures_never_count() {
2005 let url = "wss://breaker-test-demoted.example";
2006 breaker_record_at(BREAKER_TEST_GEN, url, false, false);
2007 breaker_record_at(BREAKER_TEST_GEN, url, false, false);
2008 breaker_record_at(BREAKER_TEST_GEN, url, false, false);
2009 assert!(
2010 !breaker_tripped_at(BREAKER_TEST_GEN, url),
2011 "demoted-budget failures must not trip (anti-starvation: the post-cooldown probe must stay reachable)"
2012 );
2013 }
2014
2015 #[test]
2016 fn breaker_success_resets_the_entry() {
2017 let url = "wss://breaker-test-reset.example";
2018 breaker_record_at(BREAKER_TEST_GEN, url, false, true);
2019 breaker_record_at(BREAKER_TEST_GEN, url, false, true);
2020 assert!(breaker_tripped_at(BREAKER_TEST_GEN, url));
2021 breaker_record_at(BREAKER_TEST_GEN, url, true, false);
2023 assert!(!breaker_tripped_at(BREAKER_TEST_GEN, url), "any success unconditionally resets");
2024 breaker_record_at(BREAKER_TEST_GEN, url, false, true);
2025 assert!(!breaker_tripped_at(BREAKER_TEST_GEN, url), "and the failure count restarted from zero");
2026 }
2027
2028 #[test]
2031 fn declared_evidence_stands_and_default_is_quorum() {
2032 assert_eq!(Query::default().evidence, Evidence::Quorum, "unclassified sites get Quorum");
2033 assert_eq!(
2034 effective_evidence(&Query { until: Some(1), evidence: Evidence::Fast, ..Default::default() }),
2035 Evidence::Fast,
2036 "chat pagination rides its declared tier — absence verdicts request Full themselves"
2037 );
2038 assert_eq!(
2039 effective_evidence(&Query { evidence: Evidence::Fast, ..Default::default() }),
2040 Evidence::Fast,
2041 "without `until` the declared tier stands"
2042 );
2043 }
2044}