Heterogeneous bloom filters: delta compression, variable sizing, adaptive heuristic

Replace fixed 1KB bloom filters with variable-size filters (512B-32KB)
that adapt to each node's position in the spanning tree.

Core changes:
- Internal storage: Vec<u8> → Vec<u64> for word-level operations
- Delta compression: XOR diff with word-level RLE, sequence-based
  NACK protocol for full retransmit recovery
- Size conversion: fold (large→small) and duplicate (small→large)
  with auto-converting merge for mixed-size filter combination
- Native-size storage: peer filters stored at advertised size for
  full-resolution routing queries, converted only for outgoing filter
- Adaptive sizing: outgoing fill ratio drives step-up/step-down
  between size classes with hysteresis (20%/5% thresholds)
- Filter size decoupled from FMP negotiation: announced dynamically
  in filter updates, bit 7 and TLV field 1 removed from handshake

Wire format: FilterAnnounce gains flags byte (delta bit), base_seq
field, RLE-compressed payload. New FilterNack message (0x21) for
out-of-sequence delta recovery.
This commit is contained in:
Johnathan Corgan
2026-04-11 08:16:01 +00:00
parent 8162d3c237
commit eb8479eb20
16 changed files with 1378 additions and 476 deletions
+195 -38
View File
@@ -1,15 +1,16 @@
//! Bloom filter announce send/receive logic.
//!
//! Handles building, sending, and receiving FilterAnnounce messages,
//! including debounced propagation to peers.
//! including delta compression with NACK-based recovery and debounced
//! propagation to peers.
use crate::NodeAddr;
use crate::bloom::BloomFilter;
use crate::protocol::FilterAnnounce;
use crate::protocol::{FilterAnnounce, FilterNack};
use crate::NodeAddr;
use super::{Node, NodeError};
use std::collections::HashMap;
use tracing::debug;
use tracing::{debug, trace, warn};
impl Node {
/// Collect inbound filters from full tree peers for outgoing filter computation.
@@ -33,16 +34,29 @@ impl Node {
/// Build a FilterAnnounce for a specific peer.
///
/// The outgoing filter excludes the destination peer's own filter
/// to prevent routing loops (don't tell a peer about destinations
/// reachable only through them).
/// Returns a delta (XOR diff) if we have a previous filter for this peer
/// at the same size class. Otherwise returns a full send.
fn build_filter_announce(&mut self, exclude_peer: &NodeAddr) -> FilterAnnounce {
let peer_filters = self.peer_inbound_filters();
let filter = self
.bloom_state
.compute_outgoing_filter(exclude_peer, &peer_filters);
let sequence = self.bloom_state.next_sequence();
FilterAnnounce::new(filter, sequence)
let size_class = self.bloom_state.size_class();
// Try delta if we have a previous filter for this peer at the same size
if let Some(last_filter) = self.bloom_state.last_sent_filter(exclude_peer)
&& last_filter.num_bits() == filter.num_bits()
&& let (Some(base_seq), Ok(diff)) = (
self.bloom_state.last_sent_seq(exclude_peer),
last_filter.xor_diff(&filter),
)
{
return FilterAnnounce::delta(diff, sequence, base_seq, size_class);
}
// Full send
FilterAnnounce::full(filter, sequence, size_class)
}
/// Send a FilterAnnounce to a specific peer, respecting debounce.
@@ -67,11 +81,26 @@ 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),
})?;
let is_delta = announce.is_delta;
let sent_filter = if is_delta {
// For deltas, reconstruct the actual filter for change detection:
// apply the diff to the last-sent filter
let mut reconstructed = self
.bloom_state
.last_sent_filter(peer_addr)
.cloned()
.unwrap_or_default();
let _ = reconstructed.apply_diff(&announce.filter);
reconstructed
} else {
announce.filter.clone()
};
let (encoded, stats) =
announce.encode().map_err(|e| NodeError::SendFailed {
node_addr: *peer_addr,
reason: format!("FilterAnnounce encode failed: {}", e),
})?;
// Send
if let Err(e) = self.send_encrypted_link_message(peer_addr, &encoded).await {
@@ -80,19 +109,26 @@ impl Node {
}
self.stats_mut().bloom.sent += 1;
if is_delta {
self.stats_mut().bloom.deltas_sent += 1;
} else {
self.stats_mut().bloom.full_sends += 1;
}
// Record send and store the filter for change detection
debug!(
peer = %self.peer_display_name(peer_addr),
seq = announce.sequence,
delta = is_delta,
compressed = stats.compressed_bytes,
runs = stats.run_count,
est_entries = format_args!("{:.0}", sent_filter.estimated_count()),
set_bits = sent_filter.count_ones(),
fill = format_args!("{:.1}%", sent_filter.fill_ratio() * 100.0),
tree_peer = self.is_tree_peer(peer_addr),
"Sent FilterAnnounce"
);
self.bloom_state.record_update_sent(*peer_addr, now_ms);
self.bloom_state.record_sent_filter(*peer_addr, sent_filter);
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();
}
@@ -134,10 +170,8 @@ impl Node {
/// Handle an inbound FilterAnnounce from an authenticated peer.
///
/// 1. Decode and validate the message
/// 2. Check sequence freshness (reject stale/replay)
/// 3. Store the filter on the peer
/// 4. Mark other peers for outgoing filter update
/// Supports both full sends and delta (XOR diff) updates.
/// On out-of-sequence delta, sends a NACK to request full retransmission.
pub(super) async fn handle_filter_announce(&mut self, from: &NodeAddr, payload: &[u8]) {
self.stats_mut().bloom.received += 1;
@@ -156,26 +190,22 @@ impl Node {
debug!(from = %self.peer_display_name(from), "FilterAnnounce filter/size_class mismatch");
return;
}
if !announce.is_v1_compliant() {
self.stats_mut().bloom.non_v1 += 1;
debug!(from = %self.peer_display_name(from), size_class = announce.size_class, "Non-v1 FilterAnnounce rejected");
return;
}
// Check peer exists
let current_seq = match self.peers.get(from) {
Some(peer) => peer.filter_sequence(),
let peer = match self.peers.get(from) {
Some(p) => p,
None => {
self.stats_mut().bloom.unknown_peer += 1;
debug!(from = %self.peer_display_name(from), "FilterAnnounce from unknown peer");
return;
}
};
let current_seq = peer.filter_sequence();
// Reject stale/replay
if announce.sequence <= current_seq {
self.stats_mut().bloom.stale += 1;
debug!(
trace!(
from = %self.peer_display_name(from),
received_seq = announce.sequence,
current_seq = current_seq,
@@ -184,6 +214,70 @@ impl Node {
return;
}
// Handle delta vs full
let resolved_filter = if announce.is_delta {
// Delta: apply XOR diff to stored inbound filter
let expected_base = current_seq;
if announce.base_seq != expected_base {
// Out-of-sequence delta — send NACK
debug!(
from = %self.peer_display_name(from),
expected_base = expected_base,
got_base = announce.base_seq,
"Out-of-sequence delta, sending NACK"
);
let nack = FilterNack {
expected_seq: expected_base,
};
let nack_encoded = nack.encode();
let _ = self
.send_encrypted_link_message(from, &nack_encoded)
.await;
self.stats_mut().bloom.nacks_sent += 1;
return;
}
// Apply diff to current inbound filter
match self.peers.get(from).and_then(|p| p.inbound_filter()) {
Some(current) => {
let mut result = current.clone();
if let Err(e) = result.apply_diff(&announce.filter) {
warn!(
from = %self.peer_display_name(from),
error = %e,
"Failed to apply filter delta"
);
// Send NACK to request full retransmit
let nack = FilterNack {
expected_seq: current_seq,
};
let _ = self
.send_encrypted_link_message(from, &nack.encode())
.await;
self.stats_mut().bloom.nacks_sent += 1;
return;
}
result
}
None => {
// No stored filter to apply delta to — NACK
debug!(
from = %self.peer_display_name(from),
"Delta received but no stored filter, sending NACK"
);
let nack = FilterNack { expected_seq: 0 };
let _ = self
.send_encrypted_link_message(from, &nack.encode())
.await;
self.stats_mut().bloom.nacks_sent += 1;
return;
}
}
} else {
// Full send: use directly
announce.filter.clone()
};
self.stats_mut().bloom.accepted += 1;
let now_ms = std::time::SystemTime::now()
@@ -194,31 +288,94 @@ impl Node {
debug!(
from = %self.peer_display_name(from),
seq = announce.sequence,
est_entries = format_args!("{:.0}", announce.filter.estimated_count()),
set_bits = announce.filter.count_ones(),
fill = format_args!("{:.1}%", announce.filter.fill_ratio() * 100.0),
tree_peer = self.is_tree_peer(from),
delta = announce.is_delta,
est_entries = format_args!("{:.0}", resolved_filter.estimated_count()),
fill = format_args!("{:.1}%", resolved_filter.fill_ratio() * 100.0),
"Received FilterAnnounce"
);
// Store on peer
// Store resolved filter on peer
if let Some(peer) = self.peers.get_mut(from) {
peer.update_filter(announce.filter, announce.sequence, now_ms);
peer.update_filter(resolved_filter, announce.sequence, now_ms);
}
// Check which peers' outgoing filters actually changed.
// All peers receive filters, but only tree peers' inbound filters
// are merged into outgoing computation (tree-only propagation).
// 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);
}
/// Handle an inbound FilterNack from a peer.
///
/// Clears the last-sent filter for that peer, forcing a full re-send
/// on the next tick.
pub(super) async fn handle_filter_nack(&mut self, from: &NodeAddr, payload: &[u8]) {
let nack = match FilterNack::decode(payload) {
Ok(n) => n,
Err(e) => {
debug!(from = %self.peer_display_name(from), error = %e, "Malformed FilterNack");
return;
}
};
debug!(
from = %self.peer_display_name(from),
expected_seq = nack.expected_seq,
"Received FilterNack, scheduling full re-send"
);
self.stats_mut().bloom.nacks_received += 1;
// Clear sent state for this peer → next send will be full
self.bloom_state.clear_sent_filter(from);
self.bloom_state.mark_update_needed(*from);
}
/// Evaluate adaptive filter sizing and adjust if needed.
///
/// Checks the outgoing fill ratio for a representative peer and
/// steps up or down the size class if thresholds are crossed.
/// On size change, clears all sent filters (forcing full re-sends)
/// and marks all peers for update.
fn check_adaptive_sizing(&mut self) {
// Only Full nodes participate in filter sizing
if self.node_profile != crate::protocol::NodeProfile::Full {
return;
}
// Use an arbitrary peer to compute a representative outgoing filter
let representative_peer = match self.peers.keys().next() {
Some(addr) => *addr,
None => return,
};
let peer_filters = self.peer_inbound_filters();
let outgoing = self
.bloom_state
.compute_outgoing_filter(&representative_peer, &peer_filters);
let fill = outgoing.fill_ratio();
if let Some(new_class) = self.bloom_state.evaluate_size_change(fill) {
let old_class = self.bloom_state.size_class();
debug!(
old_class = old_class,
new_class = new_class,
fill = format_args!("{:.1}%", fill * 100.0),
"Adaptive bloom filter resize"
);
self.bloom_state.set_size_class(new_class);
self.bloom_state.clear_all_sent_filters();
self.stats_mut().bloom.size_changes += 1;
let all_peers: Vec<NodeAddr> = self.peers.keys().copied().collect();
self.bloom_state.mark_all_updates_needed(all_peers);
}
}
/// Check bloom filter state on tick (called from event loop).
///
/// Sends any pending debounced filter announces.
/// Evaluates adaptive sizing, then sends any pending filter announces.
pub(super) async fn check_bloom_state(&mut self) {
self.check_adaptive_sizing();
self.send_pending_filter_announces().await;
}
}
+4
View File
@@ -42,6 +42,10 @@ impl Node {
// FilterAnnounce
self.handle_filter_announce(from, payload).await;
}
0x21 => {
// FilterNack
self.handle_filter_nack(from, payload).await;
}
0x30 => {
// LookupRequest
self.handle_lookup_request(from, payload).await;
+3 -12
View File
@@ -1159,8 +1159,6 @@ impl Node {
let remote_epoch = connection.remote_epoch();
let peer_profile = connection.peer_profile()
.unwrap_or(crate::protocol::NodeProfile::Full);
let agreed_bloom_size_class = connection.agreed_bloom_size_class()
.unwrap_or(crate::bloom::V1_SIZE_CLASS);
let peer_node_addr = *verified_identity.node_addr();
let is_outbound = connection.is_outbound();
@@ -1205,7 +1203,6 @@ impl Node {
remote_epoch,
self.node_profile,
peer_profile,
agreed_bloom_size_class,
);
new_peer.set_tree_announce_min_interval_ms(self.config.node.tree.announce_min_interval_ms);
@@ -1303,7 +1300,6 @@ impl Node {
remote_epoch,
self.node_profile,
peer_profile,
agreed_bloom_size_class,
);
new_peer.set_tree_announce_min_interval_ms(self.config.node.tree.announce_min_interval_ms);
if let Some(ts) = old_announce_ts {
@@ -1338,30 +1334,25 @@ impl Node {
/// Process an FMP negotiation payload received from a peer.
///
/// Decodes the payload, validates profile pairing, agrees on bloom
/// filter size, and stores the results on the PeerConnection.
/// Decodes the payload, validates profile pairing, and stores the
/// results on the PeerConnection.
fn process_fmp_negotiation(
our_profile: crate::protocol::NodeProfile,
conn: &mut PeerConnection,
neg_bytes: &[u8],
) -> Result<(), crate::protocol::ProtocolError> {
let our_payload = NegotiationPayload::fmp(1, 1, our_profile);
let their_payload = NegotiationPayload::decode(neg_bytes)?;
// Validate profile pairing (at least one Full)
let their_profile = their_payload.node_profile()?;
NegotiationPayload::validate_profiles(our_profile, their_profile)?;
// Agree on bloom filter size
let agreed_bloom = our_payload.agree_bloom_size(&their_payload)?;
conn.set_negotiation_results(their_profile, agreed_bloom);
conn.set_negotiation_results(their_profile);
debug!(
link_id = %conn.link_id(),
our_profile = ?our_profile,
peer_profile = ?their_profile,
agreed_bloom_size_class = agreed_bloom,
"FMP negotiation complete"
);
+17 -3
View File
@@ -208,7 +208,6 @@ pub struct BloomStats {
pub received: u64,
pub decode_error: u64,
pub invalid: u64,
pub non_v1: u64,
pub unknown_peer: u64,
pub stale: u64,
pub accepted: u64,
@@ -216,6 +215,13 @@ pub struct BloomStats {
pub sent: u64,
pub debounce_suppressed: u64,
pub send_failed: u64,
// Delta compression
pub deltas_sent: u64,
pub full_sends: u64,
pub nacks_sent: u64,
pub nacks_received: u64,
// Adaptive sizing
pub size_changes: u64,
}
impl BloomStats {
@@ -224,13 +230,17 @@ impl BloomStats {
received: self.received,
decode_error: self.decode_error,
invalid: self.invalid,
non_v1: self.non_v1,
unknown_peer: self.unknown_peer,
stale: self.stale,
accepted: self.accepted,
sent: self.sent,
debounce_suppressed: self.debounce_suppressed,
send_failed: self.send_failed,
deltas_sent: self.deltas_sent,
full_sends: self.full_sends,
nacks_sent: self.nacks_sent,
nacks_received: self.nacks_received,
size_changes: self.size_changes,
}
}
}
@@ -394,13 +404,17 @@ pub struct BloomStatsSnapshot {
pub received: u64,
pub decode_error: u64,
pub invalid: u64,
pub non_v1: u64,
pub unknown_peer: u64,
pub stale: u64,
pub accepted: u64,
pub sent: u64,
pub debounce_suppressed: u64,
pub send_failed: u64,
pub deltas_sent: u64,
pub full_sends: u64,
pub nacks_sent: u64,
pub nacks_received: u64,
pub size_changes: u64,
}
#[derive(Clone, Debug, Default, Serialize)]