1use super::config::BootstrapCacheConfig;
11use super::entry::{CachedPeer, ConnectionOutcome, PeerCapabilities, PeerSource};
12use super::persistence::{CacheData, CachePersistence};
13use super::selection::select_epsilon_greedy;
14use crate::nat_traversal_api::PeerId;
15use std::net::SocketAddr;
16use std::sync::Arc;
17use std::time::{Instant, SystemTime};
18use tokio::sync::{RwLock, broadcast};
19use tracing::{debug, info, warn};
20
21#[derive(Debug, Clone)]
23pub enum CacheEvent {
24 Updated {
26 peer_count: usize,
28 },
29 Saved,
31 Merged {
33 added: usize,
35 },
36 Cleaned {
38 removed: usize,
40 },
41}
42
43#[derive(Debug, Clone, Default)]
45pub struct CacheStats {
46 pub total_peers: usize,
48 pub relay_peers: usize,
50 pub coordinator_peers: usize,
52 pub dual_stack_relay_peers: usize,
54 pub average_quality: f64,
56 pub untested_peers: usize,
58}
59
60#[derive(Debug)]
66pub struct BootstrapCache {
67 config: BootstrapCacheConfig,
68 data: Arc<RwLock<CacheData>>,
69 persistence: Option<CachePersistence>,
71 event_tx: broadcast::Sender<CacheEvent>,
72 last_save: Arc<RwLock<Instant>>,
73 last_cleanup: Arc<RwLock<Instant>>,
74}
75
76impl BootstrapCache {
77 pub async fn open(config: BootstrapCacheConfig) -> std::io::Result<Self> {
84 let persistence = if config.persist {
85 Some(CachePersistence::new(
86 &config.cache_dir,
87 config.enable_file_locking,
88 )?)
89 } else {
90 None
91 };
92 let data = match &persistence {
93 Some(p) => p.load()?,
94 None => CacheData::new(super::persistence::generate_instance_id()),
95 };
96 let (event_tx, _) = broadcast::channel(256);
97 let now = Instant::now();
98
99 info!(
100 "Opened bootstrap cache with {} peers (persist: {})",
101 data.peers.len(),
102 config.persist
103 );
104
105 Ok(Self {
106 config,
107 data: Arc::new(RwLock::new(data)),
108 persistence,
109 event_tx,
110 last_save: Arc::new(RwLock::new(now)),
111 last_cleanup: Arc::new(RwLock::new(now)),
112 })
113 }
114
115 pub fn subscribe(&self) -> broadcast::Receiver<CacheEvent> {
117 self.event_tx.subscribe()
118 }
119
120 pub async fn peer_count(&self) -> usize {
122 self.data.read().await.peers.len()
123 }
124
125 pub async fn get_peer(&self, peer_id: &PeerId) -> Option<CachedPeer> {
127 let mut data = self.data.write().await;
128 let peer = data.peers.get_mut(&peer_id.0)?;
129 peer.capabilities
130 .refresh_direct_capabilities(self.config.reachability_ttl, SystemTime::now());
131 peer.calculate_quality(&self.config.weights);
132 Some(peer.clone())
133 }
134
135 fn refresh_cached_peer(&self, peer: &mut CachedPeer, now: SystemTime) {
136 peer.capabilities
137 .refresh_direct_capabilities(self.config.reachability_ttl, now);
138 peer.calculate_quality(&self.config.weights);
139 }
140
141 pub async fn select_peers(&self, count: usize) -> Vec<CachedPeer> {
146 let mut data = self.data.write().await;
147 let now = SystemTime::now();
148 for peer in data.peers.values_mut() {
149 self.refresh_cached_peer(peer, now);
150 }
151 let peers: Vec<CachedPeer> = data.peers.values().cloned().collect();
152 drop(data);
153
154 select_epsilon_greedy(&peers, count, self.config.epsilon)
155 .into_iter()
156 .cloned()
157 .collect()
158 }
159
160 pub async fn select_relay_peers(&self, count: usize) -> Vec<CachedPeer> {
164 let mut data = self.data.write().await;
165 let now = SystemTime::now();
166 for peer in data.peers.values_mut() {
167 self.refresh_cached_peer(peer, now);
168 }
169 let peers: Vec<CachedPeer> = data.peers.values().cloned().collect();
170 drop(data);
171
172 super::selection::select_with_capabilities(&peers, count, true, false)
173 .into_iter()
174 .cloned()
175 .collect()
176 }
177
178 pub async fn select_coordinators(&self, count: usize) -> Vec<CachedPeer> {
182 let mut data = self.data.write().await;
183 let now = SystemTime::now();
184 for peer in data.peers.values_mut() {
185 self.refresh_cached_peer(peer, now);
186 }
187 let peers: Vec<CachedPeer> = data.peers.values().cloned().collect();
188 drop(data);
189
190 super::selection::select_with_capabilities(&peers, count, false, true)
191 .into_iter()
192 .cloned()
193 .collect()
194 }
195
196 pub async fn select_relays_for_target(
206 &self,
207 count: usize,
208 target: &std::net::SocketAddr,
209 prefer_dual_stack: bool,
210 ) -> Vec<CachedPeer> {
211 use super::selection::select_relays_for_target;
212
213 let mut data = self.data.write().await;
214 let now = SystemTime::now();
215 for peer in data.peers.values_mut() {
216 self.refresh_cached_peer(peer, now);
217 }
218 let peers: Vec<CachedPeer> = data.peers.values().cloned().collect();
219 drop(data);
220
221 select_relays_for_target(&peers, count, *target, prefer_dual_stack)
222 .into_iter()
223 .cloned()
224 .collect()
225 }
226
227 pub async fn select_dual_stack_relays(&self, count: usize) -> Vec<CachedPeer> {
231 use super::selection::select_dual_stack_relays;
232
233 let mut data = self.data.write().await;
234 let now = SystemTime::now();
235 for peer in data.peers.values_mut() {
236 self.refresh_cached_peer(peer, now);
237 }
238 let peers: Vec<CachedPeer> = data.peers.values().cloned().collect();
239 drop(data);
240
241 select_dual_stack_relays(&peers, count)
242 .into_iter()
243 .cloned()
244 .collect()
245 }
246
247 pub async fn upsert(&self, peer: CachedPeer) {
251 let mut data = self.data.write().await;
252
253 if data.peers.len() >= self.config.max_peers && !data.peers.contains_key(&peer.peer_id.0) {
255 self.evict_lowest_quality(&mut data);
256 }
257
258 data.peers.insert(peer.peer_id.0, peer);
259
260 let count = data.peers.len();
261 drop(data);
262
263 let _ = self
264 .event_tx
265 .send(CacheEvent::Updated { peer_count: count });
266 }
267
268 pub async fn add_seed(&self, peer_id: PeerId, addresses: Vec<SocketAddr>) {
270 let peer = CachedPeer::new(peer_id, addresses, PeerSource::Seed);
271 self.upsert(peer).await;
272 }
273
274 pub async fn add_from_connection(
276 &self,
277 peer_id: PeerId,
278 addresses: Vec<SocketAddr>,
279 caps: Option<PeerCapabilities>,
280 ) {
281 let mut peer = CachedPeer::new(peer_id, addresses, PeerSource::Connection);
282 if let Some(caps) = caps {
283 peer.capabilities = caps;
284 }
285 self.upsert(peer).await;
286 }
287
288 pub async fn record_outcome(&self, peer_id: &PeerId, outcome: ConnectionOutcome) {
290 let mut data = self.data.write().await;
291
292 if let Some(peer) = data.peers.get_mut(&peer_id.0) {
293 if outcome.success {
294 peer.record_success(
295 outcome.rtt_ms.unwrap_or(100),
296 outcome.capabilities_discovered,
297 );
298 } else {
299 peer.record_failure();
300 }
301
302 peer.calculate_quality(&self.config.weights);
304 }
305 }
306
307 pub async fn record_success(&self, peer_id: &PeerId, rtt_ms: u32) {
309 self.record_outcome(
310 peer_id,
311 ConnectionOutcome {
312 success: true,
313 rtt_ms: Some(rtt_ms),
314 capabilities_discovered: None,
315 },
316 )
317 .await;
318 }
319
320 pub async fn record_failure(&self, peer_id: &PeerId) {
322 self.record_outcome(
323 peer_id,
324 ConnectionOutcome {
325 success: false,
326 rtt_ms: None,
327 capabilities_discovered: None,
328 },
329 )
330 .await;
331 }
332
333 pub async fn update_capabilities(&self, peer_id: &PeerId, caps: PeerCapabilities) {
335 let mut data = self.data.write().await;
336
337 if let Some(peer) = data.peers.get_mut(&peer_id.0) {
338 peer.capabilities = caps;
339 peer.calculate_quality(&self.config.weights);
340 }
341 }
342
343 pub async fn observe_inbound_peer_address(&self, peer_id: PeerId, address: SocketAddr) {
361 let mut data = self.data.write().await;
362 let now = SystemTime::now();
363 let peer = data
364 .peers
365 .entry(peer_id.0)
366 .or_insert_with(|| CachedPeer::new(peer_id, vec![address], PeerSource::Connection));
367 if !peer.addresses.contains(&address) {
368 peer.addresses.push(address);
369 }
370 peer.last_seen = now;
371 peer.stats.success_count = peer.stats.success_count.saturating_add(1);
372 self.refresh_cached_peer(peer, now);
373 let count = data.peers.len();
374 drop(data);
375 let _ = self
376 .event_tx
377 .send(CacheEvent::Updated { peer_count: count });
378 }
379
380 pub async fn observe_direct_reachability(&self, peer_id: PeerId, address: SocketAddr) {
385 let mut data = self.data.write().await;
386 let now = SystemTime::now();
387
388 let peer = data
389 .peers
390 .entry(peer_id.0)
391 .or_insert_with(|| CachedPeer::new(peer_id, vec![address], PeerSource::Connection));
392
393 if !peer.addresses.contains(&address) {
394 peer.addresses.push(address);
395 }
396
397 peer.last_seen = now;
398 peer.last_attempt = Some(now);
399 peer.stats.success_count = peer.stats.success_count.saturating_add(1);
400 peer.capabilities.record_direct_observation(address, now);
401 self.refresh_cached_peer(peer, now);
402
403 let count = data.peers.len();
404 drop(data);
405
406 let _ = self
407 .event_tx
408 .send(CacheEvent::Updated { peer_count: count });
409 }
410
411 pub async fn get(&self, peer_id: &PeerId) -> Option<CachedPeer> {
413 let mut data = self.data.write().await;
414 let peer = data.peers.get_mut(&peer_id.0)?;
415 self.refresh_cached_peer(peer, SystemTime::now());
416 Some(peer.clone())
417 }
418
419 pub async fn update_token(&self, peer_id: PeerId, token: Vec<u8>) {
421 let mut data = self.data.write().await;
422 if let Some(peer) = data.peers.get_mut(&peer_id.0) {
423 peer.token = Some(token);
424 }
425 }
426
427 pub async fn get_all_tokens(&self) -> std::collections::HashMap<PeerId, Vec<u8>> {
429 self.data
430 .read()
431 .await
432 .peers
433 .values()
434 .filter_map(|p| p.token.clone().map(|t| (p.peer_id, t)))
435 .collect()
436 }
437
438 pub async fn contains(&self, peer_id: &PeerId) -> bool {
440 self.data.read().await.peers.contains_key(&peer_id.0)
441 }
442
443 pub async fn remove(&self, peer_id: &PeerId) -> Option<CachedPeer> {
445 self.data.write().await.peers.remove(&peer_id.0)
446 }
447
448 pub async fn save(&self) -> std::io::Result<()> {
450 let Some(persistence) = &self.persistence else {
451 return Ok(());
452 };
453 let mut data = self.data.write().await;
454
455 if data.peers.len() < self.config.min_peers_to_save {
456 debug!(
457 "Skipping save: only {} peers (min: {})",
458 data.peers.len(),
459 self.config.min_peers_to_save
460 );
461 return Ok(());
462 }
463
464 persistence.save(&mut data)?;
465
466 drop(data);
467 *self.last_save.write().await = Instant::now();
468 let _ = self.event_tx.send(CacheEvent::Saved);
469
470 Ok(())
471 }
472
473 pub async fn cleanup_stale(&self) -> usize {
478 let mut data = self.data.write().await;
479 let initial_count = data.peers.len();
480
481 data.peers
482 .retain(|_, peer| !peer.is_stale(self.config.stale_threshold));
483
484 let removed = initial_count - data.peers.len();
485
486 if removed > 0 {
487 info!("Cleaned up {} stale peers", removed);
488 let _ = self.event_tx.send(CacheEvent::Cleaned { removed });
489 }
490
491 drop(data);
492 *self.last_cleanup.write().await = Instant::now();
493
494 removed
495 }
496
497 pub async fn recalculate_quality(&self) {
499 let mut data = self.data.write().await;
500
501 for peer in data.peers.values_mut() {
502 peer.calculate_quality(&self.config.weights);
503 }
504
505 let count = data.peers.len();
506 let _ = self
507 .event_tx
508 .send(CacheEvent::Updated { peer_count: count });
509 }
510
511 pub async fn stats(&self) -> CacheStats {
513 let mut data = self.data.write().await;
514 let now = SystemTime::now();
515 for peer in data.peers.values_mut() {
516 self.refresh_cached_peer(peer, now);
517 }
518
519 let relay_count = data
520 .peers
521 .values()
522 .filter(|p| p.capabilities.supports_relay)
523 .count();
524 let coord_count = data
525 .peers
526 .values()
527 .filter(|p| p.capabilities.supports_coordination)
528 .count();
529 let dual_stack_count = data
530 .peers
531 .values()
532 .filter(|p| p.capabilities.supports_relay && p.capabilities.supports_dual_stack())
533 .count();
534 let untested = data
535 .peers
536 .values()
537 .filter(|p| p.stats.success_count + p.stats.failure_count == 0)
538 .count();
539 let avg_quality = if data.peers.is_empty() {
540 0.0
541 } else {
542 data.peers.values().map(|p| p.quality_score).sum::<f64>() / data.peers.len() as f64
543 };
544
545 CacheStats {
546 total_peers: data.peers.len(),
547 relay_peers: relay_count,
548 coordinator_peers: coord_count,
549 dual_stack_relay_peers: dual_stack_count,
550 average_quality: avg_quality,
551 untested_peers: untested,
552 }
553 }
554
555 pub fn start_maintenance(self: Arc<Self>) -> tokio::task::JoinHandle<()> {
568 let cache = self;
569
570 tokio::spawn(async move {
571 let mut save_interval = tokio::time::interval(cache.config.save_interval);
572 let mut cleanup_interval = tokio::time::interval(cache.config.cleanup_interval);
573 let mut quality_interval = tokio::time::interval(cache.config.quality_update_interval);
574
575 loop {
576 tokio::select! {
577 _ = save_interval.tick() => {
578 if let Err(e) = cache.save().await {
579 warn!("Failed to save cache: {}", e);
580 }
581 }
582 _ = cleanup_interval.tick() => {
583 cache.cleanup_stale().await;
584 }
585 _ = quality_interval.tick() => {
586 cache.recalculate_quality().await;
587 }
588 }
589 }
590 })
591 }
592
593 pub async fn all_peers(&self) -> Vec<CachedPeer> {
595 let mut data = self.data.write().await;
596 let now = SystemTime::now();
597 for peer in data.peers.values_mut() {
598 self.refresh_cached_peer(peer, now);
599 }
600 data.peers.values().cloned().collect()
601 }
602
603 pub fn config(&self) -> &BootstrapCacheConfig {
605 &self.config
606 }
607
608 fn evict_lowest_quality(&self, data: &mut CacheData) {
609 let evict_count = (self.config.max_peers / 20).max(1); let mut sorted: Vec<_> = data.peers.iter().collect();
612 sorted.sort_by(|a, b| {
613 a.1.quality_score
614 .partial_cmp(&b.1.quality_score)
615 .unwrap_or(std::cmp::Ordering::Equal)
616 });
617
618 let to_remove: Vec<[u8; 32]> = sorted
619 .into_iter()
620 .take(evict_count)
621 .map(|(id, _)| *id)
622 .collect();
623
624 for id in to_remove {
625 data.peers.remove(&id);
626 }
627
628 debug!("Evicted {} lowest quality peers", evict_count);
629 }
630}
631
632#[cfg(test)]
633mod tests {
634 use super::*;
635 use tempfile::TempDir;
636
637 async fn create_test_cache(temp_dir: &TempDir) -> BootstrapCache {
638 let config = BootstrapCacheConfig::builder()
639 .cache_dir(temp_dir.path())
640 .max_peers(100)
641 .epsilon(0.0) .min_peers_to_save(1)
643 .build();
644
645 BootstrapCache::open(config).await.unwrap()
646 }
647
648 #[tokio::test]
654 async fn in_memory_cache_never_touches_disk() {
655 let temp_dir = TempDir::new().unwrap();
656 let cache_dir = temp_dir.path().join("does-not-exist");
657 let config = BootstrapCacheConfig::builder()
658 .cache_dir(&cache_dir)
659 .min_peers_to_save(1)
660 .persist(false)
661 .build();
662
663 let cache = BootstrapCache::open(config).await.unwrap();
664 cache
665 .add_seed(PeerId([7u8; 32]), vec!["127.0.0.1:9000".parse().unwrap()])
666 .await;
667 cache.save().await.unwrap();
668
669 assert!(
670 !cache_dir.exists(),
671 "in-memory cache must not create its cache directory"
672 );
673 assert_eq!(cache.peer_count().await, 1);
675 }
676
677 #[tokio::test]
678 async fn test_cache_creation() {
679 let temp_dir = TempDir::new().unwrap();
680 let cache = create_test_cache(&temp_dir).await;
681 assert_eq!(cache.peer_count().await, 0);
682 }
683
684 #[tokio::test]
685 async fn test_add_and_get() {
686 let temp_dir = TempDir::new().unwrap();
687 let cache = create_test_cache(&temp_dir).await;
688
689 let peer_id = PeerId([1u8; 32]);
690 cache
691 .add_seed(peer_id, vec!["127.0.0.1:9000".parse().unwrap()])
692 .await;
693
694 assert_eq!(cache.peer_count().await, 1);
695 assert!(cache.contains(&peer_id).await);
696
697 let peer = cache.get(&peer_id).await.unwrap();
698 assert_eq!(peer.addresses.len(), 1);
699 }
700
701 #[tokio::test]
702 async fn test_select_peers() {
703 let temp_dir = TempDir::new().unwrap();
704 let cache = create_test_cache(&temp_dir).await;
705
706 for i in 0..10usize {
708 let peer_id = PeerId([i as u8; 32]);
709 let mut peer = CachedPeer::new(
710 peer_id,
711 vec![format!("127.0.0.1:{}", 9000 + i).parse().unwrap()],
712 PeerSource::Seed,
713 );
714 peer.quality_score = i as f64 / 10.0;
715 cache.upsert(peer).await;
716 }
717
718 let selected = cache.select_peers(5).await;
720 assert_eq!(selected.len(), 5);
721 assert!(selected[0].quality_score >= selected[4].quality_score);
722 }
723
724 #[tokio::test]
725 async fn test_persistence() {
726 let temp_dir = TempDir::new().unwrap();
727
728 {
730 let cache = create_test_cache(&temp_dir).await;
731 cache
732 .add_seed(PeerId([1; 32]), vec!["127.0.0.1:9000".parse().unwrap()])
733 .await;
734 cache.save().await.unwrap();
735 }
736
737 {
739 let cache = create_test_cache(&temp_dir).await;
740 assert_eq!(cache.peer_count().await, 1);
741 assert!(cache.contains(&PeerId([1; 32])).await);
742 }
743 }
744
745 #[tokio::test]
746 async fn test_persisted_explicit_assist_hints_survive_reopen() {
747 let temp_dir = TempDir::new().unwrap();
748 let peer_id = PeerId([9; 32]);
749 let peer_addr: SocketAddr = "198.51.100.9:9000".parse().unwrap();
750
751 {
752 let cache = create_test_cache(&temp_dir).await;
753 let mut peer = CachedPeer::new(peer_id, vec![peer_addr], PeerSource::Merge);
754 peer.capabilities.record_assist_hints(true, true);
755 cache.upsert(peer).await;
756 cache.save().await.unwrap();
757 }
758
759 {
760 let cache = create_test_cache(&temp_dir).await;
761 let peer = cache.get(&peer_id).await.expect("peer should reload");
762 assert!(peer.capabilities.hinted_supports_relay);
763 assert!(peer.capabilities.hinted_supports_coordination);
764 assert!(peer.capabilities.supports_relay);
765 assert!(peer.capabilities.supports_coordination);
766 assert!(peer.addresses.contains(&peer_addr));
767 }
768 }
769
770 #[tokio::test]
771 async fn test_quality_scoring() {
772 let temp_dir = TempDir::new().unwrap();
773 let cache = create_test_cache(&temp_dir).await;
774
775 let peer_id = PeerId([1; 32]);
776 cache
777 .add_seed(peer_id, vec!["127.0.0.1:9000".parse().unwrap()])
778 .await;
779
780 let peer = cache.get(&peer_id).await.unwrap();
782 let initial_quality = peer.quality_score;
783
784 for _ in 0..5 {
786 cache.record_success(&peer_id, 50).await;
787 }
788
789 let peer = cache.get(&peer_id).await.unwrap();
790 assert!(peer.quality_score > initial_quality);
791 assert!(peer.success_rate() > 0.9);
792 }
793
794 #[tokio::test]
795 async fn test_eviction() {
796 let temp_dir = TempDir::new().unwrap();
797 let config = BootstrapCacheConfig::builder()
798 .cache_dir(temp_dir.path())
799 .max_peers(10)
800 .build();
801
802 let cache = BootstrapCache::open(config).await.unwrap();
803
804 for i in 0..15u8 {
806 let peer_id = PeerId([i; 32]);
807 let mut peer = CachedPeer::new(
808 peer_id,
809 vec![format!("127.0.0.1:{}", 9000 + i as u16).parse().unwrap()],
810 PeerSource::Seed,
811 );
812 peer.quality_score = i as f64 / 15.0;
813 cache.upsert(peer).await;
814 }
815
816 assert!(cache.peer_count().await <= 10);
818 }
819
820 #[tokio::test]
821 async fn test_stats() {
822 let temp_dir = TempDir::new().unwrap();
823 let cache = create_test_cache(&temp_dir).await;
824
825 let mut peer1 = CachedPeer::new(
827 PeerId([1; 32]),
828 vec!["203.0.113.1:9001".parse().unwrap()],
829 PeerSource::Seed,
830 );
831 peer1
832 .capabilities
833 .record_direct_observation("203.0.113.1:9001".parse().unwrap(), SystemTime::now());
834 cache.upsert(peer1).await;
835
836 let mut peer2 = CachedPeer::new(
837 PeerId([2; 32]),
838 vec!["198.51.100.2:9002".parse().unwrap()],
839 PeerSource::Seed,
840 );
841 peer2
842 .capabilities
843 .record_direct_observation("198.51.100.2:9002".parse().unwrap(), SystemTime::now());
844 cache.upsert(peer2).await;
845
846 cache
847 .add_seed(PeerId([3; 32]), vec!["127.0.0.1:9003".parse().unwrap()])
848 .await;
849
850 let stats = cache.stats().await;
851 assert_eq!(stats.total_peers, 3);
852 assert_eq!(stats.relay_peers, 2);
853 assert_eq!(stats.coordinator_peers, 2);
854 assert_eq!(stats.untested_peers, 3);
855 }
856
857 #[tokio::test]
858 async fn test_select_relay_peers() {
859 let temp_dir = TempDir::new().unwrap();
860 let cache = create_test_cache(&temp_dir).await;
861
862 for i in 0..10u8 {
864 let addr: SocketAddr = format!("127.0.0.1:{}", 9000 + i as u16).parse().unwrap();
865 let mut peer = CachedPeer::new(PeerId([i; 32]), vec![addr], PeerSource::Seed);
866 if i % 2 == 0 {
867 peer.capabilities
868 .record_direct_observation(addr, SystemTime::now());
869 }
870 peer.quality_score = i as f64 / 10.0;
871 cache.upsert(peer).await;
872 }
873
874 let relays = cache.select_relay_peers(10).await;
877 assert_eq!(relays.len(), 10); let relay_capable = relays
881 .iter()
882 .take(5)
883 .filter(|p| p.capabilities.direct_reachability_scope.is_some())
884 .count();
885 assert_eq!(
886 relay_capable, 5,
887 "Scoped direct-evidence peers should be first"
888 );
889 }
890
891 #[tokio::test]
892 async fn test_observe_direct_reachability_preserves_local_scope_without_global_promotion() {
893 let temp_dir = TempDir::new().unwrap();
894 let cache = create_test_cache(&temp_dir).await;
895 let peer_id = PeerId([9; 32]);
896 let addr: SocketAddr = "192.168.1.50:9000".parse().unwrap();
897
898 cache.observe_direct_reachability(peer_id, addr).await;
899
900 let peer = cache.get(&peer_id).await.expect("peer inserted");
901 assert!(!peer.capabilities.supports_relay);
902 assert!(!peer.capabilities.supports_coordination);
903 assert_eq!(
904 peer.capabilities.direct_reachability_scope,
905 Some(crate::reachability::ReachabilityScope::LocalNetwork)
906 );
907 assert!(peer.addresses.contains(&addr));
908 assert!(
909 peer.capabilities
910 .reachable_addresses
911 .iter()
912 .any(|entry| entry.address == addr)
913 );
914 assert!(peer.success_rate() > 0.0);
915 }
916 #[tokio::test]
920 async fn inbound_observation_caches_address_without_capability_grant() {
921 let temp_dir = TempDir::new().expect("tempdir");
922 let cache = create_test_cache(&temp_dir).await;
923 let peer_id = PeerId([41u8; 32]);
924 let addr: SocketAddr = "203.0.113.9:41000".parse().expect("addr");
925 cache.observe_inbound_peer_address(peer_id, addr).await;
926 let peer = cache.get(&peer_id).await.expect("cached");
927 assert!(peer.addresses.contains(&addr), "redial candidate kept");
928 assert!(
929 !peer.capabilities.supports_relay && !peer.capabilities.supports_coordination,
930 "inbound evidence must not grant helper capabilities"
931 );
932 }
933
934 #[tokio::test]
936 async fn outbound_observation_still_grants_capability() {
937 let temp_dir = TempDir::new().expect("tempdir");
938 let cache = create_test_cache(&temp_dir).await;
939 let peer_id = PeerId([42u8; 32]);
940 let addr: SocketAddr = "203.0.113.10:41001".parse().expect("addr");
941 cache.observe_direct_reachability(peer_id, addr).await;
942 let peer = cache.get(&peer_id).await.expect("cached");
943 assert!(
944 peer.capabilities.supports_relay && peer.capabilities.supports_coordination,
945 "outbound global evidence grants helper capabilities"
946 );
947 }
948}