Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
5 changes: 5 additions & 0 deletions protocols/gossipsub/CHANGELOG.md
Original file line number Diff line number Diff line change
@@ -1,4 +1,9 @@
## 0.50.0
- Add the gossipsub v1.4 large message extension wire types and the `largeMessageHandling`
capability advertisement behind the default-off `large_message_handling` config option
([issue 6597](https://github.com/libp2p/rust-libp2p/issues/6597)).
See [PR 6599](https://github.com/libp2p/rust-libp2p/pull/6599).

- Change default `TopicSubscriptionFilter` from `AllowAllSubscriptionFilter` to `MaxCountSubscriptionFilter<AllowAllSubscriptionFilter>`
with default limits of `100` for both `max_subscribed_topics` and `max_subscriptions_per_request`,
providing built-in protection against excessive subscription requests.
Expand Down
35 changes: 35 additions & 0 deletions protocols/gossipsub/src/behaviour.rs
Original file line number Diff line number Diff line change
Expand Up @@ -3439,6 +3439,7 @@ where

#[cfg(not(feature = "partial-messages"))]
partial_messages: None,
large_message_handling: self.config.large_message_handling().then_some(true),
}),
);
}
Expand Down Expand Up @@ -3478,6 +3479,7 @@ where
partial_messages: Some(true),
#[cfg(not(feature = "partial-messages"))]
partial_messages: None,
large_message_handling: self.config.large_message_handling().then_some(true),
}),
);
}
Expand Down Expand Up @@ -3639,6 +3641,24 @@ where
self.handle_extensions(&propagation_source, extensions);
}
}
ControlAction::Preamble(preamble) => {
// Large message handling is not yet implemented.
tracing::trace!(
peer=%propagation_source,
message=%preamble.message_id,
message_size=%preamble.message_size,
topic=%preamble.topic_hash,
"Ignoring PREAMBLE control message"
);
}
ControlAction::ImReceiving(imreceiving) => {
// Large message handling is not yet implemented.
tracing::trace!(
peer=%propagation_source,
message=%imreceiving.message_id,
"Ignoring IMRECEIVING control message"
);
}
}
}
if !ihave_msgs.is_empty() {
Expand All @@ -3651,6 +3671,21 @@ where
self.handle_prune(&propagation_source, prune_msgs);
}

// Large message handling is not yet implemented.
rpc.large_message_fragments
.into_iter()
.for_each(|fragment| {
tracing::trace!(
peer=%propagation_source,
message=%fragment.message_id,
fragment_index=%fragment.fragment_index,
total_fragments=%fragment.total_fragments,
fragment_len=%fragment.fragment_data.len(),
topic=%fragment.topic_hash,
"Ignoring large message fragment"
);
});

#[cfg(feature = "partial-messages")]
if let Some(partial_message) = rpc.partial_message {
if self
Expand Down
213 changes: 213 additions & 0 deletions protocols/gossipsub/src/behaviour/tests/extensions.rs
Original file line number Diff line number Diff line change
@@ -0,0 +1,213 @@
//! Tests for the extensions advertisement and the gossipsub v1.4 Large Message
//! Handling capability.

use std::collections::HashMap;

use asynchronous_codec::{Decoder, Encoder};
use bytes::BytesMut;
use libp2p_core::{Multiaddr, PeerId};
use libp2p_swarm::{ConnectionId, NetworkBehaviour};

use super::DefaultBehaviourTestBuilder;
use crate::{
IdentTopic as Topic, ValidationMode,
config::ConfigBuilder,
handler::HandlerEvent,
protocol::GossipsubCodec,
queue::Queue,
rpc_proto::proto,
types::{
ControlAction, Extensions, ImReceiving, LargeMessageFragment, MessageId, Preamble, RpcIn,
RpcOut,
},
};

/// Pops messages from a peer's queue until it finds the extensions
/// advertisement sent on connect.
fn advertised_extensions(queue: &mut Queue) -> Extensions {
std::iter::from_fn(|| queue.try_pop())
.find_map(|rpc| {
if let RpcOut::Extensions(extensions) = rpc {
Some(extensions)
} else {
None
}
})
.expect("Extensions message should be sent on connect")
}

/// Verifies that a peer advertising `largeMessageHandling` is tracked as
/// supporting it, independently of any topic subscriptions.
#[test]
fn test_peer_advertised_extensions_are_tracked() {
let (mut gs, peers, _, _) = DefaultBehaviourTestBuilder::default()
.peer_no(1)
.create_network();
let peer_id = peers[0];
assert_eq!(gs.connected_peers.get(&peer_id).unwrap().extensions, None);

let extensions = Extensions {
partial_messages: None,
large_message_handling: Some(true),
};
gs.on_connection_handler_event(
peer_id,
ConnectionId::new_unchecked(0),
HandlerEvent::Message {
rpc: RpcIn {
messages: vec![],
subscriptions: vec![],
control_msgs: vec![ControlAction::Extensions(Some(extensions))],
large_message_fragments: vec![],
#[cfg(feature = "partial-messages")]
partial_message: None,
},
invalid_messages: vec![],
},
);

assert_eq!(
gs.connected_peers.get(&peer_id).unwrap().extensions,
Some(extensions)
);
}

/// Verifies that with the config option unset the extensions advertisement
/// does not include the `largeMessageHandling` flag.
#[test]
fn test_large_message_handling_not_advertised_by_default() {
let (mut gs, _, _, _) = DefaultBehaviourTestBuilder::default().create_network();
let peer_id = PeerId::random();
gs.handle_established_inbound_connection(
ConnectionId::new_unchecked(0),
peer_id,
&Multiaddr::empty(),
&Multiaddr::empty(),
)
.unwrap();

let mut queue = gs.connected_peers.get(&peer_id).unwrap().messages.clone();
let extensions = advertised_extensions(&mut queue);
assert_eq!(extensions.large_message_handling, None);
}

/// Verifies that with the config option set the extensions advertisement
/// includes `largeMessageHandling`.
#[test]
fn test_large_message_handling_advertised_when_enabled() {
let config = ConfigBuilder::default()
.large_message_handling(true)
.build()
.unwrap();
let (mut gs, _, _, _) = DefaultBehaviourTestBuilder::default()
.gs_config(config)
.create_network();
let peer_id = PeerId::random();
gs.handle_established_inbound_connection(
ConnectionId::new_unchecked(0),
peer_id,
&Multiaddr::empty(),
&Multiaddr::empty(),
)
.unwrap();

let mut queue = gs.connected_peers.get(&peer_id).unwrap().messages.clone();
let extensions = advertised_extensions(&mut queue);
assert_eq!(extensions.large_message_handling, Some(true));
}

/// Verifies that an RPC carrying PREAMBLE, IMRECEIVING and large message
/// fragment entries decodes without error and produces no events.
#[test]
fn test_large_message_rpc_decodes_and_produces_no_events() {
let message_id = MessageId::new(&[1, 2, 3, 4]);
let topic_hash = Topic::new("large-message-topic").hash();

let rpc = proto::Rpc {
publish: vec![],
subscriptions: vec![],
control: Some(proto::ControlMessage {
ihave: vec![],
iwant: vec![],
graft: vec![],
prune: vec![],
idontwant: vec![],
extensions: None,
preamble: vec![proto::ControlPreamble {
message_id: Some(message_id.0.clone()),
message_size: Some(1 << 20),
topic_id: Some(topic_hash.clone().into_string()),
}],
imreceiving: vec![proto::ControlImReceiving {
message_id: Some(message_id.0.clone()),
}],
}),
partial: None,
large_message_fragments: vec![proto::LargeMessageFragment {
message_id: Some(message_id.0.clone()),
fragment_index: Some(0),
total_fragments: Some(4),
fragment_data: Some(vec![7u8; 128]),
topic_id: Some(topic_hash.clone().into_string()),
}],
};

let mut codec = GossipsubCodec::new(
u32::MAX as usize,
ValidationMode::Strict,
HashMap::new(),
5000,
5000,
);
let mut buf = BytesMut::new();
codec.encode(rpc, &mut buf).unwrap();
let event = codec.decode(&mut buf).unwrap().unwrap();

let HandlerEvent::Message {
rpc,
invalid_messages,
} = event
else {
panic!("Expected message event");
};
assert!(invalid_messages.is_empty());
assert!(
rpc.control_msgs
.contains(&ControlAction::Preamble(Preamble {
message_id: message_id.clone(),
message_size: 1 << 20,
topic_hash: topic_hash.clone(),
}))
);
assert!(
rpc.control_msgs
.contains(&ControlAction::ImReceiving(ImReceiving {
message_id: message_id.clone(),
}))
);
assert_eq!(
rpc.large_message_fragments,
vec![LargeMessageFragment {
message_id,
fragment_index: 0,
total_fragments: 4,
fragment_data: vec![7u8; 128],
topic_hash,
}]
);

// Delivering the decoded RPC to the behaviour is a no-op for now.
let (mut gs, peers, _, _) = DefaultBehaviourTestBuilder::default()
.peer_no(1)
.create_network();
gs.events.clear();
gs.on_connection_handler_event(
peers[0],
ConnectionId::new_unchecked(0),
HandlerEvent::Message {
rpc,
invalid_messages: vec![],
},
);
assert!(gs.events.is_empty());
}
1 change: 1 addition & 0 deletions protocols/gossipsub/src/behaviour/tests/gossip.rs
Original file line number Diff line number Diff line change
Expand Up @@ -207,6 +207,7 @@ fn test_handle_iwant_msg_but_already_sent_idontwant() {
let rpc = RpcIn {
messages: vec![],
subscriptions: vec![],
large_message_fragments: vec![],
#[cfg(feature = "partial-messages")]
partial_message: None,
control_msgs: vec![ControlAction::IDontWant(IDontWant {
Expand Down
1 change: 1 addition & 0 deletions protocols/gossipsub/src/behaviour/tests/idontwant.rs
Original file line number Diff line number Diff line change
Expand Up @@ -199,6 +199,7 @@ fn parses_idontwant() {
let rpc = RpcIn {
messages: vec![],
subscriptions: vec![],
large_message_fragments: vec![],
#[cfg(feature = "partial-messages")]
partial_message: None,
control_msgs: vec![ControlAction::IDontWant(IDontWant {
Expand Down
2 changes: 2 additions & 0 deletions protocols/gossipsub/src/behaviour/tests/mod.rs
Original file line number Diff line number Diff line change
Expand Up @@ -31,6 +31,7 @@
//! - Event helpers: [`count_control_msgs`], [`flush_events`]

mod explicit_peers;
mod extensions;
mod floodsub;
mod gossip;
mod graft_prune;
Expand Down Expand Up @@ -543,6 +544,7 @@ pub(super) fn proto_to_message(rpc: &proto::Rpc) -> RpcIn {
})
.collect(),
control_msgs,
large_message_fragments: vec![],
#[cfg(feature = "partial-messages")]
partial_message: None,
}
Expand Down
2 changes: 2 additions & 0 deletions protocols/gossipsub/src/behaviour/tests/partial.rs
Original file line number Diff line number Diff line change
Expand Up @@ -1632,6 +1632,7 @@ fn test_partial_messages_two_node_exchange() {
messages: vec![],
subscriptions: vec![],
control_msgs: vec![],
large_message_fragments: vec![],
partial_message: Some(node2_partial),
},
invalid_messages: vec![],
Expand Down Expand Up @@ -1908,6 +1909,7 @@ fn test_partial_messages_no_redundant_response_on_receive() {
messages: vec![],
subscriptions: vec![],
control_msgs: vec![],
large_message_fragments: vec![],
partial_message: Some(peer_partial),
},
invalid_messages: vec![],
Expand Down
3 changes: 3 additions & 0 deletions protocols/gossipsub/src/behaviour/tests/scoring.rs
Original file line number Diff line number Diff line change
Expand Up @@ -680,6 +680,7 @@ fn test_ignore_rpc_from_peers_below_graylist_threshold() {
messages: vec![raw_message1],
subscriptions: vec![subscription.clone()],
control_msgs: vec![control_action],
large_message_fragments: vec![],
#[cfg(feature = "partial-messages")]
partial_message: None,
},
Expand Down Expand Up @@ -708,6 +709,7 @@ fn test_ignore_rpc_from_peers_below_graylist_threshold() {
messages: vec![raw_message3],
subscriptions: vec![subscription],
control_msgs: vec![control_action],
large_message_fragments: vec![],
#[cfg(feature = "partial-messages")]
partial_message: None,
},
Expand Down Expand Up @@ -1306,6 +1308,7 @@ fn test_scoring_p4_invalid_signature() {
messages: vec![],
subscriptions: vec![],
control_msgs: vec![],
large_message_fragments: vec![],
#[cfg(feature = "partial-messages")]
partial_message: None,
},
Expand Down
4 changes: 4 additions & 0 deletions protocols/gossipsub/src/behaviour/tests/topic_config.rs
Original file line number Diff line number Diff line change
Expand Up @@ -636,6 +636,7 @@ fn test_validation_error_message_size_too_large_topic_specific() {
messages: vec![raw_message],
subscriptions: vec![],
control_msgs: vec![],
large_message_fragments: vec![],
#[cfg(feature = "partial-messages")]
partial_message: None,
},
Expand Down Expand Up @@ -685,6 +686,7 @@ fn test_validation_error_message_size_too_large_topic_specific() {
subscriptions: vec![],
control: None,
partial: None,
large_message_fragments: vec![],
};
codec.encode(rpc, &mut buf).unwrap();

Expand Down Expand Up @@ -745,6 +747,7 @@ fn test_validation_message_size_within_topic_specific() {
messages: vec![raw_message],
subscriptions: vec![],
control_msgs: vec![],
large_message_fragments: vec![],
#[cfg(feature = "partial-messages")]
partial_message: None,
},
Expand Down Expand Up @@ -794,6 +797,7 @@ fn test_validation_message_size_within_topic_specific() {
subscriptions: vec![],
control: None,
partial: None,
large_message_fragments: vec![],
};
codec.encode(rpc, &mut buf).unwrap();

Expand Down
Loading