mirror of
https://github.com/jmcorgan/fips.git
synced 2026-07-22 07:48:26 +00:00
Fix bloom filter cascade: gate re-announcements on content change
FilterAnnounce messages never settled in steady state because handle_filter_announce unconditionally marked all peers for re-announcement on every inbound filter, creating a perpetual ping-pong at ~1 message/sec. Added last_sent_filters tracking to BloomState. mark_changed_peers() computes outgoing filters and compares against what was last sent, only marking peers whose filter content actually differs.
This commit is contained in:
41
src/bloom.rs
41
src/bloom.rs
@@ -306,6 +306,8 @@ pub struct BloomState {
|
||||
pending_updates: HashSet<NodeAddr>,
|
||||
/// Current sequence number for outgoing filters.
|
||||
sequence: u64,
|
||||
/// Last outgoing filter sent to each peer (for change detection).
|
||||
last_sent_filters: HashMap<NodeAddr, BloomFilter>,
|
||||
}
|
||||
|
||||
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<NodeAddr, BloomFilter>,
|
||||
) {
|
||||
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:
|
||||
|
||||
@@ -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<NodeAddr> = 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<NodeAddr> = 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).
|
||||
|
||||
@@ -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<NodeAddr> = self.peers.keys().copied().collect();
|
||||
self.bloom_state.mark_all_updates_needed(remaining_peers);
|
||||
|
||||
|
||||
Reference in New Issue
Block a user