1use std::collections::{HashMap, HashSet};
7use std::sync::atomic::{AtomicU64, AtomicUsize, Ordering};
8use std::sync::{Arc, Mutex, Weak};
9use std::time::Instant;
10
11use nostr_sdk::prelude::*;
12use std::sync::LazyLock;
13
14use crate::state::nostr_client;
15use crate::ClientRelayExt;
16
17pub struct EventPublishTracker {
47 event_id: EventId,
48 successes: Mutex<Vec<RelayUrl>>,
51 notify: tokio::sync::Notify,
52 in_flight: AtomicUsize,
56}
57
58impl EventPublishTracker {
59 fn new(event_id: EventId, initial_in_flight: usize) -> Arc<Self> {
60 Arc::new(Self {
61 event_id,
62 successes: Mutex::new(Vec::new()),
63 notify: tokio::sync::Notify::new(),
64 in_flight: AtomicUsize::new(initial_in_flight),
65 })
66 }
67
68 fn note_success(&self, url: RelayUrl) {
70 self.successes.lock().unwrap().push(url);
71 self.notify.notify_waiters();
72 }
73
74 fn note_settled(&self) {
78 let mut trackers = PUBLISH_TRACKERS.lock().unwrap();
82 if self.in_flight.fetch_sub(1, Ordering::SeqCst) == 1 {
83 self.notify.notify_waiters();
84 match trackers.get(&self.event_id) {
85 Some(current) if std::ptr::eq(Arc::as_ptr(current), self) => {
86 trackers.remove(&self.event_id);
87 }
88 _ => {}
89 }
90 }
91 }
92
93 pub async fn next_success(&self, cursor: &mut usize) -> Option<RelayUrl> {
99 loop {
100 let notified = self.notify.notified();
104 tokio::pin!(notified);
105 notified.as_mut().enable();
106
107 let (next, done) = {
108 let successes = self.successes.lock().unwrap();
109 let next = successes.get(*cursor).cloned();
110 let done = self.in_flight.load(Ordering::SeqCst) == 0
111 && *cursor >= successes.len();
112 (next, done)
113 };
114
115 if let Some(url) = next {
116 *cursor += 1;
117 return Some(url);
118 }
119 if done {
120 return None;
121 }
122
123 notified.await;
124 }
125 }
126}
127
128static PUBLISH_TRACKERS: LazyLock<Mutex<HashMap<EventId, Arc<EventPublishTracker>>>> =
131 LazyLock::new(|| Mutex::new(HashMap::new()));
132
133pub fn get_publish_tracker(event_id: &EventId) -> Option<Arc<EventPublishTracker>> {
139 PUBLISH_TRACKERS.lock().unwrap().get(event_id).cloned()
140}
141
142pub fn spawn_tracked_publish(
153 resolved: Vec<(RelayUrl, Relay)>,
154 event: Event,
155) -> Vec<tokio::task::JoinHandle<(RelayUrl, Result<EventId, String>)>> {
156 let event_id = event.id;
157 if resolved.is_empty() {
160 return Vec::new();
161 }
162 let tracker = {
166 let mut trackers = PUBLISH_TRACKERS.lock().unwrap();
167 match trackers.get(&event_id) {
168 Some(existing) if existing.in_flight.load(Ordering::SeqCst) > 0 => {
169 existing.in_flight.fetch_add(resolved.len(), Ordering::SeqCst);
170 existing.clone()
171 }
172 _ => {
173 let t = EventPublishTracker::new(event_id, resolved.len());
174 trackers.insert(event_id, t.clone());
175 t
176 }
177 }
178 };
179
180 let mut handles = Vec::with_capacity(resolved.len());
181 for (url, relay) in resolved {
182 let event = event.clone();
183 let tracker = tracker.clone();
184 handles.push(tokio::spawn(async move {
185 let result = relay
186 .send_event(&event)
187 .await
188 .map(|o| *o.id())
189 .map_err(|e| e.to_string());
190 if result.is_ok() {
191 tracker.note_success(url.clone());
192 }
193 tracker.note_settled();
194 (url, result)
195 }));
196 }
197 handles
198}
199
200const CACHE_TTL_SECS: u64 = 3600; const CACHE_TTL_ERROR_SECS: u64 = 60; struct CachedRelays {
211 relays: Vec<String>,
212 fetched_at: Instant,
213 fetch_ok: bool,
216}
217
218static INBOX_RELAY_CACHE: LazyLock<Mutex<HashMap<PublicKey, CachedRelays>>> =
219 LazyLock::new(|| Mutex::new(HashMap::new()));
220
221pub fn clear_inbox_relay_cache() {
227 if let Ok(mut cache) = INBOX_RELAY_CACHE.lock() {
228 cache.clear();
229 }
230}
231
232static FETCH_LOCKS: LazyLock<Mutex<HashMap<PublicKey, Weak<tokio::sync::Mutex<()>>>>> =
239 LazyLock::new(|| Mutex::new(HashMap::new()));
240
241static PRUNE_COUNTER: AtomicU64 = AtomicU64::new(0);
245
246#[cfg(not(test))]
249const PRUNE_INTERVAL: u64 = 100;
250
251#[cfg(test)]
254const PRUNE_INTERVAL: u64 = 1;
255
256struct FetchLockEntryCleanup {
259 pubkey: PublicKey,
260 key_lock: Arc<tokio::sync::Mutex<()>>,
261}
262
263impl FetchLockEntryCleanup {
264 fn new(pubkey: PublicKey, key_lock: Arc<tokio::sync::Mutex<()>>) -> Self {
265 Self { pubkey, key_lock }
266 }
267}
268
269impl Drop for FetchLockEntryCleanup {
270 fn drop(&mut self) {
271 let mut locks = match FETCH_LOCKS.lock() {
272 Ok(locks) => locks,
273 Err(_) => return, };
275
276 let should_remove = match locks.get(&self.pubkey).and_then(|weak| weak.upgrade()) {
277 Some(current) => {
278 Arc::ptr_eq(¤t, &self.key_lock) && Arc::strong_count(¤t) == 2
282 }
283 None => false,
284 };
285 if should_remove {
286 locks.remove(&self.pubkey);
287 }
288 }
289}
290
291pub fn normalize_relay_url(s: &str) -> String {
299 s.trim_end_matches('/').to_ascii_lowercase()
300}
301
302async fn inbox_query_targets(client: &Client) -> Vec<RelayUrl> {
307 let discovery: HashSet<String> = crate::state::discovery_relay_iter()
308 .map(normalize_relay_url)
309 .collect();
310 client
311 .relays().all()
312 .await
313 .iter()
314 .filter(|(url, relay)| {
315 relay.capabilities().load().can_read() || discovery.contains(&normalize_relay_url(url.as_str()))
316 })
317 .map(|(url, _)| url.clone())
318 .collect()
319}
320
321struct FetchResult {
323 relays: Vec<String>,
324 fetch_ok: bool,
326}
327
328async fn fetch_inbox_relays(client: &Client, pubkey: &PublicKey) -> FetchResult {
332 let filter = Filter::new()
333 .author(*pubkey)
334 .kind(Kind::Custom(10050))
335 .limit(1);
336
337 let targets = inbox_query_targets(client).await;
338 let fetched = if targets.is_empty() {
339 client
340 .fetch_events(filter).timeout(std::time::Duration::from_secs(5))
341 .await
342 } else {
343 client
344 .fetch_events(nostr_sdk::prelude::ReqTarget::manual(
345 targets.into_iter().map(|u| (u, vec![filter.clone()])),
346 ))
347 .timeout(std::time::Duration::from_secs(5))
348 .await
349 };
350 let events = match fetched {
351 Ok(events) => events,
352 Err(e) => {
353 eprintln!("[InboxRelays] Failed to fetch 10050 for {}: {}", pubkey, e);
354 return FetchResult { relays: Vec::new(), fetch_ok: false };
355 }
356 };
357
358 let event = match events.into_iter().max_by_key(|e| e.created_at) {
361 Some(e) => e,
362 None => return FetchResult { relays: Vec::new(), fetch_ok: true },
363 };
364
365 FetchResult { relays: parse_relay_tags(&event.tags), fetch_ok: true }
366}
367
368fn parse_relay_tags(tags: &Tags) -> Vec<String> {
371 tags.iter()
372 .filter_map(|tag| {
373 let values: Vec<&str> = tag.as_slice().iter().map(|s| s.as_str()).collect();
374 if values.len() >= 2 && values[0] == "relay" {
375 Some(values[1].to_string())
376 } else {
377 None
378 }
379 })
380 .collect()
381}
382
383async fn get_or_fetch_with_lock<F, Fut>(pubkey: &PublicKey, fetch_fn: F) -> Vec<String>
388where
389 F: FnOnce() -> Fut,
390 Fut: std::future::Future<Output = FetchResult>,
391{
392 {
394 let cache = INBOX_RELAY_CACHE.lock().unwrap();
395 if let Some(entry) = cache.get(pubkey) {
396 let ttl = if entry.fetch_ok { CACHE_TTL_SECS } else { CACHE_TTL_ERROR_SECS };
397 if entry.fetched_at.elapsed().as_secs() < ttl {
398 return entry.relays.clone();
399 }
400 }
401 }
402
403 let cleanup_guard = {
406 let mut locks = FETCH_LOCKS.lock().unwrap();
407
408 if PRUNE_COUNTER.fetch_add(1, Ordering::Relaxed) % PRUNE_INTERVAL == 0 {
412 locks.retain(|_, weak| Weak::strong_count(weak) > 0);
413 }
414
415 let weak = locks.entry(*pubkey).or_insert_with(|| Weak::new());
416 let key_lock = match weak.upgrade() {
419 Some(arc) => arc,
420 None => {
421 let new_arc = Arc::new(tokio::sync::Mutex::new(()));
422 *weak = Arc::downgrade(&new_arc);
423 new_arc
424 }
425 };
426 FetchLockEntryCleanup::new(*pubkey, key_lock)
428 };
429 let relays = {
430 let _guard = cleanup_guard.key_lock.lock().await;
431
432 let cached_relays = {
434 let cache = INBOX_RELAY_CACHE.lock().unwrap();
435 if let Some(entry) = cache.get(pubkey) {
436 let ttl = if entry.fetch_ok { CACHE_TTL_SECS } else { CACHE_TTL_ERROR_SECS };
437 if entry.fetched_at.elapsed().as_secs() < ttl {
438 Some(entry.relays.clone())
439 } else {
440 None
441 }
442 } else {
443 None
444 }
445 };
446
447 match cached_relays {
448 Some(relays) => relays,
449 None => {
450 let result = fetch_fn().await;
452
453 {
455 let mut cache = INBOX_RELAY_CACHE.lock().unwrap();
456 cache.insert(
457 *pubkey,
458 CachedRelays {
459 relays: result.relays.clone(),
460 fetched_at: Instant::now(),
461 fetch_ok: result.fetch_ok,
462 },
463 );
464 }
465
466 result.relays
467 }
468 }
469 }; drop(cleanup_guard);
474 relays
475}
476
477async fn get_or_fetch_inbox_relays(client: &Client, pubkey: &PublicKey) -> Vec<String> {
479 get_or_fetch_with_lock(pubkey, || fetch_inbox_relays(client, pubkey)).await
480}
481
482static TRUSTED_RELAY_URLS: LazyLock<Vec<RelayUrl>> = LazyLock::new(|| {
488 crate::state::TRUSTED_RELAYS
489 .iter()
490 .filter_map(|s| RelayUrl::parse(s).ok())
491 .collect()
492});
493
494pub fn trusted_relay_urls() -> Vec<RelayUrl> {
496 TRUSTED_RELAY_URLS.clone()
497}
498
499pub async fn send_event_first_ok(
509 client: &Client,
510 urls: Vec<RelayUrl>,
511 event: &Event,
512) -> Result<nostr_sdk::prelude::SendEventOutput, nostr_sdk::prelude::Error> {
513 let pool = client;
514 let relays = pool.relays().await;
515 let event_id = event.id;
516
517 let mut resolved: Vec<(RelayUrl, Relay)> = Vec::new();
519 for url in urls {
520 if let Some(relay) = relays.get(&url) {
521 resolved.push((url, relay.clone()));
522 }
523 }
524
525 if resolved.is_empty() {
526 return client.send_event(event).await;
527 }
528
529 let handles = spawn_tracked_publish(resolved, event.clone());
533
534 let mut output = Output::new(event_id);
536
537 let mut remaining = handles;
538 while !remaining.is_empty() {
539 let (result, _index, rest) = futures_util::future::select_all(remaining).await;
540 remaining = rest;
541
542 if let Ok((url, relay_result)) = result {
543 match relay_result {
544 Ok(_) => {
545 output.success.insert(url, nostr_sdk::prelude::EventSendStatus::Sent);
546 drop(remaining);
550 return Ok(output);
551 }
552 Err(e) => {
553 output.failed.insert(url, e);
554 }
555 }
556 }
557 }
558
559 Ok(output)
561}
562
563pub async fn send_event_pool_first_ok(
566 client: &Client,
567 event: &Event,
568) -> Result<nostr_sdk::prelude::SendEventOutput, nostr_sdk::prelude::Error> {
569 let pool = client;
570 let relays = pool.relays().await;
571 let write_urls: Vec<RelayUrl> = relays
572 .iter()
573 .filter(|(_, r)| r.capabilities().load().can_write())
574 .map(|(url, _)| url.clone())
575 .collect();
576 send_event_first_ok(&client, write_urls, event).await
577}
578
579pub fn wrap_with_retained_key(
589 receiver: &PublicKey,
590 seal: &Event,
591 extra_tags: impl IntoIterator<Item = Tag>,
592) -> Result<(Event, SecretKey), String> {
593 use nostr_sdk::prelude::nip44;
594
595 if seal.kind != Kind::Seal {
596 return Err(format!("expected Seal kind, got {:?}", seal.kind));
597 }
598 let keys = Keys::generate();
599 let secret = keys.secret_key().clone();
600 let content = nip44::encrypt(
601 keys.secret_key(),
602 receiver,
603 seal.as_json(),
604 nip44::Version::default(),
605 )
606 .map_err(|e| format!("nip44 encrypt: {}", e))?;
607 let mut tags: Vec<Tag> = extra_tags.into_iter().collect();
608 tags.push(Tag::public_key(*receiver));
609 let event = EventBuilder::new(Kind::GiftWrap, content)
610 .tags(tags)
611 .custom_created_at(crate::sending::tweaked_timestamp())
612 .finalize(&keys)
613 .map_err(|e| format!("sign wrap: {}", e))?;
614 Ok((event, secret))
615}
616
617pub struct GiftWrapSendOutcome {
621 pub output: nostr_sdk::prelude::SendEventOutput,
622 pub wrap_event_id: EventId,
623 pub wrap_secret: SecretKey,
624 pub targeted_relays: Vec<String>,
627}
628
629pub struct BuiltGiftWrap {
636 pub event: Event,
637 pub secret: SecretKey,
638}
639
640pub async fn build_gift_wrap_retained(
642 _client: &Client,
643 recipient: &PublicKey,
644 rumor: UnsignedEvent,
645 extra_tags: impl IntoIterator<Item = Tag>,
646) -> Result<BuiltGiftWrap, String> {
647 let signer = crate::signer::active_signer().map_err(|e| e.to_string())?;
648 let seal: Event = nostr_sdk::prelude::GiftWrapSealBuilder::new(rumor, *recipient)
649 .finalize_async(&signer)
650 .await
651 .map_err(|e| e.to_string())?;
652 let (event, secret) = wrap_with_retained_key(recipient, &seal, extra_tags)?;
653 Ok(BuiltGiftWrap { event, secret })
654}
655
656pub struct GiftWrapTargets {
662 pub resolved: Vec<(RelayUrl, Relay)>,
663 pub targeted_relays: Vec<String>,
666 transient_added: Vec<RelayUrl>,
667}
668
669pub async fn send_gift_wrap_retained(
680 client: &Client,
681 recipient: &PublicKey,
682 rumor: UnsignedEvent,
683 extra_tags: impl IntoIterator<Item = Tag>,
684) -> Result<GiftWrapSendOutcome, String> {
685 let built = build_gift_wrap_retained(client, recipient, rumor, extra_tags).await?;
686 let targets = resolve_gift_wrap_targets(client, recipient).await;
687 let publish_result = publish_gift_wrap_to_targets(client, &targets, &built.event).await;
688 teardown_gift_wrap_targets(client, &targets).await;
689 Ok(GiftWrapSendOutcome {
690 output: publish_result?,
691 wrap_event_id: built.event.id,
692 wrap_secret: built.secret,
693 targeted_relays: targets.targeted_relays,
694 })
695}
696
697pub async fn resolve_gift_wrap_targets(
702 client: &Client,
703 recipient: &PublicKey,
704) -> GiftWrapTargets {
705 let inbox_strs = get_or_fetch_inbox_relays(client, recipient).await;
706 let targeted_strs: Vec<String> = if !inbox_strs.is_empty() {
707 inbox_strs.clone()
708 } else {
709 let pool = client;
710 let relays = pool.relays().await;
711 relays.iter()
712 .filter(|(_, r)| r.capabilities().load().can_write())
713 .map(|(url, _)| url.to_string())
714 .collect()
715 };
716 use normalize_relay_url as normalize_url_for_match;
723 let pool = client;
724 let pool_relays = pool.relays().all().await;
728 let pool_norm: Vec<(String, RelayUrl, Relay)> = pool_relays.iter()
729 .map(|(url, relay)| (
730 normalize_url_for_match(&url.to_string()),
731 url.clone(),
732 relay.clone(),
733 ))
734 .collect();
735 let mut resolved: Vec<(RelayUrl, Relay)> = targeted_strs
736 .iter()
737 .filter_map(|s| {
738 let norm = normalize_url_for_match(s);
739 pool_norm.iter()
740 .find(|(pnorm, _, _)| pnorm == &norm)
741 .map(|(_, url, relay)| (url.clone(), relay.clone()))
742 })
743 .collect();
744
745 let mut transient_added: Vec<RelayUrl> = Vec::new();
751 if !inbox_strs.is_empty() {
752 for s in &targeted_strs {
753 let norm = normalize_url_for_match(s);
754 let in_pool = pool_norm.iter().any(|(p, _, _)| p == &norm);
755 let already_added = transient_added.iter()
756 .any(|u| normalize_url_for_match(&u.to_string()) == norm);
757 if in_pool || already_added { continue; }
758 if pool.add_managed_relay(s.as_str()).await.is_ok() {
759 if let Ok(Some(relay)) = pool.relay(s.as_str()).await {
760 let _ = relay.try_connect().timeout(crate::relay_connect_timeout(std::time::Duration::from_secs(6))).await;
761 transient_added.push(relay.url().clone());
762 resolved.push((relay.url().clone(), relay));
763 }
764 }
765 }
766 if !transient_added.is_empty() {
767 crate::log_info!(
768 "[InboxRelays] on-demand connected {} inbox relay(s) for {} (transient)",
769 transient_added.len(),
770 recipient,
771 );
772 }
773 }
774
775 if !inbox_strs.is_empty() {
776 println!(
777 "[InboxRelays] Routing gift-wrap to {} inbox relays for {}",
778 resolved.len(),
779 recipient
780 );
781 }
782
783 GiftWrapTargets {
784 resolved,
785 targeted_relays: targeted_strs,
786 transient_added,
787 }
788}
789
790pub async fn reconnect_gift_wrap_targets(targets: &GiftWrapTargets) {
795 let stale: Vec<&Relay> = targets.resolved.iter()
796 .filter(|(_, r)| r.status() != RelayStatus::Connected)
797 .map(|(_, r)| r)
798 .collect();
799 if stale.is_empty() {
800 return;
801 }
802 futures_util::future::join_all(stale.into_iter().map(|r| async move {
804 r.try_connect().timeout(crate::relay_connect_timeout(std::time::Duration::from_secs(6))).await
805 }))
806 .await;
807}
808
809pub async fn publish_gift_wrap_to_targets(
821 client: &Client,
822 targets: &GiftWrapTargets,
823 event: &Event,
824) -> Result<nostr_sdk::prelude::SendEventOutput, String> {
825 if targets.resolved.is_empty() {
828 return client
831 .send_event(event)
832 .await
833 .map_err(|e| e.to_string());
834 }
835
836 let handles = spawn_tracked_publish(targets.resolved.clone(), event.clone());
837
838 let mut output = Output::new(event.id);
843 let mut remaining = handles;
844 while !remaining.is_empty() {
845 let (result, _idx, rest) = futures_util::future::select_all(remaining).await;
846 remaining = rest;
847 if let Ok((url, relay_result)) = result {
848 match relay_result {
849 Ok(_) => {
850 output.success.insert(url, nostr_sdk::prelude::EventSendStatus::Sent);
851 drop(remaining);
852 break;
853 }
854 Err(e) => {
855 output.failed.insert(url, e.to_string());
856 }
857 }
858 }
859 }
860 Ok(output)
861}
862
863pub async fn teardown_gift_wrap_targets(client: &Client, targets: &GiftWrapTargets) {
868 let pool = client;
869 for url in &targets.transient_added {
870 let _ = pool.remove_relay(url).await;
871 }
872}
873
874pub async fn send_gift_wrap(
887 client: &Client,
888 recipient: &PublicKey,
889 rumor: UnsignedEvent,
890 extra_tags: impl IntoIterator<Item = Tag>,
891) -> Result<nostr_sdk::prelude::SendEventOutput, String> {
892 let outcome = send_gift_wrap_retained(client, recipient, rumor, extra_tags).await?;
893 Ok(outcome.output)
894}
895
896const CONTRIBUTED_KEY: &str = "dm_relays_contributed";
906
907const MAX_FOREIGN_RELAYS: usize = 10;
910
911fn load_contributed() -> HashSet<String> {
912 crate::db::get_sql_setting(CONTRIBUTED_KEY.to_string())
913 .ok()
914 .flatten()
915 .and_then(|json| serde_json::from_str::<Vec<String>>(&json).ok())
916 .map(|v| v.into_iter().map(|s| normalize_relay_url(&s)).collect())
917 .unwrap_or_default()
918}
919
920fn store_contributed(contributed: &[String]) {
921 if let Ok(json) = serde_json::to_string(contributed) {
922 let _ = crate::db::set_sql_setting(CONTRIBUTED_KEY.to_string(), json);
923 }
924}
925
926struct MergePlan {
928 list: Vec<String>,
931 changed: bool,
933 contributed: Vec<String>,
936}
937
938fn merge_inbox_relays(
943 remote: &[String],
944 contributed_before: &HashSet<String>,
945 ours: &[String],
946) -> MergePlan {
947 let mut seen: HashSet<String> = HashSet::new();
948 let mut list: Vec<String> = Vec::new();
949 let mut foreign_norm: HashSet<String> = HashSet::new();
950 let mut dropped_foreign = 0usize;
951
952 for url in remote {
953 let norm = normalize_relay_url(url);
954 if seen.contains(&norm) || contributed_before.contains(&norm) {
955 continue;
956 }
957 if foreign_norm.len() >= MAX_FOREIGN_RELAYS {
958 dropped_foreign += 1;
959 continue;
960 }
961 seen.insert(norm.clone());
962 foreign_norm.insert(norm);
963 list.push(url.clone());
964 }
965 if dropped_foreign > 0 {
966 crate::log_warn!(
967 "[InboxRelays] remote 10050 over the {}-relay foreign cap, dropped {}",
968 MAX_FOREIGN_RELAYS,
969 dropped_foreign
970 );
971 }
972
973 let mut contributed: Vec<String> = Vec::new();
974 for url in ours {
975 let norm = normalize_relay_url(url);
976 if seen.insert(norm.clone()) {
977 list.push(url.clone());
978 }
979 if !foreign_norm.contains(&norm) && !contributed.contains(&norm) {
980 contributed.push(norm);
981 }
982 }
983
984 let remote_set: HashSet<String> = remote.iter().map(|s| normalize_relay_url(s)).collect();
989 let ours_norm: HashSet<String> = ours.iter().map(|s| normalize_relay_url(s)).collect();
990 let has_addition = ours_norm.iter().any(|n| !remote_set.contains(n));
991 let has_removal = contributed_before
992 .iter()
993 .any(|n| remote_set.contains(n) && !ours_norm.contains(n));
994 MergePlan { list, changed: has_addition || has_removal, contributed }
995}
996
997pub async fn fetch_own_inbox_list(client: &Client) -> Result<Option<(Vec<String>, u64)>, String> {
1003 let me = crate::state::my_public_key().ok_or("no active pubkey")?;
1004 let targets = inbox_query_targets(client).await;
1005 if targets.is_empty() {
1006 return Err("no query targets in pool".to_string());
1007 }
1008
1009 let discovery: HashSet<String> = crate::state::discovery_relay_iter()
1016 .map(normalize_relay_url)
1017 .collect();
1018 let deadline = Instant::now() + std::time::Duration::from_secs(8);
1019 loop {
1020 let relays = client.relays().all().await;
1021 let connected: Vec<&RelayUrl> = targets
1022 .iter()
1023 .filter(|url| {
1024 relays
1025 .get(url)
1026 .map(|r| r.status() == RelayStatus::Connected)
1027 .unwrap_or(false)
1028 })
1029 .collect();
1030 let discovery_up = connected
1031 .iter()
1032 .any(|url| discovery.contains(&normalize_relay_url(url.as_str())));
1033 if discovery_up {
1034 break;
1035 }
1036 if Instant::now() >= deadline {
1037 if connected.is_empty() {
1038 return Err("no query target connected".to_string());
1039 }
1040 break;
1041 }
1042 tokio::time::sleep(std::time::Duration::from_millis(250)).await;
1043 }
1044
1045 let filter = Filter::new().author(me).kind(Kind::Custom(10050)).limit(1);
1046 let events = client
1047 .fetch_events(nostr_sdk::prelude::ReqTarget::manual(
1048 targets.iter().cloned().map(|u| (u, vec![filter.clone()])),
1049 ))
1050 .timeout(std::time::Duration::from_secs(6))
1051 .await
1052 .map_err(|e| e.to_string())?;
1053 let newest = events
1055 .into_iter()
1056 .max_by(|a, b| a.created_at.cmp(&b.created_at).then(b.id.cmp(&a.id)))
1057 .map(|e| (parse_relay_tags(&e.tags), e.created_at.as_secs()));
1058
1059 if newest.is_none() {
1066 let has_discovery_target = targets
1067 .iter()
1068 .any(|url| discovery.contains(&normalize_relay_url(url.as_str())));
1069 if has_discovery_target {
1070 let relays = client.relays().all().await;
1071 let discovery_answered = targets.iter().any(|url| {
1072 discovery.contains(&normalize_relay_url(url.as_str()))
1073 && relays
1074 .get(url)
1075 .map(|r| r.status() == RelayStatus::Connected)
1076 .unwrap_or(false)
1077 });
1078 if !discovery_answered {
1079 return Err(
1080 "no 10050 found and no Discovery Relay connected; refusing to bootstrap"
1081 .to_string(),
1082 );
1083 }
1084 }
1085 }
1086 Ok(newest)
1087}
1088
1089const LIST_SEEN_TS_KEY: &str = "dm_list_last_ts";
1094
1095fn load_list_seen() -> u64 {
1096 crate::db::get_sql_setting(LIST_SEEN_TS_KEY.to_string())
1097 .ok()
1098 .flatten()
1099 .and_then(|v| v.parse::<u64>().ok())
1100 .unwrap_or(0)
1101}
1102
1103pub fn note_contributed(urls: &[String]) {
1108 if urls.is_empty() {
1109 return;
1110 }
1111 let mut set = load_contributed();
1112 for url in urls {
1113 set.insert(normalize_relay_url(url));
1114 }
1115 let list: Vec<String> = set.into_iter().collect();
1116 store_contributed(&list);
1117}
1118
1119pub fn note_list_seen(ts: u64) {
1121 if ts > load_list_seen() {
1122 let _ = crate::db::set_sql_setting(LIST_SEEN_TS_KEY.to_string(), ts.to_string());
1123 }
1124}
1125
1126#[derive(Debug, Default, PartialEq)]
1129pub struct InboundReconcile {
1130 pub adopt: Vec<String>,
1132 pub revive: Vec<String>,
1134 pub retire: Vec<String>,
1137}
1138
1139pub fn plan_inbound_reconcile(
1143 remote: &[String],
1144 remote_ts: u64,
1145 ours: &[String],
1146 declined: &[String],
1147) -> InboundReconcile {
1148 plan_inbound_reconcile_pure(
1149 remote,
1150 remote_ts,
1151 ours,
1152 declined,
1153 &load_contributed(),
1154 load_list_seen(),
1155 )
1156}
1157
1158fn plan_inbound_reconcile_pure(
1159 remote: &[String],
1160 remote_ts: u64,
1161 ours: &[String],
1162 declined: &[String],
1163 contributed_before: &HashSet<String>,
1164 last_seen_ts: u64,
1165) -> InboundReconcile {
1166 if remote_ts <= last_seen_ts {
1167 return InboundReconcile::default();
1168 }
1169 let ours_norm: HashSet<String> = ours.iter().map(|s| normalize_relay_url(s)).collect();
1170 let declined_norm: HashSet<String> = declined.iter().map(|s| normalize_relay_url(s)).collect();
1171 let remote_norm: HashSet<String> = remote.iter().map(|s| normalize_relay_url(s)).collect();
1172
1173 let mut seen: HashSet<String> = HashSet::new();
1174 let mut adopt: Vec<String> = Vec::new();
1175 let mut revive: Vec<String> = Vec::new();
1176 for url in remote {
1177 let norm = normalize_relay_url(url);
1178 if !seen.insert(norm.clone()) {
1179 continue;
1180 }
1181 if ours_norm.contains(&norm) {
1182 continue;
1183 }
1184 if declined_norm.contains(&norm) {
1185 revive.push(url.clone());
1186 } else if adopt.len() < MAX_FOREIGN_RELAYS
1187 && url.starts_with("wss://")
1188 && url.len() <= 256
1189 {
1190 adopt.push(url.clone());
1191 }
1192 }
1193
1194 let retire: Vec<String> = ours
1195 .iter()
1196 .filter(|url| {
1197 let norm = normalize_relay_url(url);
1198 contributed_before.contains(&norm) && !remote_norm.contains(&norm)
1199 })
1200 .cloned()
1201 .collect();
1202
1203 InboundReconcile { adopt, revive, retire }
1204}
1205
1206static PUBLISH_MUTEX: tokio::sync::Mutex<()> = tokio::sync::Mutex::const_new(());
1210
1211pub async fn publish_inbox_relays(client: &Client) -> Result<(), String> {
1217 let remote = fetch_own_inbox_list(client).await?;
1220 publish_inbox_relays_synced(client, remote, None).await
1221}
1222
1223pub async fn publish_inbox_relays_synced(
1230 client: &Client,
1231 remote: Option<(Vec<String>, u64)>,
1232 ours_override: Option<Vec<String>>,
1233) -> Result<(), String> {
1234 let _serial = PUBLISH_MUTEX.lock().await;
1235 let session = crate::state::SessionGuard::capture();
1236
1237 let ours: Vec<String> = match ours_override {
1239 Some(list) => list,
1240 None => client
1241 .relays()
1242 .await
1243 .iter()
1244 .filter(|(_, relay)| relay.capabilities().load().can_read())
1245 .map(|(url, _)| url.to_string())
1246 .collect(),
1247 };
1248
1249 let remote_found = remote.is_some();
1250 let (remote, remote_ts) = remote.unwrap_or_default();
1251 if remote_found && remote_ts < load_list_seen() {
1256 return Err("stale 10050 fetch (older than last seen), skipping publish".to_string());
1257 }
1258
1259 let plan = merge_inbox_relays(&remote, &load_contributed(), &ours);
1260
1261 if !session.is_valid() {
1262 return Ok(());
1263 }
1264 store_contributed(&plan.contributed);
1265 if remote_found {
1266 note_list_seen(remote_ts);
1267 }
1268
1269 if remote_found && !plan.changed {
1270 crate::log_info!(
1271 "[InboxRelays] kind 10050 already in sync ({} relay(s)), not publishing",
1272 plan.list.len()
1273 );
1274 return Ok(());
1275 }
1276 if plan.list.is_empty() && !remote_found {
1277 return Ok(());
1279 }
1280
1281 let mut builder = EventBuilder::new(Kind::Custom(10050), "");
1282 for url in &plan.list {
1283 builder = builder.tag(Tag::custom("relay", vec![url.clone()]));
1284 }
1285 let event = crate::sign_builder(builder)
1286 .await
1287 .map_err(|e| format!("Failed to sign inbox relays: {}", e))?;
1288
1289 if !session.is_valid() {
1290 return Ok(());
1291 }
1292 let pool_send = client.send_event(&event).await;
1293
1294 let discovery: HashSet<String> = crate::state::DISCOVERY_RELAYS
1298 .iter()
1299 .map(|s| normalize_relay_url(s))
1300 .collect();
1301 let discovery_targets: Vec<RelayUrl> = client
1302 .relays().all()
1303 .await
1304 .iter()
1305 .filter(|(url, relay)| {
1306 !relay.capabilities().load().can_write() && discovery.contains(&normalize_relay_url(url.as_str()))
1307 })
1308 .map(|(url, _)| url.clone())
1309 .collect();
1310 let mut discovery_ok = false;
1311 if !discovery_targets.is_empty() {
1312 if let Ok(out) = client.send_event(&event).to(discovery_targets).await {
1313 discovery_ok = !out.success.is_empty();
1314 }
1315 }
1316
1317 let pool_ok = matches!(&pool_send, Ok(out) if !out.success.is_empty());
1318 if !pool_ok {
1319 if !discovery_ok {
1320 return Err(match pool_send {
1321 Err(e) => format!("Failed to publish inbox relays: {}", e),
1322 Ok(_) => "Failed to publish inbox relays: no relay accepted it".to_string(),
1323 });
1324 }
1325 crate::log_warn!(
1326 "[InboxRelays] pool publish failed, list delivered via Discovery Relays only"
1327 );
1328 }
1329 if session.is_valid() {
1332 note_list_seen(event.created_at.as_secs().max(remote_ts));
1333 }
1334
1335 println!(
1336 "[InboxRelays] Published kind 10050 with {} relay(s) ({} foreign preserved)",
1337 plan.list.len(),
1338 plan.list.len().saturating_sub(plan.contributed.len())
1339 );
1340 Ok(())
1341}
1342
1343static REPUBLISH_GEN: AtomicU64 = AtomicU64::new(0);
1346
1347#[cfg(test)]
1349static DEBOUNCE_PASS_COUNT: AtomicU64 = AtomicU64::new(0);
1350
1351pub fn republish_inbox_relays_debounced() {
1355 let gen = REPUBLISH_GEN.fetch_add(1, Ordering::SeqCst) + 1;
1356 let session = crate::state::SessionGuard::capture();
1361 tokio::spawn(async move {
1362 tokio::time::sleep(std::time::Duration::from_millis(800)).await;
1365 if REPUBLISH_GEN.load(Ordering::SeqCst) != gen {
1366 return; }
1368 if !session.is_valid() {
1369 return; }
1371 #[cfg(test)]
1372 DEBOUNCE_PASS_COUNT.fetch_add(1, Ordering::SeqCst);
1373 let client = match nostr_client() {
1374 Some(c) => c,
1375 None => return,
1376 };
1377 if let Err(e) = publish_inbox_relays(&client).await {
1378 eprintln!("[InboxRelays] Failed to republish after config change: {}", e);
1379 }
1380 });
1381}
1382
1383#[cfg(test)]
1384mod tests {
1385 use super::*;
1386
1387 fn strs(v: &[&str]) -> Vec<String> {
1390 v.iter().map(|s| s.to_string()).collect()
1391 }
1392
1393 fn norm_set(v: &[&str]) -> HashSet<String> {
1394 v.iter().map(|s| normalize_relay_url(s)).collect()
1395 }
1396
1397 #[test]
1398 fn merge_preserves_foreign_entries() {
1399 let remote = strs(&["wss://other-app.example", "wss://alice.example"]);
1400 let ours = strs(&["wss://vector.example"]);
1401 let plan = merge_inbox_relays(&remote, &HashSet::new(), &ours);
1402 assert!(plan.changed);
1403 assert_eq!(plan.list, strs(&["wss://other-app.example", "wss://alice.example", "wss://vector.example"]));
1404 assert_eq!(plan.contributed, strs(&["wss://vector.example"]));
1405 }
1406
1407 #[test]
1408 fn merge_noop_when_remote_covers_ours() {
1409 let remote = strs(&["wss://other-app.example", "wss://vector.example/"]);
1410 let ours = strs(&["wss://vector.example"]);
1411 let plan = merge_inbox_relays(&remote, &HashSet::new(), &ours);
1412 assert!(!plan.changed, "trailing-slash variants are the same relay");
1413 assert_eq!(plan.list.len(), 2);
1414 }
1415
1416 #[test]
1417 fn merge_drops_only_our_own_removed_contribution() {
1418 let remote = strs(&["wss://foreign.example", "wss://x.example"]);
1421 let contributed = norm_set(&["wss://x.example"]);
1422 let ours = strs(&["wss://new.example"]);
1423 let plan = merge_inbox_relays(&remote, &contributed, &ours);
1424 assert!(plan.changed);
1425 assert_eq!(plan.list, strs(&["wss://foreign.example", "wss://new.example"]));
1426 }
1427
1428 #[test]
1429 fn merge_never_clears_a_foreign_list() {
1430 let remote = strs(&["wss://foreign.example"]);
1432 let plan = merge_inbox_relays(&remote, &HashSet::new(), &[]);
1433 assert!(!plan.changed);
1434 assert_eq!(plan.list, remote);
1435 assert!(plan.contributed.is_empty());
1436 }
1437
1438 #[test]
1439 fn merge_contributed_excludes_foreign_overlap() {
1440 let remote = strs(&["wss://shared.example"]);
1446 let ours = strs(&["wss://shared.example", "wss://mine.example"]);
1447 let plan = merge_inbox_relays(&remote, &HashSet::new(), &ours);
1448 assert_eq!(plan.contributed, strs(&["wss://mine.example"]));
1449 let next = merge_inbox_relays(
1450 &plan.list,
1451 &plan.contributed.iter().cloned().collect(),
1452 &[],
1453 );
1454 assert!(next.list.contains(&"wss://shared.example".to_string()));
1455 assert!(!next.list.contains(&"wss://mine.example".to_string()));
1456 }
1457
1458 #[test]
1459 fn merge_caps_foreign_bloat_without_publishing() {
1460 let remote: Vec<String> = (0..30).map(|i| format!("wss://r{}.example", i)).collect();
1464 let plan = merge_inbox_relays(&remote, &HashSet::new(), &[]);
1465 assert_eq!(plan.list.len(), MAX_FOREIGN_RELAYS);
1466 assert!(!plan.changed, "a trim alone must not drive a publish");
1467 }
1468
1469 #[test]
1470 fn merge_cap_applies_when_own_diff_publishes() {
1471 let remote: Vec<String> = (0..30).map(|i| format!("wss://r{}.example", i)).collect();
1472 let ours = strs(&["wss://mine.example"]);
1473 let plan = merge_inbox_relays(&remote, &HashSet::new(), &ours);
1474 assert!(plan.changed, "our addition is a real diff");
1475 assert_eq!(plan.list.len(), MAX_FOREIGN_RELAYS + 1);
1476 assert!(plan.list.contains(&"wss://mine.example".to_string()));
1477 }
1478
1479 #[test]
1480 fn merge_two_devices_reach_fixpoint() {
1481 let ours_a = strs(&["wss://a1.example", "wss://shared.example"]);
1484 let ours_b = strs(&["wss://b1.example", "wss://shared.example"]);
1485 let mut network = strs(&["wss://foreign.example"]);
1486 let mut contributed_a: HashSet<String> = HashSet::new();
1487 let mut contributed_b: HashSet<String> = HashSet::new();
1488 let mut publishes = 0;
1489 for round in 0..6 {
1490 for device in 0..2 {
1491 let (ours, contributed) = if device == 0 {
1492 (&ours_a, &mut contributed_a)
1493 } else {
1494 (&ours_b, &mut contributed_b)
1495 };
1496 let plan = merge_inbox_relays(&network, contributed, ours);
1497 *contributed = plan.contributed.iter().cloned().collect();
1498 if plan.changed {
1499 publishes += 1;
1500 network = plan.list;
1501 }
1502 if round >= 2 {
1503 assert!(!plan.changed, "no publish after convergence (round {round})");
1504 }
1505 }
1506 }
1507 assert!(publishes <= 2, "one publish per device to converge, got {publishes}");
1508 for url in ["wss://foreign.example", "wss://a1.example", "wss://b1.example", "wss://shared.example"] {
1509 assert!(network.contains(&url.to_string()), "union must hold {url}");
1510 }
1511 }
1512
1513 #[test]
1514 fn merge_first_run_publishes_ours() {
1515 let ours = strs(&["wss://a.example", "wss://b.example"]);
1516 let plan = merge_inbox_relays(&[], &HashSet::new(), &ours);
1517 assert!(plan.changed);
1518 assert_eq!(plan.list, ours);
1519 assert_eq!(plan.contributed, ours);
1520 }
1521
1522 #[test]
1525 fn reconcile_stale_remote_is_a_no_op() {
1526 let remote = strs(&["wss://foreign.example"]);
1527 let plan = plan_inbound_reconcile_pure(&remote, 100, &[], &[], &HashSet::new(), 100);
1528 assert_eq!(plan, InboundReconcile::default(), "ts <= last_seen must not act");
1529 }
1530
1531 #[test]
1532 fn reconcile_adopts_unknown_entries_capped_and_wss_only() {
1533 let mut remote: Vec<String> = (0..12).map(|i| format!("wss://r{}.example", i)).collect();
1534 remote.push("ws://plaintext.example".to_string());
1535 remote.push("http://nope.example".to_string());
1536 let plan = plan_inbound_reconcile_pure(&remote, 200, &[], &[], &HashSet::new(), 100);
1537 assert_eq!(plan.adopt.len(), MAX_FOREIGN_RELAYS);
1538 assert!(plan.adopt.iter().all(|u| u.starts_with("wss://")));
1539 assert!(plan.revive.is_empty() && plan.retire.is_empty());
1540 }
1541
1542 #[test]
1543 fn reconcile_revives_locally_disabled_entry() {
1544 let remote = strs(&["wss://back.example"]);
1545 let declined = strs(&["wss://back.example/"]);
1546 let plan = plan_inbound_reconcile_pure(&remote, 200, &[], &declined, &HashSet::new(), 100);
1547 assert_eq!(plan.revive, strs(&["wss://back.example"]));
1548 assert!(plan.adopt.is_empty());
1549 }
1550
1551 #[test]
1552 fn reconcile_retires_contributed_entry_dropped_by_newer_remote() {
1553 let remote = strs(&["wss://keep.example"]);
1554 let ours = strs(&["wss://keep.example", "wss://gone.example"]);
1555 let contributed = norm_set(&["wss://keep.example", "wss://gone.example"]);
1556 let plan = plan_inbound_reconcile_pure(&remote, 200, &ours, &[], &contributed, 100);
1557 assert_eq!(plan.retire, strs(&["wss://gone.example"]));
1558 }
1559
1560 #[test]
1561 fn reconcile_never_retires_unpublished_local_addition() {
1562 let remote = strs(&["wss://old.example"]);
1565 let ours = strs(&["wss://old.example", "wss://just-added.example"]);
1566 let contributed = norm_set(&["wss://old.example"]);
1567 let plan = plan_inbound_reconcile_pure(&remote, 200, &ours, &[], &contributed, 100);
1568 assert!(plan.retire.is_empty());
1569 }
1570
1571 #[test]
1572 fn reconcile_two_devices_propagates_default_disable() {
1573 #[derive(Clone)]
1577 struct Device {
1578 ours: Vec<String>,
1579 declined: Vec<String>,
1580 contributed: HashSet<String>,
1581 last_seen: u64,
1582 }
1583 impl Device {
1584 fn new(defaults: &[&str]) -> Self {
1585 Device {
1586 ours: strs(defaults),
1587 declined: Vec::new(),
1588 contributed: HashSet::new(),
1589 last_seen: 0,
1590 }
1591 }
1592 fn sync(&mut self, network: &mut (Vec<String>, u64)) -> bool {
1594 let (remote, ts) = network.clone();
1595 for u in &self.ours {
1596 if remote.iter().any(|r| normalize_relay_url(r) == normalize_relay_url(u)) {
1597 self.contributed.insert(normalize_relay_url(u));
1598 }
1599 }
1600 let plan = plan_inbound_reconcile_pure(
1601 &remote, ts, &self.ours, &self.declined, &self.contributed, self.last_seen,
1602 );
1603 for u in &plan.retire {
1604 self.ours.retain(|o| o != u);
1605 self.declined.push(u.clone());
1606 }
1607 for u in &plan.revive {
1608 self.declined.retain(|d| normalize_relay_url(d) != normalize_relay_url(u));
1609 self.ours.push(u.clone());
1610 self.contributed.insert(normalize_relay_url(u));
1611 }
1612 for u in &plan.adopt {
1613 self.ours.push(u.clone());
1614 self.contributed.insert(normalize_relay_url(u));
1615 }
1616 self.last_seen = self.last_seen.max(ts);
1617 let m = merge_inbox_relays(&remote, &self.contributed, &self.ours);
1618 self.contributed = m.contributed.iter().cloned().collect();
1619 if m.changed {
1620 network.1 += 1;
1621 network.0 = m.list;
1622 self.last_seen = network.1;
1623 }
1624 m.changed
1625 }
1626 }
1627
1628 const DEFAULTS: &[&str] = &["wss://d1.example", "wss://d2.example"];
1629 let mut network: (Vec<String>, u64) = (Vec::new(), 0);
1630 let mut a = Device::new(DEFAULTS);
1631 let mut b = Device::new(DEFAULTS);
1632
1633 assert!(a.sync(&mut network), "first device bootstraps the list");
1634 assert!(!b.sync(&mut network), "second device is already in sync");
1635
1636 b.ours.retain(|u| u != "wss://d2.example");
1638 b.declined.push("wss://d2.example".to_string());
1639 assert!(b.sync(&mut network), "disable must publish");
1640 assert!(!network.0.contains(&"wss://d2.example".to_string()));
1641
1642 assert!(!a.sync(&mut network), "A must adopt the removal, not republish d2");
1644 assert!(a.declined.contains(&"wss://d2.example".to_string()));
1645 assert!(!network.0.contains(&"wss://d2.example".to_string()), "no resurrection");
1646
1647 b.declined.retain(|u| u != "wss://d2.example");
1649 b.ours.push("wss://d2.example".to_string());
1650 assert!(b.sync(&mut network), "re-enable must publish");
1651 assert!(!a.sync(&mut network), "revive is inbound-only, no republish");
1652 assert!(a.ours.contains(&"wss://d2.example".to_string()), "A revives d2");
1653
1654 for _ in 0..3 {
1656 assert!(!a.sync(&mut network));
1657 assert!(!b.sync(&mut network));
1658 }
1659 }
1660
1661 #[test]
1664 fn parse_relay_tags_extracts_urls() {
1665 let tags = Tags::from_list(vec![
1666 Tag::custom("relay", vec!["wss://relay.example.com"]),
1667 Tag::custom("relay", vec!["wss://other.example.com"]),
1668 ]);
1669 let result = parse_relay_tags(&tags);
1670 assert_eq!(result, vec![
1671 "wss://relay.example.com".to_string(),
1672 "wss://other.example.com".to_string(),
1673 ]);
1674 }
1675
1676 #[test]
1677 fn parse_relay_tags_ignores_non_relay_tags() {
1678 let tags = Tags::from_list(vec![
1679 Tag::custom("relay", vec!["wss://good.example.com"]),
1680 Tag::custom("p", vec!["deadbeef"]),
1681 Tag::custom("e", vec!["cafebabe"]),
1682 ]);
1683 let result = parse_relay_tags(&tags);
1684 assert_eq!(result, vec!["wss://good.example.com".to_string()]);
1685 }
1686
1687 #[test]
1688 fn parse_relay_tags_empty() {
1689 let tags = Tags::new();
1690 let result = parse_relay_tags(&tags);
1691 assert!(result.is_empty());
1692 }
1693
1694 #[test]
1695 fn parse_relay_tags_ignores_relay_tag_without_value() {
1696 let tags = Tags::from_list(vec![
1698 Tag::custom("relay", Vec::<String>::new()),
1699 ]);
1700 let result = parse_relay_tags(&tags);
1701 assert!(result.is_empty());
1702 }
1703
1704 fn test_pubkey() -> PublicKey {
1707 let keys = Keys::generate();
1708 keys.public_key()
1709 }
1710
1711 static TEST_GLOBALS_LOCK: LazyLock<tokio::sync::Mutex<()>> =
1713 LazyLock::new(|| tokio::sync::Mutex::new(()));
1714
1715 #[test]
1716 fn cache_stores_and_retrieves() {
1717 let _guard = TEST_GLOBALS_LOCK.blocking_lock();
1718 let pk = test_pubkey();
1719 let relays = vec!["wss://a.example.com".to_string()];
1720
1721 {
1722 let mut cache = INBOX_RELAY_CACHE.lock().unwrap();
1723 cache.insert(pk, CachedRelays {
1724 relays: relays.clone(),
1725 fetched_at: Instant::now(),
1726 fetch_ok: true,
1727 });
1728 }
1729
1730 let cache = INBOX_RELAY_CACHE.lock().unwrap();
1731 let entry = cache.get(&pk).unwrap();
1732 assert_eq!(entry.relays, relays);
1733 assert!(entry.fetch_ok);
1734 assert!(entry.fetched_at.elapsed().as_secs() < CACHE_TTL_SECS);
1735 }
1736
1737 #[test]
1738 fn cache_expires_after_ttl() {
1739 let _guard = TEST_GLOBALS_LOCK.blocking_lock();
1740 let pk = test_pubkey();
1741
1742 {
1743 let mut cache = INBOX_RELAY_CACHE.lock().unwrap();
1744 cache.insert(pk, CachedRelays {
1745 relays: vec!["wss://stale.example.com".to_string()],
1746 fetched_at: Instant::now() - std::time::Duration::from_secs(CACHE_TTL_SECS + 1),
1747 fetch_ok: true,
1748 });
1749 }
1750
1751 let cache = INBOX_RELAY_CACHE.lock().unwrap();
1752 let entry = cache.get(&pk).unwrap();
1753 assert!(entry.fetched_at.elapsed().as_secs() >= CACHE_TTL_SECS);
1754 }
1755
1756 #[test]
1757 fn cache_stores_empty_results() {
1758 let _guard = TEST_GLOBALS_LOCK.blocking_lock();
1759 let pk = test_pubkey();
1760
1761 {
1762 let mut cache = INBOX_RELAY_CACHE.lock().unwrap();
1763 cache.insert(pk, CachedRelays {
1764 relays: vec![],
1765 fetched_at: Instant::now(),
1766 fetch_ok: true,
1767 });
1768 }
1769
1770 let cache = INBOX_RELAY_CACHE.lock().unwrap();
1771 let entry = cache.get(&pk).unwrap();
1772 assert!(entry.relays.is_empty());
1773 assert!(entry.fetch_ok);
1774 assert!(entry.fetched_at.elapsed().as_secs() < CACHE_TTL_SECS);
1775 }
1776
1777 #[test]
1778 fn cache_error_uses_short_ttl() {
1779 let _guard = TEST_GLOBALS_LOCK.blocking_lock();
1780 let pk = test_pubkey();
1781
1782 {
1783 let mut cache = INBOX_RELAY_CACHE.lock().unwrap();
1784 cache.insert(pk, CachedRelays {
1785 relays: vec![],
1786 fetched_at: Instant::now() - std::time::Duration::from_secs(120),
1788 fetch_ok: false,
1789 });
1790 }
1791
1792 let cache = INBOX_RELAY_CACHE.lock().unwrap();
1793 let entry = cache.get(&pk).unwrap();
1794 assert!(!entry.fetch_ok);
1795 assert!(entry.fetched_at.elapsed().as_secs() >= CACHE_TTL_ERROR_SECS);
1797 assert!(entry.fetched_at.elapsed().as_secs() < CACHE_TTL_SECS);
1799 }
1800
1801 #[tokio::test]
1804 async fn concurrent_fetches_for_same_pubkey_serialize() {
1805 let _guard = TEST_GLOBALS_LOCK.lock().await;
1806 let pk = test_pubkey();
1807
1808 {
1810 let mut cache = INBOX_RELAY_CACHE.lock().unwrap();
1811 cache.remove(&pk);
1812 }
1813
1814 let fetch_counter = Arc::new(AtomicU64::new(0));
1815
1816 let mut handles = vec![];
1819 for _ in 0..10 {
1820 let counter = fetch_counter.clone();
1821 let handle = tokio::spawn(async move {
1822 get_or_fetch_with_lock(&pk, || async {
1823 counter.fetch_add(1, Ordering::SeqCst);
1824 tokio::time::sleep(std::time::Duration::from_millis(50)).await;
1826 FetchResult {
1827 relays: vec!["wss://test.example.com".to_string()],
1828 fetch_ok: true,
1829 }
1830 })
1831 .await
1832 });
1833 handles.push(handle);
1834 }
1835
1836 let results = futures_util::future::join_all(handles).await;
1838
1839 for result in &results {
1841 assert!(result.is_ok());
1842 let relays = result.as_ref().unwrap();
1843 assert_eq!(relays, &vec!["wss://test.example.com".to_string()]);
1844 }
1845
1846 assert_eq!(
1848 fetch_counter.load(Ordering::SeqCst),
1849 1,
1850 "Expected exactly 1 fetch for 10 concurrent requests to same pubkey"
1851 );
1852
1853 let locks_after = {
1854 let locks = FETCH_LOCKS.lock().unwrap();
1855 locks.len()
1856 };
1857 assert_eq!(locks_after, 0, "Lock entry should be removed after all waiters complete");
1858 }
1859
1860 #[tokio::test]
1861 async fn fetch_locks_do_not_accumulate_after_calls_complete() {
1862 let _guard = TEST_GLOBALS_LOCK.lock().await;
1863
1864 let pk1 = test_pubkey();
1868 let pk2 = test_pubkey();
1869 let pk3 = test_pubkey();
1870
1871 {
1873 let mut cache = INBOX_RELAY_CACHE.lock().unwrap();
1874 cache.clear();
1875 }
1876 {
1877 let mut locks = FETCH_LOCKS.lock().unwrap();
1878 locks.clear();
1879 }
1880
1881 get_or_fetch_with_lock(&pk1, || async {
1883 FetchResult {
1884 relays: vec!["wss://relay1.example.com".to_string()],
1885 fetch_ok: true,
1886 }
1887 })
1888 .await;
1889
1890 let locks_after_pk1 = {
1893 let locks = FETCH_LOCKS.lock().unwrap();
1894 locks.len()
1895 };
1896 assert_eq!(locks_after_pk1, 0, "No lock entries should remain after pk1 call");
1897
1898 get_or_fetch_with_lock(&pk2, || async {
1900 FetchResult {
1901 relays: vec!["wss://relay2.example.com".to_string()],
1902 fetch_ok: true,
1903 }
1904 })
1905 .await;
1906
1907 let locks_after_pk2 = {
1908 let locks = FETCH_LOCKS.lock().unwrap();
1909 locks.len()
1910 };
1911 assert_eq!(locks_after_pk2, 0, "No lock entries should remain after pk2 call");
1912
1913 get_or_fetch_with_lock(&pk3, || async {
1915 FetchResult {
1916 relays: vec!["wss://relay3.example.com".to_string()],
1917 fetch_ok: true,
1918 }
1919 })
1920 .await;
1921
1922 let locks_after_pk3 = {
1923 let locks = FETCH_LOCKS.lock().unwrap();
1924 locks.len()
1925 };
1926 assert_eq!(locks_after_pk3, 0, "No lock entries should remain after pk3 call");
1927 }
1928
1929 #[tokio::test]
1930 async fn cancelled_fetch_cleans_up_lock_entry() {
1931 let _guard = TEST_GLOBALS_LOCK.lock().await;
1932 let pk = test_pubkey();
1933
1934 {
1935 let mut cache = INBOX_RELAY_CACHE.lock().unwrap();
1936 cache.clear();
1937 }
1938 {
1939 let mut locks = FETCH_LOCKS.lock().unwrap();
1940 locks.clear();
1941 }
1942
1943 let (started_tx, started_rx) = tokio::sync::oneshot::channel::<()>();
1944 let task_pk = pk;
1945 let handle = tokio::spawn(async move {
1946 get_or_fetch_with_lock(&task_pk, || async move {
1947 let _ = started_tx.send(());
1948 tokio::time::sleep(std::time::Duration::from_secs(30)).await;
1949 FetchResult { relays: Vec::new(), fetch_ok: false }
1950 })
1951 .await
1952 });
1953
1954 started_rx.await.expect("fetch closure should start before abort");
1955 handle.abort();
1956 let _ = handle.await;
1957 tokio::task::yield_now().await;
1958
1959 let locks_after = {
1960 let locks = FETCH_LOCKS.lock().unwrap();
1961 locks.len()
1962 };
1963 assert_eq!(
1964 locks_after, 0,
1965 "Lock entry should be removed even if fetch task is cancelled"
1966 );
1967 }
1968
1969 #[tokio::test(start_paused = true)]
1976 async fn debounce_coalesces_rapid_calls_into_one() {
1977 let gen_before = REPUBLISH_GEN.load(Ordering::SeqCst);
1979 let pass_before = DEBOUNCE_PASS_COUNT.load(Ordering::SeqCst);
1980
1981 republish_inbox_relays_debounced();
1983 republish_inbox_relays_debounced();
1984 republish_inbox_relays_debounced();
1985
1986 let gen_after = REPUBLISH_GEN.load(Ordering::SeqCst);
1987 assert_eq!(gen_after, gen_before + 3);
1988
1989 tokio::time::sleep(std::time::Duration::from_millis(1000)).await;
1991
1992 let pass_after = DEBOUNCE_PASS_COUNT.load(Ordering::SeqCst);
1993 assert_eq!(pass_after - pass_before, 1);
1997 }
1998}