diff --git a/src/bloom.rs b/src/bloom.rs index 8244ee6..11cfcb8 100644 --- a/src/bloom.rs +++ b/src/bloom.rs @@ -306,6 +306,8 @@ pub struct BloomState { pending_updates: HashSet, /// Current sequence number for outgoing filters. sequence: u64, + /// Last outgoing filter sent to each peer (for change detection). + last_sent_filters: HashMap, } impl BloomState { @@ -319,6 +321,7 @@ impl BloomState { last_update_sent: HashMap::new(), pending_updates: HashSet::new(), sequence: 0, + last_sent_filters: HashMap::new(), } } @@ -418,6 +421,44 @@ impl BloomState { self.pending_updates.clear(); } + /// Record the outgoing filter that was sent to a peer. + pub fn record_sent_filter(&mut self, peer_id: NodeAddr, filter: BloomFilter) { + self.last_sent_filters.insert(peer_id, filter); + } + + /// Remove stored filter state for a peer that was removed. + pub fn remove_peer_state(&mut self, peer_id: &NodeAddr) { + self.last_sent_filters.remove(peer_id); + self.last_update_sent.remove(peer_id); + self.pending_updates.remove(peer_id); + } + + /// Mark only peers whose outgoing filter has actually changed. + /// + /// Computes the outgoing filter for each peer and compares it + /// against what was last sent. Only marks peers where the filter + /// differs. This prevents cascading update loops in steady state. + pub fn mark_changed_peers( + &mut self, + exclude_from: &NodeAddr, + peer_addrs: &[NodeAddr], + peer_filters: &HashMap, + ) { + for peer_addr in peer_addrs { + if peer_addr == exclude_from { + continue; + } + let new_filter = self.compute_outgoing_filter(peer_addr, peer_filters); + let changed = match self.last_sent_filters.get(peer_addr) { + Some(last) => *last != new_filter, + None => true, // never sent → must send + }; + if changed { + self.pending_updates.insert(*peer_addr); + } + } + } + /// Compute the outgoing filter for a specific peer. /// /// The filter includes: diff --git a/src/node/bloom.rs b/src/node/bloom.rs index f3e24c0..3c6e154 100644 --- a/src/node/bloom.rs +++ b/src/node/bloom.rs @@ -61,6 +61,7 @@ impl Node { // Build and encode let announce = self.build_filter_announce(peer_addr); + let sent_filter = announce.filter.clone(); let encoded = announce.encode().map_err(|e| NodeError::SendFailed { node_addr: *peer_addr, reason: format!("FilterAnnounce encode failed: {}", e), @@ -69,8 +70,9 @@ impl Node { // Send self.send_encrypted_link_message(peer_addr, &encoded).await?; - // Record send + // Record send and store the filter for change detection self.bloom_state.record_update_sent(*peer_addr, now_ms); + self.bloom_state.record_sent_filter(*peer_addr, sent_filter); if let Some(peer) = self.peers.get_mut(peer_addr) { peer.clear_filter_update_needed(); } @@ -165,14 +167,11 @@ impl Node { "Received FilterAnnounce" ); - // Our outgoing filter changed — mark all other peers for update - let other_peers: Vec = self - .peers - .keys() - .filter(|addr| *addr != from) - .copied() - .collect(); - self.bloom_state.mark_all_updates_needed(other_peers); + // Check which peers' outgoing filters actually changed + let peer_addrs: Vec = self.peers.keys().copied().collect(); + let peer_filters = self.peer_inbound_filters(); + self.bloom_state + .mark_changed_peers(from, &peer_addrs, &peer_filters); } /// Check bloom filter state on tick (called from event loop). diff --git a/src/node/handlers/dispatch.rs b/src/node/handlers/dispatch.rs index a1397bf..819be30 100644 --- a/src/node/handlers/dispatch.rs +++ b/src/node/handlers/dispatch.rs @@ -107,7 +107,8 @@ impl Node { } } - // Bloom filter cleanup: our outgoing filter changed (lost a peer's filter) + // Bloom filter cleanup: clear state for removed peer, mark remaining + self.bloom_state.remove_peer_state(node_addr); let remaining_peers: Vec = self.peers.keys().copied().collect(); self.bloom_state.mark_all_updates_needed(remaining_peers);