Skip to main content

ant_quic/bootstrap_cache/
cache.rs

1// Copyright 2024 Saorsa Labs Ltd.
2//
3// This Saorsa Network Software is licensed under the General Public License (GPL), version 3.
4// Please see the file LICENSE-GPL, or visit <http://www.gnu.org/licenses/> for the full text.
5//
6// Full details available at https://saorsalabs.com/licenses
7
8//! Main bootstrap cache implementation.
9
10use 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/// Bootstrap cache event for notifications
22#[derive(Debug, Clone)]
23pub enum CacheEvent {
24    /// Cache was updated (peers added/removed/modified)
25    Updated {
26        /// Current peer count
27        peer_count: usize,
28    },
29    /// Cache was saved to disk
30    Saved,
31    /// Cache was merged from another source
32    Merged {
33        /// Number of peers added from merge
34        added: usize,
35    },
36    /// Stale peers were cleaned up
37    Cleaned {
38        /// Number of peers removed
39        removed: usize,
40    },
41}
42
43/// Cache statistics
44#[derive(Debug, Clone, Default)]
45pub struct CacheStats {
46    /// Total number of cached peers
47    pub total_peers: usize,
48    /// Peers that support relay
49    pub relay_peers: usize,
50    /// Peers that support NAT coordination
51    pub coordinator_peers: usize,
52    /// Peers that support dual-stack (IPv4 + IPv6) bridging
53    pub dual_stack_relay_peers: usize,
54    /// Average quality score across all peers
55    pub average_quality: f64,
56    /// Number of untested peers
57    pub untested_peers: usize,
58}
59
60/// Greedy bootstrap cache with quality-based peer selection.
61///
62/// This cache stores peer information with quality metrics and provides
63/// epsilon-greedy selection to balance exploitation (using known-good peers)
64/// with exploration (trying new peers to discover potentially better ones).
65#[derive(Debug)]
66pub struct BootstrapCache {
67    config: BootstrapCacheConfig,
68    data: Arc<RwLock<CacheData>>,
69    /// `None` when the cache is in-memory only (`config.persist == false`).
70    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    // ... (existing open/subscribe methods)
78    /// Open or create a bootstrap cache.
79    ///
80    /// Loads existing cache data from disk if available, otherwise starts fresh.
81    /// With `config.persist == false` the cache is in-memory only: nothing is
82    /// loaded, nothing is ever written, and no cache directory is created.
83    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    /// Subscribe to cache events
116    pub fn subscribe(&self) -> broadcast::Receiver<CacheEvent> {
117        self.event_tx.subscribe()
118    }
119
120    /// Get the number of cached peers
121    pub async fn peer_count(&self) -> usize {
122        self.data.read().await.peers.len()
123    }
124
125    /// Get a specific peer from the cache
126    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    /// Select peers for bootstrap using epsilon-greedy strategy.
142    ///
143    /// Returns up to `count` peers, balancing exploitation of known-good peers
144    /// with exploration of untested peers based on the configured epsilon.
145    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    /// Select peers that support relay functionality.
161    ///
162    /// Returns peers sorted by quality score, preferring observed relay capability.
163    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    /// Select peers that support NAT coordination.
179    ///
180    /// Returns peers sorted by quality score, preferring observed coordination capability.
181    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    /// Select relay peers that can reach a target IP version.
197    ///
198    /// Returns relays sorted by quality that can bridge traffic to the target.
199    /// Dual-stack relays are preferred as they can reach any target.
200    ///
201    /// # Arguments
202    /// * `count` - Maximum number of relays to return
203    /// * `target` - The target address to reach
204    /// * `prefer_dual_stack` - If true, prioritize dual-stack relays
205    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    /// Select relay peers that support dual-stack (IPv4 + IPv6) bridging.
228    ///
229    /// These peers are valuable for bridging between IPv4-only and IPv6-only networks.
230    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    /// Add or update a peer in the cache.
248    ///
249    /// If the cache is at capacity, evicts the lowest quality peers.
250    pub async fn upsert(&self, peer: CachedPeer) {
251        let mut data = self.data.write().await;
252
253        // Evict lowest quality if at capacity
254        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    /// Add a seed peer (user-provided bootstrap node).
269    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    /// Add a peer discovered from an active connection.
275    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    /// Record a connection attempt result.
289    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            // Recalculate quality score
303            peer.calculate_quality(&self.config.weights);
304        }
305    }
306
307    /// Record successful connection.
308    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    /// Record failed connection.
321    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    /// Update peer capabilities.
334    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    /// Record that a peer was directly reachable from this node.
344    ///
345    /// This is observer-scoped evidence. A peer is considered suitable for
346    /// relay/bootstrap/coordinator selection only after a direct connection to
347    /// one of its addresses succeeds without coordinator or relay assistance.
348    /// Record an inbound peer's source address as a redial candidate WITHOUT
349    /// granting reachability-derived capabilities (x0x#398 / #262 fix 4).
350    ///
351    /// ant-quic uses one socket for inbound and outbound, so an inbound
352    /// peer's observed source port is its listening port and is a sound
353    /// redial candidate. It is NOT evidence the peer is globally reachable:
354    /// a NATed desktop dialling out through a cone NAT shows a global source
355    /// address that only *we* (or nobody, for symmetric NATs) can reach.
356    /// Granting `supports_relay`/`supports_coordination` from inbound
357    /// evidence made every desktop that ever dialled a bootstrap advertise
358    /// as a relay/coordinator, poisoning helper selection for third-party
359    /// hole punches.
360    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    /// Record OUTBOUND-verified direct reachability: we dialled `address` and
381    /// the connection succeeded, so the peer is genuinely reachable there.
382    /// This is the only path that feeds capability derivation
383    /// (`supports_relay` / `supports_coordination`).
384    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    /// Get a specific peer.
412    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    /// Update the address validation token for a peer
420    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    /// Get all tokens from cached peers (for initializing TokenStore)
428    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    /// Check if peer exists in cache.
439    pub async fn contains(&self, peer_id: &PeerId) -> bool {
440        self.data.read().await.peers.contains_key(&peer_id.0)
441    }
442
443    /// Remove a peer from cache.
444    pub async fn remove(&self, peer_id: &PeerId) -> Option<CachedPeer> {
445        self.data.write().await.peers.remove(&peer_id.0)
446    }
447
448    /// Save cache to disk. No-op for in-memory caches (`persist: false`).
449    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    /// Cleanup stale peers.
474    ///
475    /// Removes peers that haven't been seen within the stale threshold.
476    /// Returns the number of peers removed.
477    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    /// Recalculate quality scores for all peers.
498    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    /// Get cache statistics.
512    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    /// Start background maintenance tasks.
556    ///
557    /// Spawns a task that periodically:
558    /// - Saves the cache to disk (no-op for in-memory caches)
559    /// - Cleans up stale peers
560    /// - Recalculates quality scores
561    ///
562    /// Returns a handle that can be used to cancel the task.
563    ///
564    /// Note: [`crate::p2p_endpoint::P2pEndpoint`] starts maintenance
565    /// automatically for its own cache — embedders sharing the endpoint's
566    /// cache (via [`crate::Node::bootstrap_cache`]) must not call this again.
567    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    /// Get all cached peers (for export/debug).
594    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    /// Get the configuration.
604    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); // Evict ~5%
610
611        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) // Pure exploitation for predictable tests
642            .min_peers_to_save(1)
643            .build();
644
645        BootstrapCache::open(config).await.unwrap()
646    }
647
648    /// `persist: false` must make the cache purely in-memory: opening in a
649    /// directory that does not exist succeeds, and adding peers + saving
650    /// never creates the directory, cache file, or lock file. This is what
651    /// lets multiple node instances on one host opt out of the shared
652    /// default cache dir without cross-instance pollution.
653    #[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        // Runtime behaviour is unchanged: the peer is still selectable.
674        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        // Add peers with different quality
707        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        // Select should return highest quality first (epsilon=0)
719        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        // Create and populate cache
729        {
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        // Reopen and verify
738        {
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        // Initial quality should be neutral
781        let peer = cache.get(&peer_id).await.unwrap();
782        let initial_quality = peer.quality_score;
783
784        // Record successes - quality should improve
785        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        // Add 15 peers
805        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        // Should have evicted some
817        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        // Add some peers with capabilities
826        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        // Add mix of relay and non-relay peers
863        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        // v0.13.0+: Measure, don't trust - returns all peers but prefers
875        // those with observed relay capability.
876        let relays = cache.select_relay_peers(10).await;
877        assert_eq!(relays.len(), 10); // All peers are candidates
878
879        // First 5 should have direct reachability evidence (prioritized)
880        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    /// #262 fix 4: an INBOUND peer's source address is a redial candidate but
917    /// must not grant relay/coordinator capability — a NATed desktop dialling
918    /// out shows a global source addr that is not proof of reachability.
919    #[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    /// #262 fix 4 counterpart: OUTBOUND-verified reachability still grants.
935    #[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}