1mod affinity;
107
108use crate::config::*;
109
110use affinity::AffinityFilter;
111use codec::{Compact, Decode, Encode, MaxEncodedLen};
112use futures::{
113 channel::oneshot,
114 future::{pending, FusedFuture},
115 prelude::*,
116 stream::FuturesUnordered,
117};
118use governor::{
119 clock::{Clock, DefaultClock},
120 middleware::NoOpMiddleware,
121 state::{InMemoryState, NotKeyed},
122 Quota, RateLimiter,
123};
124use prometheus_endpoint::{
125 exponential_buckets, register, Counter, CounterVec, Gauge, GaugeVec, Histogram, HistogramOpts,
126 HistogramVec, Opts, PrometheusError, Registry, U64,
127};
128use rand::seq::IteratorRandom;
129use sc_network::{
130 config::{NonReservedPeerMode, SetConfig},
131 error, multiaddr,
132 peer_store::PeerStoreProvider,
133 service::{
134 traits::{NotificationEvent, NotificationService, ValidationResult},
135 NotificationMetrics,
136 },
137 types::ProtocolName,
138 utils::interval,
139 NetworkBackend, NetworkEventStream, NetworkPeers,
140};
141use sc_network_sync::{SyncEvent, SyncEventStream};
142use sc_network_types::PeerId;
143use sp_runtime::{
144 traits::{Block as BlockT, ConstU32},
145 BoundedVec,
146};
147use sp_statement_store::{
148 AdmittedBatch, FilterDecision, Hash, Statement, StatementSource, StatementStore, SubmitResult,
149};
150use std::{
151 collections::{hash_map::Entry, HashMap, HashSet, VecDeque},
152 fmt, iter,
153 num::NonZeroU32,
154 pin::Pin,
155 sync::Arc,
156 time::Instant,
157};
158use tokio::time::timeout;
159pub mod config;
160
161pub type Statements = Vec<Statement>;
163
164type StatementBatch = BoundedVec<Statement, ConstU32<{ MAX_STATEMENTS_PER_NOTIFICATION as u32 }>>;
165
166#[derive(Debug, Clone, Copy, PartialEq, Eq)]
168enum PeerProtocolVersion {
169 V1,
171 V2,
173}
174
175impl PeerProtocolVersion {
176 fn envelope_overhead(&self) -> usize {
178 match self {
179 PeerProtocolVersion::V1 => V1_ENVELOPE_OVERHEAD,
180 PeerProtocolVersion::V2 => V2_ENVELOPE_OVERHEAD,
181 }
182 }
183}
184
185#[derive(Debug, Encode, Decode)]
186enum StatementMessage {
187 #[codec(index = 0)]
188 Statements(StatementBatch),
189 #[codec(index = 1)]
191 ExplicitTopicAffinity(AffinityFilter),
192}
193
194const STATEMENTS_VARIANT_INDEX: u8 = 0;
196
197impl StatementMessage {
198 fn encode_statement_refs(statements: &[&Statement]) -> Vec<u8> {
201 let mut out = Vec::new();
202 STATEMENTS_VARIANT_INDEX.encode_to(&mut out);
203 statements.encode_to(&mut out);
204 out
205 }
206}
207
208pub type StatementImportFuture = oneshot::Receiver<SubmitResult>;
210
211mod rep {
212 use sc_network::ReputationChange as Rep;
213 pub const ANY_STATEMENT: Rep = Rep::new(-(1 << 4), "Any statement");
218 pub const ANY_STATEMENT_REFUND: Rep = Rep::new(1 << 4, "Any statement (refund)");
220 pub const GOOD_STATEMENT: Rep = Rep::new(1 << 8, "Good statement");
222 pub const INVALID_STATEMENT: Rep = Rep::new(-(1 << 12), "Invalid statement");
224 pub const DUPLICATE_STATEMENT: Rep = Rep::new(-(1 << 7), "Duplicate statement");
226 pub const STATEMENT_FLOODING: Rep = Rep::new_fatal("Statement flooding");
228 pub const BAD_MESSAGE: Rep = Rep::new(-(1 << 12), "Bad statement message");
230}
231
232const LOG_TARGET: &str = "statement-gossip";
233const STATEMENT_PROTOCOL_V2: &str = "statement/2";
236const STATEMENT_PROTOCOL_V1: &str = "statement/1";
238const SEND_TIMEOUT: std::time::Duration = std::time::Duration::from_secs(10);
240const INITIAL_SYNC_BURST_INTERVAL: std::time::Duration = std::time::Duration::from_millis(10);
242
243const INITIAL_SYNC_SCAN_LIMIT: usize = 4096;
246const PENDING_AFFINITIES_INTERVAL: std::time::Duration = std::time::Duration::from_secs(1);
248const SYNC_RECOVERY_READD_DELAY: std::time::Duration = std::time::Duration::from_secs(60);
250
251struct Metrics {
252 propagated_statements: Counter<U64>,
253 known_statements_received: Counter<U64>,
254 skipped_oversized_statements: Counter<U64>,
255 propagated_statements_chunks: HistogramVec,
256 pending_statements: Gauge<U64>,
257 ignored_statements: Counter<U64>,
258 peers_connected: GaugeVec<U64>,
259 statements_received: Counter<U64>,
260 bytes_sent_total: Counter<U64>,
261 bytes_received_total: Counter<U64>,
262 sent_latency_seconds: Histogram,
263 initial_sync_statements_sent: Counter<U64>,
264 initial_sync_bursts_total: Counter<U64>,
265 initial_sync_in_flight_bytes: Gauge<U64>,
266 propagation_in_flight_bytes: Gauge<U64>,
267 initial_sync_peers_active: Gauge<U64>,
268 initial_sync_duration_seconds: HistogramVec,
269 statement_flooding_detected: Counter<U64>,
270 send_failures: CounterVec<U64>,
271 undelivered_statements: CounterVec<U64>,
272}
273
274mod send_failure {
275 pub const NETWORK: &str = "network";
277 pub const TIMEOUT: &str = "timeout";
279 pub const NO_SINK: &str = "no_sink";
281 pub const OUTBOX_FULL: &str = "outbox_full";
283}
284
285mod sync_outcome {
286 pub const COMPLETED: &str = "completed";
288 pub const ABANDONED: &str = "abandoned";
290}
291
292impl Metrics {
293 fn register(r: &Registry) -> Result<Self, PrometheusError> {
294 let peers_connected = register(
295 GaugeVec::new(
296 Opts::new(
297 "substrate_sync_statement_peers_connected",
298 "Number of peers connected using the statement protocol by kind",
299 ),
300 &["kind"],
301 )?,
302 r,
303 )?;
304 peers_connected.with_label_values(&["full"]).set(0);
305 peers_connected.with_label_values(&["light"]).set(0);
306
307 Ok(Self {
308 propagated_statements: register(
309 Counter::new(
310 "substrate_sync_propagated_statements",
311 "Total statements propagated to peers, counted once per recipient (a statement sent to N peers increments by N)",
312 )?,
313 r,
314 )?,
315 known_statements_received: register(
316 Counter::new(
317 "substrate_sync_known_statement_received",
318 "Number of statements received via gossiping that were already in the statement store",
319 )?,
320 r,
321 )?,
322 skipped_oversized_statements: register(
323 Counter::new(
324 "substrate_sync_skipped_oversized_statements",
325 "Number of oversized statements that were skipped to be gossiped",
326 )?,
327 r,
328 )?,
329 propagated_statements_chunks: register(
330 HistogramVec::new(
331 HistogramOpts::new(
332 "substrate_sync_propagated_statements_chunks",
333 "Distribution of chunk sizes when sending statements, by send path. Initial \
334 sync fills every chunk to the size limit, propagation does not",
335 )
336 .buckets(exponential_buckets(1.0, 2.0, 14)?),
337 &["kind"],
338 )?,
339 r,
340 )?,
341 pending_statements: register(
342 Gauge::new(
343 "substrate_sync_pending_statement_validations",
344 "Number of pending statement validations, sampled once per propagation tick",
345 )?,
346 r,
347 )?,
348 ignored_statements: register(
349 Counter::new(
350 "substrate_sync_ignored_statements",
351 "Number of statements ignored due to exceeding MAX_PENDING_STATEMENTS limit",
352 )?,
353 r,
354 )?,
355 peers_connected,
356 statements_received: register(
357 Counter::new(
358 "substrate_sync_statements_received",
359 "Total number of statements received from peers",
360 )?,
361 r,
362 )?,
363 bytes_sent_total: register(
364 Counter::new(
365 "substrate_sync_statement_bytes_sent_total",
366 "Total bytes sent for statement protocol messages",
367 )?,
368 r,
369 )?,
370 bytes_received_total: register(
371 Counter::new(
372 "substrate_sync_statement_bytes_received_total",
373 "Total bytes received for statement protocol messages (includes bytes from notifications that are later discarded — e.g. while major-syncing)",
374 )?,
375 r,
376 )?,
377 sent_latency_seconds: register(
378 Histogram::with_opts(
379 HistogramOpts::new(
380 "substrate_sync_statement_sent_latency_seconds",
381 "Time to send statement messages to peers",
382 )
383 .buckets(vec![0.000_001, 0.000_01, 0.000_1, 0.001, 0.01, 0.1, 1.0]),
385 )?,
386 r,
387 )?,
388 initial_sync_statements_sent: register(
389 Counter::new(
390 "substrate_sync_initial_sync_statements_sent",
391 "Total statements sent during initial sync bursts to newly connected peers",
392 )?,
393 r,
394 )?,
395 initial_sync_bursts_total: register(
396 Counter::new(
397 "substrate_sync_initial_sync_bursts_total",
398 "Total initial-sync burst rounds attempted (includes rounds that return early with no hashes left)",
399 )?,
400 r,
401 )?,
402 initial_sync_in_flight_bytes: register(
403 Gauge::new(
404 "substrate_sync_initial_sync_in_flight_bytes",
405 "Encoded bytes of initial-sync chunks currently queued for sending",
406 )?,
407 r,
408 )?,
409 propagation_in_flight_bytes: register(
410 Gauge::new(
411 "substrate_sync_propagation_in_flight_bytes",
412 "Encoded bytes of propagation chunks currently queued for sending",
413 )?,
414 r,
415 )?,
416 initial_sync_peers_active: register(
417 Gauge::new(
418 "substrate_sync_initial_sync_peers_active",
419 "Number of peers currently being synced via initial sync",
420 )?,
421 r,
422 )?,
423 initial_sync_duration_seconds: register(
424 HistogramVec::new(
425 HistogramOpts::new(
426 "substrate_sync_initial_sync_duration_seconds",
427 "Per-peer duration of initial sync, by outcome: completed (backlog drained) or abandoned (ended with statements still queued)",
428 )
429 .buckets(vec![0.01, 0.05, 0.1, 0.25, 0.5, 1.0, 2.5, 5.0, 10.0, 30.0, 60.0]),
430 &["outcome"],
431 )?,
432 r,
433 )?,
434 statement_flooding_detected: register(
435 Counter::new(
436 "substrate_sync_statement_flooding_detected",
437 "Number of peers disconnected for exceeding statement rate limits",
438 )?,
439 r,
440 )?,
441 send_failures: register(
442 CounterVec::new(
443 Opts::new(
444 "substrate_sync_statement_send_failures_total",
445 "Total statement sends that never reached the peer, by reason",
446 ),
447 &["reason"],
448 )?,
449 r,
450 )?,
451 undelivered_statements: register(
452 CounterVec::new(
453 Opts::new(
454 "substrate_sync_statement_undelivered_total",
455 "Total statements whose delivery was abandoned, so the peer never received them, by reason",
456 ),
457 &["reason"],
458 )?,
459 r,
460 )?,
461 })
462 }
463}
464
465pub struct StatementHandlerPrototype {
467 protocol_name: ProtocolName,
468 notification_service: Box<dyn NotificationService>,
469}
470
471impl StatementHandlerPrototype {
472 pub fn new<
474 Hash: AsRef<[u8]>,
475 Block: BlockT,
476 Net: NetworkBackend<Block, <Block as BlockT>::Hash>,
477 >(
478 genesis_hash: Hash,
479 fork_id: Option<&str>,
480 metrics: NotificationMetrics,
481 peer_store_handle: Arc<dyn PeerStoreProvider>,
482 ) -> (Self, Net::NotificationProtocolConfig) {
483 let genesis_hash = genesis_hash.as_ref();
484 let hex = array_bytes::bytes2hex("", genesis_hash);
485 let (protocol_name, fallback_name) = if let Some(fork_id) = fork_id {
486 (
487 format!("/{hex}/{fork_id}/{STATEMENT_PROTOCOL_V2}"),
488 format!("/{hex}/{fork_id}/{STATEMENT_PROTOCOL_V1}"),
489 )
490 } else {
491 (format!("/{hex}/{STATEMENT_PROTOCOL_V2}"), format!("/{hex}/{STATEMENT_PROTOCOL_V1}"))
492 };
493 let (config, notification_service) = Net::notification_config(
494 protocol_name.clone().into(),
495 vec![fallback_name.into()],
496 MAX_STATEMENT_NOTIFICATION_SIZE,
497 None,
498 SetConfig {
499 in_peers: 0,
500 out_peers: 0,
501 reserved_nodes: Vec::new(),
502 non_reserved_mode: NonReservedPeerMode::Deny,
503 },
504 metrics,
505 peer_store_handle,
506 );
507
508 (Self { protocol_name: protocol_name.into(), notification_service }, config)
509 }
510
511 pub fn build<
516 N: NetworkPeers + NetworkEventStream,
517 S: SyncEventStream + sp_consensus::SyncOracle,
518 >(
519 self,
520 network: N,
521 sync: S,
522 statement_store: Arc<dyn StatementStore>,
523 metrics_registry: Option<&Registry>,
524 executor: impl Fn(Pin<Box<dyn Future<Output = ()> + Send>>) + Send,
525 mut num_submission_workers: usize,
526 statements_per_second: u32,
527 ) -> error::Result<StatementHandler<N, S>> {
528 let sync_event_stream = sync.event_stream("statement-handler-sync");
529 let (queue_sender, queue_receiver) = async_channel::unbounded();
531
532 if num_submission_workers == 0 {
533 log::warn!(
534 target: LOG_TARGET,
535 "num_submission_workers is 0, defaulting to 1"
536 );
537 num_submission_workers = 1;
538 }
539
540 let statements_per_second = match NonZeroU32::new(statements_per_second) {
541 Some(rate) => rate,
542 None => {
543 log::warn!(
544 target: LOG_TARGET,
545 "statements_per_second is 0, defaulting to {}",
546 DEFAULT_STATEMENTS_PER_SECOND
547 );
548 NonZeroU32::new(DEFAULT_STATEMENTS_PER_SECOND)
549 .expect("DEFAULT_STATEMENTS_PER_SECOND is nonzero")
550 },
551 };
552
553 let metrics =
554 if let Some(r) = metrics_registry { Some(Metrics::register(r)?) } else { None };
555
556 for _ in 0..num_submission_workers {
557 let store = statement_store.clone();
558 let mut queue_receiver = queue_receiver.clone();
559 executor(
560 async move {
561 loop {
562 let task: Option<(Statement, oneshot::Sender<SubmitResult>)> =
563 queue_receiver.next().await;
564 match task {
565 None => return,
566 Some((statement, completion)) => {
567 let result = store.submit(statement, StatementSource::Network);
568 if completion.send(result).is_err() {
569 log::debug!(
570 target: LOG_TARGET,
571 "Error sending validation completion"
572 );
573 }
574 },
575 }
576 }
577 }
578 .boxed(),
579 );
580 }
581
582 let handler = StatementHandler {
583 protocol_name: self.protocol_name,
584 notification_service: self.notification_service,
585 propagate_timeout: (Box::pin(interval(PROPAGATE_TIMEOUT))
586 as Pin<Box<dyn Stream<Item = ()> + Send>>)
587 .fuse(),
588 pending_statements: FuturesUnordered::new(),
589 pending_statements_peers: HashMap::new(),
590 recently_received_statements: HashMap::new(),
591 network,
592 sync,
593 sync_event_stream: sync_event_stream.fuse(),
594 peers: HashMap::new(),
595 statement_store,
596 queue_sender,
597 statements_per_second,
598 metrics,
599 initial_sync_timeout: Box::pin(tokio::time::sleep(INITIAL_SYNC_BURST_INTERVAL).fuse()),
600 pending_affinities_timeout: Box::pin(
601 tokio::time::sleep(PENDING_AFFINITIES_INTERVAL).fuse(),
602 ),
603 pending_initial_syncs: HashMap::new(),
604 initial_sync_peer_queue: VecDeque::new(),
605 next_initial_sync_id: 0,
606 initial_sync_in_flight_bytes: 0,
607 propagation_outboxes: HashMap::new(),
608 in_flight_chunks: HashMap::new(),
609 next_chunk_id: 0,
610 propagation_in_flight_bytes: 0,
611 parked_propagations: VecDeque::new(),
612 pending_sends: FuturesUnordered::new(),
613 deferred_peers: HashSet::new(),
614 dropped_statements_during_sync: false,
615 sync_recovery_peer: None,
616 sync_recovery_readd_timeout: Box::pin(pending().fuse()),
617 };
618
619 Ok(handler)
620 }
621}
622
623pub struct StatementHandler<
625 N: NetworkPeers + NetworkEventStream,
626 S: SyncEventStream + sp_consensus::SyncOracle,
627> {
628 protocol_name: ProtocolName,
629 propagate_timeout: stream::Fuse<Pin<Box<dyn Stream<Item = ()> + Send>>>,
631 pending_statements:
633 FuturesUnordered<Pin<Box<dyn Future<Output = (Hash, Option<SubmitResult>)> + Send>>>,
634 pending_statements_peers: HashMap<Hash, HashSet<PeerId>>,
639 recently_received_statements: HashMap<Hash, HashSet<PeerId>>,
644 network: N,
646 sync: S,
648 sync_event_stream: stream::Fuse<Pin<Box<dyn Stream<Item = SyncEvent> + Send>>>,
650 notification_service: Box<dyn NotificationService>,
652 peers: HashMap<PeerId, Peer>,
654 statement_store: Arc<dyn StatementStore>,
655 queue_sender: async_channel::Sender<(Statement, oneshot::Sender<SubmitResult>)>,
656 statements_per_second: NonZeroU32,
658 metrics: Option<Metrics>,
660 initial_sync_timeout: Pin<Box<dyn FusedFuture<Output = ()> + Send>>,
662 pending_affinities_timeout: Pin<Box<dyn FusedFuture<Output = ()> + Send>>,
664 pending_initial_syncs: HashMap<PeerId, PendingInitialSync>,
666 initial_sync_peer_queue: VecDeque<PeerId>,
668 next_initial_sync_id: u64,
670 initial_sync_in_flight_bytes: u64,
673 propagation_outboxes: HashMap<PeerId, VecDeque<Hash>>,
676 in_flight_chunks: HashMap<PeerId, u64>,
679 next_chunk_id: u64,
681 propagation_in_flight_bytes: u64,
684 parked_propagations: VecDeque<PeerId>,
688 pending_sends: PendingSends,
690 deferred_peers: HashSet<PeerId>,
693 dropped_statements_during_sync: bool,
695 sync_recovery_peer: Option<PeerId>,
697 sync_recovery_readd_timeout: Pin<Box<dyn FusedFuture<Output = ()> + Send>>,
699}
700
701trait TokenBucket: fmt::Debug + Send + Sync {
703 fn would_exceed(&self, n: NonZeroU32) -> bool;
705}
706
707impl<C> TokenBucket for RateLimiter<NotKeyed, InMemoryState, C, NoOpMiddleware<C::Instant>>
708where
709 C: Clock + fmt::Debug + Send + Sync,
710 C::Instant: fmt::Debug + Send + Sync,
711{
712 fn would_exceed(&self, n: NonZeroU32) -> bool {
713 !matches!(self.check_n(n), Ok(Ok(())))
714 }
715}
716
717#[derive(Debug)]
722struct PeerRateLimiter {
723 bucket: Box<dyn TokenBucket>,
724}
725
726impl PeerRateLimiter {
727 fn new(statements_per_second: NonZeroU32, burst: NonZeroU32) -> Self {
728 Self::with_clock(statements_per_second, burst, &DefaultClock::default())
729 }
730
731 fn with_clock<C>(statements_per_second: NonZeroU32, burst: NonZeroU32, clock: &C) -> Self
733 where
734 C: Clock + fmt::Debug + Send + Sync + 'static,
735 C::Instant: fmt::Debug + Send + Sync,
736 {
737 let quota = Quota::per_second(statements_per_second).allow_burst(burst);
738 Self { bucket: Box::new(RateLimiter::direct_with_clock(quota, clock)) }
739 }
740
741 fn is_flooding(&self, count: usize) -> bool {
743 if count > u32::MAX as usize {
744 return true;
745 }
746
747 let Some(n) = NonZeroU32::new(count as u32) else {
748 return false;
749 };
750 self.bucket.would_exceed(n)
751 }
752}
753
754#[cfg_attr(not(any(test, feature = "test-helpers")), doc(hidden))]
756#[derive(Debug)]
757pub struct Peer {
758 rate_limiter: PeerRateLimiter,
760 protocol_version: PeerProtocolVersion,
762 topic_affinity: Option<AffinityFilter>,
765 is_light: bool,
768 pending_topic_affinity: Option<AffinityFilter>,
772 sync_watermark: u64,
774}
775
776struct PendingInitialSync {
779 cursor: u64,
781 watermark: u64,
783 started_at: Instant,
784 sync_id: u64,
787}
788
789enum SendOutcome {
790 Sent,
792 NetworkError(error::Error),
794 TimedOut,
796}
797
798enum SendKind {
799 Propagation,
800 InitialSync {
801 sync_id: u64,
802 next_cursor: u64,
804 },
805}
806
807impl SendKind {
808 fn label(&self) -> &'static str {
809 match self {
810 Self::Propagation => "propagation",
811 Self::InitialSync { .. } => "initial_sync",
812 }
813 }
814}
815
816struct PendingSendResult {
818 peer: PeerId,
819 statement_count: usize,
820 bytes_sent: u64,
821 result: SendOutcome,
822 kind: SendKind,
823 chunk_id: u64,
826}
827
828type PendingSends =
830 FuturesUnordered<Pin<Box<dyn Future<Output = PendingSendResult> + Send + 'static>>>;
831
832const V1_ENVELOPE_OVERHEAD: usize = 5;
834
835const V2_ENVELOPE_OVERHEAD: usize = 1 + V1_ENVELOPE_OVERHEAD;
837
838fn max_statement_payload_size(envelope_overhead: usize) -> usize {
841 debug_assert_eq!(
842 V1_ENVELOPE_OVERHEAD,
843 Compact::<u32>::max_encoded_len(),
844 "V1_ENVELOPE_OVERHEAD must equal Compact::<u32>::max_encoded_len()"
845 );
846 MAX_STATEMENT_NOTIFICATION_SIZE as usize - envelope_overhead
847}
848
849fn unix_timestamp_secs() -> u64 {
850 std::time::SystemTime::now()
851 .duration_since(std::time::UNIX_EPOCH)
852 .unwrap_or_default()
853 .as_secs()
854}
855
856fn fetch_admitted_chunk(
864 store: &dyn StatementStore,
865 recently_received_statements: &HashMap<Hash, HashSet<PeerId>>,
866 pending_statements_peers: &HashMap<Hash, HashSet<PeerId>>,
867 who: &PeerId,
868 peer_data: &Peer,
869 cursor: u64,
870 watermark: u64,
871 max_size: usize,
872) -> sp_statement_store::Result<(AdmittedBatch, usize)> {
873 let now = unix_timestamp_secs();
874 let mut accumulated_size = 0;
875 let batch = store.admitted_statements(
876 cursor,
877 watermark,
878 INITIAL_SYNC_SCAN_LIMIT,
879 &mut |hash, encoded, stmt| {
880 if stmt.is_expired(now) {
881 return FilterDecision::Skip;
882 }
883 if peer_data.topic_affinity.as_ref().is_some_and(|a| !a.matches_statement(stmt)) {
884 return FilterDecision::Skip;
885 }
886 if has_received_from(recently_received_statements, pending_statements_peers, hash, who)
888 {
889 return FilterDecision::Skip;
890 }
891 if accumulated_size > 0 && accumulated_size + encoded.len() > max_size {
892 return FilterDecision::Abort;
893 }
894 accumulated_size += encoded.len();
895 FilterDecision::Take
896 },
897 )?;
898 Ok((batch, accumulated_size))
899}
900
901fn fetch_statement_chunk(
904 store: &dyn StatementStore,
905 recently_received_statements: &HashMap<Hash, HashSet<PeerId>>,
906 pending_statements_peers: &HashMap<Hash, HashSet<PeerId>>,
907 who: &PeerId,
908 peer_data: &Peer,
909 hashes: &[Hash],
910 max_size: usize,
911) -> sp_statement_store::Result<(Vec<(Hash, Statement)>, usize, usize)> {
912 let now = unix_timestamp_secs();
913 let mut accumulated_size = 0;
914 let (statements, processed) =
915 store.statements_by_hashes(hashes, &mut |hash, encoded, stmt| {
916 if stmt.is_expired(now) {
917 return FilterDecision::Skip;
918 }
919 if peer_data.topic_affinity.as_ref().is_some_and(|a| !a.matches_statement(stmt)) {
920 return FilterDecision::Skip;
921 }
922 if has_received_from(recently_received_statements, pending_statements_peers, hash, who)
924 {
925 return FilterDecision::Skip;
926 }
927 if accumulated_size > 0 && accumulated_size + encoded.len() > max_size {
928 return FilterDecision::Abort;
929 }
930 accumulated_size += encoded.len();
931 FilterDecision::Take
932 })?;
933 Ok((statements, processed, accumulated_size))
934}
935
936async fn send_with_timeout<F>(send: F) -> SendOutcome
937where
938 F: Future<Output = Result<(), error::Error>>,
939{
940 match timeout(SEND_TIMEOUT, send).await {
941 Ok(Ok(())) => SendOutcome::Sent,
942 Ok(Err(error)) => SendOutcome::NetworkError(error),
943 Err(_elapsed) => SendOutcome::TimedOut,
944 }
945}
946
947fn has_received_from(
951 recently_received_statements: &HashMap<Hash, HashSet<PeerId>>,
952 pending_statements_peers: &HashMap<Hash, HashSet<PeerId>>,
953 hash: &Hash,
954 who: &PeerId,
955) -> bool {
956 recently_received_statements.get(hash).is_some_and(|peers| peers.contains(who)) ||
957 pending_statements_peers.get(hash).is_some_and(|peers| peers.contains(who))
958}
959
960impl Peer {
961 #[cfg(any(test, feature = "test-helpers"))]
963 pub fn new_for_testing(statements_per_second: NonZeroU32, burst: NonZeroU32) -> Self {
964 Self {
965 rate_limiter: PeerRateLimiter::new(statements_per_second, burst),
966 protocol_version: PeerProtocolVersion::V1,
967 topic_affinity: None,
968 is_light: false,
969 pending_topic_affinity: None,
970 sync_watermark: 0,
971 }
972 }
973
974 fn can_receive(&self) -> bool {
978 !(self.is_light &&
979 self.protocol_version == PeerProtocolVersion::V2 &&
980 self.topic_affinity.is_none())
981 }
982
983 fn kind(&self) -> &'static str {
984 if self.is_light {
985 "light"
986 } else {
987 "full"
988 }
989 }
990}
991
992impl<N, S> StatementHandler<N, S>
993where
994 N: NetworkPeers + NetworkEventStream,
995 S: SyncEventStream + sp_consensus::SyncOracle,
996{
997 #[cfg(any(test, feature = "test-helpers"))]
999 pub fn new_for_testing(
1000 protocol_name: ProtocolName,
1001 notification_service: Box<dyn NotificationService>,
1002 propagate_timeout: stream::Fuse<Pin<Box<dyn Stream<Item = ()> + Send>>>,
1003 network: N,
1004 sync: S,
1005 sync_event_stream: stream::Fuse<Pin<Box<dyn Stream<Item = SyncEvent> + Send>>>,
1006 peers: HashMap<PeerId, Peer>,
1007 statement_store: Arc<dyn StatementStore>,
1008 queue_sender: async_channel::Sender<(Statement, oneshot::Sender<SubmitResult>)>,
1009 statements_per_second: NonZeroU32,
1010 ) -> Self {
1011 Self {
1012 protocol_name,
1013 notification_service,
1014 propagate_timeout,
1015 pending_statements: FuturesUnordered::new(),
1016 pending_statements_peers: HashMap::new(),
1017 recently_received_statements: HashMap::new(),
1018 network,
1019 sync,
1020 sync_event_stream,
1021 peers,
1022 statement_store,
1023 queue_sender,
1024 statements_per_second,
1025 metrics: None,
1026 initial_sync_timeout: Box::pin(pending().fuse()),
1027 pending_affinities_timeout: Box::pin(pending().fuse()),
1028 pending_initial_syncs: HashMap::new(),
1029 initial_sync_peer_queue: VecDeque::new(),
1030 next_initial_sync_id: 0,
1031 initial_sync_in_flight_bytes: 0,
1032 propagation_outboxes: HashMap::new(),
1033 in_flight_chunks: HashMap::new(),
1034 next_chunk_id: 0,
1035 propagation_in_flight_bytes: 0,
1036 parked_propagations: VecDeque::new(),
1037 pending_sends: FuturesUnordered::new(),
1038 deferred_peers: HashSet::new(),
1039 dropped_statements_during_sync: false,
1040 sync_recovery_peer: None,
1041 sync_recovery_readd_timeout: Box::pin(pending().fuse()),
1042 }
1043 }
1044
1045 #[cfg(any(test, feature = "test-helpers"))]
1047 pub fn pending_statements_mut(
1048 &mut self,
1049 ) -> &mut FuturesUnordered<Pin<Box<dyn Future<Output = (Hash, Option<SubmitResult>)> + Send>>>
1050 {
1051 &mut self.pending_statements
1052 }
1053
1054 pub async fn run(mut self) {
1057 loop {
1058 futures::select_biased! {
1059 send_result = self.pending_sends.select_next_some() => {
1060 self.handle_send_result(send_result);
1061 },
1062 _ = self.propagate_timeout.next() => {
1063 self.propagate_statements().await;
1064 self.shrink_pending_statements_peers();
1065 self.metrics.as_ref().map(|metrics| {
1066 metrics.pending_statements.set(self.pending_statements.len() as u64);
1067 });
1068 },
1069 (hash, result) = self.pending_statements.select_next_some() => {
1070 self.on_statement_submit_result(hash, result);
1071 },
1072 sync_event = self.sync_event_stream.next() => {
1073 if let Some(sync_event) = sync_event {
1074 self.handle_sync_event(sync_event);
1075 } else {
1076 return;
1078 }
1079 }
1080 event = self.notification_service.next_event().fuse() => {
1081 if let Some(event) = event {
1082 self.handle_notification_event(event).await
1083 } else {
1084 return
1086 }
1087 }
1088 _ = &mut self.initial_sync_timeout => {
1089 self.process_initial_sync_burst();
1090 self.initial_sync_timeout =
1091 Box::pin(tokio::time::sleep(INITIAL_SYNC_BURST_INTERVAL).fuse());
1092 },
1093 _ = &mut self.pending_affinities_timeout => {
1094 self.process_pending_affinities();
1095 self.pending_affinities_timeout =
1096 Box::pin(tokio::time::sleep(PENDING_AFFINITIES_INTERVAL).fuse());
1097 },
1098 _ = &mut self.sync_recovery_readd_timeout => {
1099 self.try_readd_sync_recovery_peer();
1100 self.sync_recovery_readd_timeout = Box::pin(pending().fuse());
1101 },
1102 }
1103
1104 if !self.sync.is_major_syncing() {
1105 self.drain_deferred_peers();
1106 self.start_sync_recovery();
1107 }
1108 }
1109 }
1110
1111 fn shrink_pending_statements_peers(&mut self) {
1113 const MIN_RETAINED_CAPACITY: usize = 1024;
1114 let map = &mut self.pending_statements_peers;
1115 if map.capacity() > MIN_RETAINED_CAPACITY && map.capacity() / 4 > map.len() {
1116 map.shrink_to(MIN_RETAINED_CAPACITY);
1117 }
1118 }
1119
1120 fn record_send_failure(&self, reason: &str) {
1122 self.metrics.as_ref().map(|metrics| {
1123 metrics.send_failures.with_label_values(&[reason]).inc();
1124 });
1125 }
1126
1127 fn record_abandoned_send(&self, reason: &str, statement_count: usize) {
1129 self.record_send_failure(reason);
1130 self.metrics.as_ref().map(|metrics| {
1131 metrics
1132 .undelivered_statements
1133 .with_label_values(&[reason])
1134 .inc_by(statement_count as u64);
1135 });
1136 }
1137
1138 fn drain_deferred_peers(&mut self) {
1140 if self.deferred_peers.is_empty() {
1141 return;
1142 }
1143
1144 log::debug!(
1145 target: LOG_TARGET,
1146 "Major sync complete, adding {} deferred statement peers",
1147 self.deferred_peers.len(),
1148 );
1149
1150 let addrs: HashSet<multiaddr::Multiaddr> = self
1151 .deferred_peers
1152 .drain()
1153 .map(|p| {
1154 iter::once(multiaddr::Protocol::P2p(p.into())).collect::<multiaddr::Multiaddr>()
1155 })
1156 .collect();
1157
1158 if let Err(err) = self.network.add_peers_to_reserved_set(self.protocol_name.clone(), addrs)
1159 {
1160 log::warn!(target: LOG_TARGET, "Failed to add deferred peers: {err}");
1161 }
1162 }
1163
1164 fn start_sync_recovery(&mut self) {
1169 if !self.dropped_statements_during_sync {
1170 return;
1171 }
1172 self.dropped_statements_during_sync = false;
1173
1174 if self.sync_recovery_peer.is_some() {
1175 return;
1176 }
1177
1178 let Some(&peer_id) = self.peers.keys().choose(&mut rand::thread_rng()) else {
1179 return;
1180 };
1181
1182 log::trace!(
1183 target: LOG_TARGET,
1184 "Major sync complete, force-reconnecting {peer_id} for statement recovery",
1185 );
1186
1187 if let Err(err) = self.network.remove_peers_from_reserved_set(
1188 self.protocol_name.clone(),
1189 iter::once(peer_id).collect(),
1190 ) {
1191 log::warn!(target: LOG_TARGET, "Failed to remove peer {peer_id} for sync recovery: {err}");
1192 return;
1193 }
1194
1195 self.sync_recovery_peer = Some(peer_id);
1196 self.sync_recovery_readd_timeout =
1197 Box::pin(tokio::time::sleep(SYNC_RECOVERY_READD_DELAY).fuse());
1198 }
1199
1200 fn try_readd_sync_recovery_peer(&mut self) {
1202 let Some(peer_id) = self.sync_recovery_peer.take() else { return };
1203 log::trace!(
1204 target: LOG_TARGET,
1205 "Re-adding {peer_id} to reserved set after sync recovery window",
1206 );
1207 let addr =
1208 iter::once(multiaddr::Protocol::P2p(peer_id.into())).collect::<multiaddr::Multiaddr>();
1209 if let Err(err) = self
1210 .network
1211 .add_peers_to_reserved_set(self.protocol_name.clone(), iter::once(addr).collect())
1212 {
1213 log::warn!(target: LOG_TARGET, "Failed to re-add sync recovery peer {peer_id}: {err}");
1214 }
1215 }
1216
1217 fn handle_sync_event(&mut self, event: SyncEvent) {
1226 match event {
1227 SyncEvent::PeerConnected { peer_id: remote, roles: _ } => {
1228 if self.sync.is_major_syncing() {
1229 log::trace!(
1230 target: LOG_TARGET,
1231 "Major sync in progress, deferring connection to {remote}",
1232 );
1233 self.deferred_peers.insert(remote);
1234 return;
1235 }
1236 let addr = iter::once(multiaddr::Protocol::P2p(remote.into()))
1237 .collect::<multiaddr::Multiaddr>();
1238 let result = self.network.add_peers_to_reserved_set(
1239 self.protocol_name.clone(),
1240 iter::once(addr).collect(),
1241 );
1242 if let Err(err) = result {
1243 log::error!(target: LOG_TARGET, "Add reserved peer failed: {}", err);
1244 }
1245 },
1246 SyncEvent::PeerDisconnected(remote) => {
1247 if self.deferred_peers.remove(&remote) {
1248 return;
1249 }
1250 let result = self.network.remove_peers_from_reserved_set(
1251 self.protocol_name.clone(),
1252 iter::once(remote).collect(),
1253 );
1254 if let Err(err) = result {
1255 log::error!(target: LOG_TARGET, "Failed to remove reserved peer: {err}");
1256 }
1257 },
1258 }
1259 }
1260
1261 async fn handle_notification_event(&mut self, event: NotificationEvent) {
1271 match event {
1272 NotificationEvent::ValidateInboundSubstream { peer, handshake, result_tx, .. } => {
1273 let result = self
1275 .network
1276 .peer_role(peer, handshake)
1277 .map_or(ValidationResult::Reject, |_| ValidationResult::Accept);
1278 let _ = result_tx.send(result);
1279 },
1280 NotificationEvent::NotificationStreamOpened {
1281 peer,
1282 negotiated_fallback,
1283 handshake,
1284 ..
1285 } => {
1286 let protocol_version = if negotiated_fallback.is_some() {
1289 PeerProtocolVersion::V1
1290 } else {
1291 PeerProtocolVersion::V2
1292 };
1293 let Some(peer_role) = self.network.peer_role(peer, handshake) else {
1294 log::debug!(
1295 target: LOG_TARGET,
1296 "Peer {peer} connected but role could not be determined, ignoring"
1297 );
1298 return;
1299 };
1300 let is_light = peer_role.is_light();
1301 log::debug!(
1302 target: LOG_TARGET,
1303 "Peer {peer} connected with statement protocol {protocol_version:?}, role={peer_role:?}"
1304 );
1305 let _was_in = self.peers.insert(
1306 peer,
1307 Peer {
1308 rate_limiter: PeerRateLimiter::new(
1309 self.statements_per_second,
1310 NonZeroU32::new(
1311 self.statements_per_second.get() *
1312 config::STATEMENTS_BURST_COEFFICIENT,
1313 )
1314 .expect("burst capacity is nonzero"),
1315 ),
1316 protocol_version,
1317 topic_affinity: None,
1318 is_light,
1319 pending_topic_affinity: None,
1320 sync_watermark: 0,
1321 },
1322 );
1323 debug_assert!(_was_in.is_none());
1324
1325 self.metrics.as_ref().map(|metrics| {
1326 if let Some(peer) = self.peers.get(&peer) {
1327 metrics.peers_connected.with_label_values(&[peer.kind()]).inc();
1328 }
1329 });
1330
1331 if self.peers.get(&peer).map_or(false, |p| p.can_receive()) {
1334 self.schedule_initial_sync_for_peer(peer);
1335 }
1336 },
1337 NotificationEvent::NotificationStreamClosed { peer } => {
1338 let removed_peer = self.peers.remove(&peer);
1339 debug_assert!(removed_peer.is_some());
1340
1341 if let Some(removed_peer) = removed_peer {
1342 self.metrics.as_ref().map(|metrics| {
1343 metrics.peers_connected.with_label_values(&[removed_peer.kind()]).dec();
1344 });
1345 }
1346
1347 if let Some(pending) = self.pending_initial_syncs.remove(&peer) {
1348 self.record_initial_sync_completion(
1349 sync_outcome::ABANDONED,
1350 pending.started_at,
1351 );
1352 }
1353 self.initial_sync_peer_queue.retain(|p| *p != peer);
1354 self.propagation_outboxes.remove(&peer);
1355 self.in_flight_chunks.remove(&peer);
1356 },
1357 NotificationEvent::NotificationReceived { peer, notification } => {
1358 let bytes_received = notification.len() as u64;
1359 self.metrics.as_ref().map(|metrics| {
1360 metrics.bytes_received_total.inc_by(bytes_received);
1361 });
1362
1363 if self.sync.is_major_syncing() {
1365 log::trace!(
1366 target: LOG_TARGET,
1367 "{peer}: Ignoring statements while major syncing or offline"
1368 );
1369 self.dropped_statements_during_sync = true;
1370 return;
1371 }
1372
1373 let Some(peer_data) = self.peers.get(&peer) else {
1374 log::error!(target: LOG_TARGET, "Received notification from unknown peer {peer}");
1375 return;
1376 };
1377
1378 match peer_data.protocol_version {
1379 PeerProtocolVersion::V1 => {
1380 if let Ok(statements) =
1382 <StatementBatch as Decode>::decode(&mut notification.as_ref())
1383 {
1384 self.on_statements(peer, statements.into_inner());
1385 } else {
1386 log::debug!(
1387 target: LOG_TARGET,
1388 "Failed to decode v1 statement list from {peer}"
1389 );
1390 self.network.report_peer(peer, rep::BAD_MESSAGE);
1391 }
1392 },
1393 PeerProtocolVersion::V2 => {
1394 if let Ok(message) = StatementMessage::decode(&mut notification.as_ref()) {
1396 match message {
1397 StatementMessage::Statements(statements) => {
1398 self.on_statements(peer, statements.into_inner())
1399 },
1400 StatementMessage::ExplicitTopicAffinity(filter) => {
1401 if let Some(peer_data) = self.peers.get_mut(&peer) {
1402 if peer_data.rate_limiter.is_flooding(1) {
1403 log::debug!(
1404 target: LOG_TARGET,
1405 "Rate-limiting ExplicitTopicAffinity from {peer}"
1406 );
1407 self.network.report_peer(peer, rep::BAD_MESSAGE);
1408 } else {
1409 log::debug!(
1410 target: LOG_TARGET,
1411 "Received topic affinity filter from {peer}"
1412 );
1413 peer_data.pending_topic_affinity = Some(filter);
1416 }
1417 }
1418 },
1419 }
1420 } else {
1421 log::debug!(
1422 target: LOG_TARGET,
1423 "Failed to decode v2 statement message from {peer}"
1424 );
1425 self.network.report_peer(peer, rep::BAD_MESSAGE);
1426 }
1427 },
1428 }
1429 },
1430 }
1431 }
1432
1433 #[cfg_attr(not(any(test, feature = "test-helpers")), doc(hidden))]
1444 pub fn on_statements(&mut self, who: PeerId, statements: Statements) {
1445 log::trace!(target: LOG_TARGET, "Received {} statements from {}", statements.len(), who);
1446
1447 self.metrics.as_ref().map(|metrics| {
1448 metrics.statements_received.inc_by(statements.len() as u64);
1449 });
1450
1451 if let Some(ref mut peer) = self.peers.get_mut(&who) {
1452 if peer.rate_limiter.is_flooding(statements.len()) {
1453 log::warn!(
1454 target: LOG_TARGET,
1455 "Peer {} exceeded statement rate limit ({} statements/sec). Disconnecting.",
1456 who,
1457 self.statements_per_second
1458 );
1459
1460 self.network.report_peer(who, rep::STATEMENT_FLOODING);
1461
1462 self.network.disconnect_peer(who, self.protocol_name.clone());
1464
1465 if let Some(ref metrics) = self.metrics {
1466 metrics.statement_flooding_detected.inc();
1467 }
1468
1469 return;
1470 }
1471
1472 let mut statements_left = statements.len() as u64;
1473 for s in statements {
1474 if self.pending_statements.len() > MAX_PENDING_STATEMENTS {
1475 log::debug!(
1476 target: LOG_TARGET,
1477 "Ignoring {} statements that exceed `MAX_PENDING_STATEMENTS`({}) limit",
1478 statements_left,
1479 MAX_PENDING_STATEMENTS,
1480 );
1481 self.metrics.as_ref().map(|metrics| {
1482 metrics.ignored_statements.inc_by(statements_left);
1483 });
1484 break;
1485 }
1486
1487 let hash = s.hash();
1488
1489 if self.statement_store.has_statement(&hash) {
1490 self.metrics.as_ref().map(|metrics| {
1491 metrics.known_statements_received.inc();
1492 });
1493
1494 if let Some(peers) = self.recently_received_statements.get_mut(&hash) {
1500 peers.insert(who);
1501 }
1502
1503 if let Some(peers) = self.pending_statements_peers.get_mut(&hash) {
1509 if peers.insert(who) {
1510 self.network.report_peer(who, rep::ANY_STATEMENT);
1511 } else {
1512 log::trace!(
1513 target: LOG_TARGET,
1514 "Already received the statement from the same peer {who}.",
1515 );
1516 self.network.report_peer(who, rep::DUPLICATE_STATEMENT);
1517 }
1518 }
1519 continue;
1520 }
1521
1522 self.network.report_peer(who, rep::ANY_STATEMENT);
1523
1524 match self.pending_statements_peers.entry(hash) {
1525 Entry::Vacant(entry) => {
1526 let (completion_sender, completion_receiver) = oneshot::channel();
1527 match self.queue_sender.try_send((s, completion_sender)) {
1528 Ok(()) => {
1529 self.pending_statements.push(
1530 async move {
1531 let res = completion_receiver.await;
1532 (hash, res.ok())
1533 }
1534 .boxed(),
1535 );
1536 entry.insert(HashSet::from_iter([who]));
1537 },
1538 Err(async_channel::TrySendError::Full(_)) => {
1539 log::debug!(
1540 target: LOG_TARGET,
1541 "Dropped statement because validation channel is full",
1542 );
1543 },
1544 Err(async_channel::TrySendError::Closed(_)) => {
1545 log::trace!(
1546 target: LOG_TARGET,
1547 "Dropped statement because validation channel is closed",
1548 );
1549 },
1550 }
1551 },
1552 Entry::Occupied(mut entry) => {
1553 if !entry.get_mut().insert(who) {
1554 self.network.report_peer(who, rep::DUPLICATE_STATEMENT);
1556 }
1557 },
1558 }
1559
1560 statements_left -= 1;
1561 }
1562 }
1563 }
1564
1565 fn on_handle_statement_import(&mut self, who: PeerId, import: &SubmitResult) {
1580 match import {
1581 SubmitResult::New => self.network.report_peer(who, rep::GOOD_STATEMENT),
1582 SubmitResult::Known => self.network.report_peer(who, rep::ANY_STATEMENT_REFUND),
1583 SubmitResult::KnownExpired => {},
1584 SubmitResult::Rejected(_) => {},
1585 SubmitResult::Invalid(_) => self.network.report_peer(who, rep::INVALID_STATEMENT),
1586 SubmitResult::InternalError(_) => {},
1587 }
1588 }
1589
1590 fn on_statement_submit_result(&mut self, hash: Hash, result: Option<SubmitResult>) {
1594 if let Some(peers) = self.pending_statements_peers.remove(&hash) {
1595 if let Some(result) = result {
1596 for peer in &peers {
1597 self.on_handle_statement_import(*peer, &result);
1598 }
1599 if matches!(result, SubmitResult::New | SubmitResult::Known) {
1602 self.recently_received_statements.entry(hash).or_default().extend(peers);
1603 }
1604 }
1605 } else {
1606 log::warn!(target: LOG_TARGET, "Inconsistent state, no peers for pending statement!");
1607 }
1608 }
1609
1610 fn queue_statements_for_peer(&mut self, who: &PeerId, statements: &[(u64, Hash, Statement)]) {
1616 let Self {
1617 peers,
1618 propagation_outboxes,
1619 recently_received_statements,
1620 pending_statements_peers,
1621 ..
1622 } = self;
1623 let Some(peer) = peers.get(who) else {
1624 return;
1625 };
1626
1627 if !peer.can_receive() {
1628 return;
1629 }
1630
1631 let to_send = statements.iter().filter_map(|(seq, hash, stmt)| {
1632 if *seq < peer.sync_watermark {
1633 return None;
1634 }
1635 if has_received_from(recently_received_statements, pending_statements_peers, hash, who)
1637 {
1638 return None;
1639 }
1640 if peer.topic_affinity.as_ref().is_some_and(|a| !a.matches_statement(stmt)) {
1642 return None;
1643 }
1644 Some(*hash)
1645 });
1646
1647 let outbox = propagation_outboxes.entry(*who).or_default();
1648 let mut queued = 0;
1649 let mut overflow = 0;
1650 for hash in to_send {
1651 if outbox.len() == MAX_PROPAGATION_OUTBOX_LEN {
1654 outbox.pop_front();
1655 overflow += 1;
1656 }
1657 outbox.push_back(hash);
1658 queued += 1;
1659 }
1660
1661 log::trace!(target: LOG_TARGET, "We have {queued} statements that the peer doesn't know about");
1662
1663 if overflow > 0 {
1664 self.record_abandoned_send(send_failure::OUTBOX_FULL, overflow);
1665 }
1666 self.try_send_next_chunk(*who);
1667 }
1668
1669 fn try_send_next_chunk(&mut self, who: PeerId) {
1675 if self.in_flight_chunks.contains_key(&who) {
1676 return;
1677 }
1678
1679 loop {
1680 let Some(outbox) = self.propagation_outboxes.get(&who) else {
1681 return;
1682 };
1683 if outbox.is_empty() {
1684 self.propagation_outboxes.remove(&who);
1685 return;
1686 }
1687 let Some(peer_data) = self.peers.get(&who) else {
1688 self.propagation_outboxes.remove(&who);
1689 return;
1690 };
1691 if self.send_in_flight_bytes() >= self.propagation_send_budget() {
1694 if !self.parked_propagations.contains(&who) {
1697 self.parked_propagations.push_back(who);
1698 }
1699 return;
1700 }
1701 let peer_version = peer_data.protocol_version;
1702 let max_size = max_statement_payload_size(peer_version.envelope_overhead());
1703 let Some(outbox) = self.propagation_outboxes.get_mut(&who) else {
1704 return;
1705 };
1706 let (statements, processed, accumulated_size) = match fetch_statement_chunk(
1707 &*self.statement_store,
1708 &self.recently_received_statements,
1709 &self.pending_statements_peers,
1710 &who,
1711 peer_data,
1712 outbox.make_contiguous(),
1713 max_size,
1714 ) {
1715 Ok(result) => result,
1716 Err(e) => {
1717 log::warn!(
1721 target: LOG_TARGET,
1722 "Failed to fetch statements for propagation to {who}, retaining {} queued hashes: {e:?}",
1723 outbox.len(),
1724 );
1725 if !self.parked_propagations.contains(&who) {
1726 self.parked_propagations.push_back(who);
1727 }
1728 return;
1729 },
1730 };
1731
1732 debug_assert!(
1733 processed > 0,
1734 "a fetch from a non-empty outbox consumes at least one hash"
1735 );
1736 if processed == 0 {
1737 return;
1738 }
1739
1740 outbox.drain(..processed);
1743
1744 if accumulated_size > max_size {
1745 log::warn!(target: LOG_TARGET, "Statement too large, skipping");
1746 self.metrics.as_ref().map(|metrics| {
1747 metrics.skipped_oversized_statements.inc();
1748 });
1749 continue;
1750 }
1751
1752 if statements.is_empty() {
1753 continue;
1756 }
1757
1758 let statement_count = statements.len();
1759 let send_stmts: Vec<_> = statements.iter().map(|(_, stmt)| stmt).collect();
1760 let encoded = match peer_version {
1761 PeerProtocolVersion::V1 => send_stmts.encode(),
1762 PeerProtocolVersion::V2 => StatementMessage::encode_statement_refs(&send_stmts),
1763 };
1764 let bytes_sent = encoded.len() as u64;
1765 let Some(message_sink) = self.notification_service.message_sink(&who) else {
1766 let abandoned = statement_count +
1767 self.propagation_outboxes.get(&who).map_or(0, |outbox| outbox.len());
1768 log::debug!(
1769 target: LOG_TARGET,
1770 "Failed to get message sink for peer {who}, abandoning {abandoned} statements ({bytes_sent} bytes in the current chunk)",
1771 );
1772 self.record_abandoned_send(send_failure::NO_SINK, abandoned);
1773 self.propagation_outboxes.remove(&who);
1774 return;
1775 };
1776 let chunk_id = self.occupy_send_slot(who);
1777 let in_flight = self.propagation_in_flight_bytes.saturating_add(bytes_sent);
1778 self.set_propagation_in_flight_bytes(in_flight);
1779 let sent_latency =
1780 self.metrics.as_ref().map(|metrics| metrics.sent_latency_seconds.clone());
1781 self.pending_sends.push(Box::pin(async move {
1782 let sent_latency_timer = sent_latency.map(|metric| metric.start_timer());
1783 let result = send_with_timeout(message_sink.send_async_notification(encoded)).await;
1784 drop(sent_latency_timer);
1785 PendingSendResult {
1786 peer: who,
1787 statement_count,
1788 bytes_sent,
1789 result,
1790 kind: SendKind::Propagation,
1791 chunk_id,
1792 }
1793 }));
1794 return;
1795 }
1796 }
1797
1798 fn handle_send_result(&mut self, send_result: PendingSendResult) {
1799 let peer = send_result.peer;
1800 let slot_freed = self.process_send_result(send_result);
1801 self.fill_parked_propagations();
1802 if slot_freed {
1803 self.try_send_next_chunk(peer);
1804 }
1805 }
1806
1807 fn process_send_result(&mut self, send_result: PendingSendResult) -> bool {
1809 let PendingSendResult { peer, statement_count, bytes_sent, result, kind, chunk_id } =
1810 send_result;
1811
1812 let kind_label = kind.label();
1813 match kind {
1814 SendKind::Propagation => {
1815 debug_assert!(
1816 self.propagation_in_flight_bytes >= bytes_sent,
1817 "propagation in-flight byte counter underflow"
1818 );
1819 let in_flight = self.propagation_in_flight_bytes.saturating_sub(bytes_sent);
1820 self.set_propagation_in_flight_bytes(in_flight);
1821 },
1822 SendKind::InitialSync { .. } => {
1823 debug_assert!(
1824 self.initial_sync_in_flight_bytes >= bytes_sent,
1825 "initial-sync in-flight byte counter underflow"
1826 );
1827 let in_flight = self.initial_sync_in_flight_bytes.saturating_sub(bytes_sent);
1828 self.set_initial_sync_in_flight_bytes(in_flight);
1829 },
1830 }
1831
1832 let failure = match result {
1833 SendOutcome::Sent => {
1834 log::trace!(target: LOG_TARGET, "Sent {} statements to {}", statement_count, peer);
1835 self.metrics.as_ref().map(|metrics| {
1836 metrics.propagated_statements.inc_by(statement_count as u64);
1837 metrics.bytes_sent_total.inc_by(bytes_sent);
1838 metrics
1839 .propagated_statements_chunks
1840 .with_label_values(&[kind_label])
1841 .observe(statement_count as f64);
1842 });
1843 None
1844 },
1845 SendOutcome::NetworkError(error) => {
1846 log::debug!(
1847 target: LOG_TARGET,
1848 "Failed to send {statement_count} statements ({bytes_sent} bytes) to {peer}: {error}",
1849 );
1850 Some(send_failure::NETWORK)
1851 },
1852 SendOutcome::TimedOut => {
1853 log::warn!(
1854 target: LOG_TARGET,
1855 "Send of {statement_count} statements ({bytes_sent} bytes) to {peer} timed out after {SEND_TIMEOUT:?}",
1856 );
1857 Some(send_failure::TIMEOUT)
1858 },
1859 };
1860
1861 if let Some(reason) = failure {
1862 match kind {
1863 SendKind::Propagation => self.record_abandoned_send(reason, statement_count),
1864 SendKind::InitialSync { .. } => self.record_send_failure(reason),
1865 }
1866 }
1867
1868 let slot_freed = self.in_flight_chunks.get(&peer) == Some(&chunk_id);
1871 if slot_freed {
1872 self.in_flight_chunks.remove(&peer);
1873 }
1874
1875 let SendKind::InitialSync { sync_id, next_cursor } = kind else { return slot_freed };
1876
1877 if self.pending_initial_syncs.get(&peer).map(|pending| pending.sync_id) != Some(sync_id) {
1880 return slot_freed;
1881 }
1882
1883 if failure.is_some() {
1884 self.initial_sync_peer_queue.push_back(peer);
1886 return slot_freed;
1887 }
1888
1889 if let Some(pending) = self.pending_initial_syncs.get_mut(&peer) {
1890 pending.cursor = next_cursor;
1891 }
1892 self.metrics.as_ref().map(|metrics| {
1893 metrics.initial_sync_statements_sent.inc_by(statement_count as u64);
1894 });
1895 self.initial_sync_peer_queue.push_back(peer);
1899 slot_freed
1900 }
1901
1902 #[cfg(test)]
1903 async fn flush_pending_sends(&mut self) {
1904 while let Some(result) = self.pending_sends.next().await {
1905 self.handle_send_result(result);
1906 }
1907 }
1908
1909 fn do_propagate_statements(&mut self, statements: &[(u64, Hash, Statement)]) {
1910 log::debug!(target: LOG_TARGET, "Propagating {} statements for {} peers", statements.len(), self.peers.len());
1911 let peers: Vec<_> = self.peers.keys().copied().collect();
1912 for who in peers {
1913 log::trace!(target: LOG_TARGET, "Start propagating statements for {}", who);
1914 self.queue_statements_for_peer(&who, statements);
1915 }
1916 log::trace!(target: LOG_TARGET, "Statements queued for propagation to all peers");
1917 }
1918
1919 async fn propagate_statements(&mut self) {
1921 if self.sync.is_major_syncing() {
1923 return;
1924 }
1925
1926 self.fill_parked_propagations();
1929
1930 let Ok(statements) = self.statement_store.take_recent_statements() else { return };
1931 if !statements.is_empty() {
1932 self.do_propagate_statements(&statements);
1933 }
1934 self.recently_received_statements.clear();
1938 }
1939
1940 fn schedule_initial_sync_for_peer(&mut self, peer: PeerId) {
1946 if !self.peers.contains_key(&peer) {
1948 return;
1949 }
1950 let watermark = match self.statement_store.admission_watermark() {
1953 Ok(watermark) => watermark,
1954 Err(e) => {
1955 log::warn!(
1956 target: LOG_TARGET,
1957 "Failed to read the admission watermark, skipping initial sync for {peer}: {e:?}",
1958 );
1959 return;
1960 },
1961 };
1962 let sync_id = self.next_initial_sync_id;
1963 self.next_initial_sync_id = self.next_initial_sync_id.saturating_add(1);
1964 if let Some(pending) = self.pending_initial_syncs.remove(&peer) {
1965 self.record_initial_sync_completion(sync_outcome::ABANDONED, pending.started_at);
1966 self.initial_sync_peer_queue.retain(|p| *p != peer);
1967 }
1968 if watermark > 0 {
1969 if let Some(peer_data) = self.peers.get_mut(&peer) {
1970 peer_data.sync_watermark = peer_data.sync_watermark.max(watermark);
1971 }
1972 self.propagation_outboxes.remove(&peer);
1976 self.pending_initial_syncs.insert(
1977 peer,
1978 PendingInitialSync { cursor: 0, watermark, started_at: Instant::now(), sync_id },
1979 );
1980 self.initial_sync_peer_queue.push_back(peer);
1981 self.metrics.as_ref().map(|metrics| {
1982 metrics.initial_sync_peers_active.inc();
1983 });
1984 }
1985 }
1986
1987 fn process_pending_affinities(&mut self) {
1993 let ready_peers: Vec<PeerId> = self
1994 .peers
1995 .iter()
1996 .filter(|(peer_id, peer_data)| {
1997 peer_data.pending_topic_affinity.is_some() &&
1998 !self.pending_initial_syncs.contains_key(peer_id)
1999 })
2000 .map(|(peer_id, _)| *peer_id)
2001 .collect();
2002
2003 for peer_id in ready_peers {
2004 if let Some(peer_data) = self.peers.get_mut(&peer_id) {
2005 peer_data.topic_affinity = peer_data.pending_topic_affinity.take();
2006 }
2007 self.schedule_initial_sync_for_peer(peer_id);
2008 }
2009 }
2010
2011 fn set_initial_sync_in_flight_bytes(&mut self, bytes: u64) {
2013 self.initial_sync_in_flight_bytes = bytes;
2014 self.metrics
2015 .as_ref()
2016 .map(|metrics| metrics.initial_sync_in_flight_bytes.set(bytes));
2017 }
2018
2019 fn set_propagation_in_flight_bytes(&mut self, bytes: u64) {
2021 self.propagation_in_flight_bytes = bytes;
2022 self.metrics
2023 .as_ref()
2024 .map(|metrics| metrics.propagation_in_flight_bytes.set(bytes));
2025 }
2026
2027 fn send_in_flight_bytes(&self) -> u64 {
2030 self.initial_sync_in_flight_bytes
2031 .saturating_add(self.propagation_in_flight_bytes)
2032 }
2033
2034 fn propagation_send_budget(&self) -> u64 {
2041 if self.pending_initial_syncs.is_empty() {
2042 MAX_SEND_IN_FLIGHT_BYTES
2043 } else {
2044 MAX_SEND_IN_FLIGHT_BYTES - INITIAL_SYNC_RESERVED_BYTES
2045 }
2046 }
2047
2048 fn fill_parked_propagations(&mut self) {
2055 for _ in 0..self.parked_propagations.len() {
2056 if self.send_in_flight_bytes() >= self.propagation_send_budget() {
2057 return;
2058 }
2059 let Some(peer) = self.parked_propagations.pop_front() else { return };
2060 self.try_send_next_chunk(peer);
2061 }
2062 }
2063
2064 fn occupy_send_slot(&mut self, peer: PeerId) -> u64 {
2066 let chunk_id = self.next_chunk_id;
2067 self.next_chunk_id = self.next_chunk_id.saturating_add(1);
2068 self.in_flight_chunks.insert(peer, chunk_id);
2069 chunk_id
2070 }
2071
2072 fn record_initial_sync_completion(&self, outcome: &str, started_at: Instant) {
2074 self.metrics.as_ref().map(|metrics| {
2075 metrics.initial_sync_peers_active.dec();
2076 metrics
2077 .initial_sync_duration_seconds
2078 .with_label_values(&[outcome])
2079 .observe(started_at.elapsed().as_secs_f64());
2080 });
2081 }
2082
2083 fn process_initial_sync_burst(&mut self) {
2085 if self.sync.is_major_syncing() {
2086 return;
2087 }
2088
2089 if self.send_in_flight_bytes() >= MAX_SEND_IN_FLIGHT_BYTES {
2090 log::debug!(
2091 target: LOG_TARGET,
2092 "Skipping initial sync burst, {} bytes still in flight",
2093 self.send_in_flight_bytes(),
2094 );
2095 return;
2096 }
2097
2098 let Some(pos) = self
2101 .initial_sync_peer_queue
2102 .iter()
2103 .position(|peer| !self.in_flight_chunks.contains_key(peer))
2104 else {
2105 return;
2106 };
2107 self.initial_sync_peer_queue.rotate_left(pos);
2108 let Some(peer_id) = self.initial_sync_peer_queue.pop_front() else {
2109 return;
2110 };
2111
2112 let Entry::Occupied(mut entry) = self.pending_initial_syncs.entry(peer_id) else {
2113 return;
2114 };
2115 let sync_id = entry.get().sync_id;
2116
2117 self.metrics.as_ref().map(|metrics| {
2118 metrics.initial_sync_bursts_total.inc();
2119 });
2120
2121 if entry.get().cursor >= entry.get().watermark {
2122 let started_at = entry.get().started_at;
2123 entry.remove();
2124 self.record_initial_sync_completion(sync_outcome::COMPLETED, started_at);
2125 return;
2126 }
2127
2128 let Some(peer_data) = self.peers.get(&peer_id) else {
2131 log::error!(target: LOG_TARGET, "Peer {peer_id} has pending initial sync but is not in peers map");
2132 let pending = entry.remove();
2133 self.record_initial_sync_completion(sync_outcome::ABANDONED, pending.started_at);
2134 return;
2135 };
2136 let peer_version = peer_data.protocol_version;
2137 let envelope_overhead = peer_version.envelope_overhead();
2138 let max_size = max_statement_payload_size(envelope_overhead);
2139 let (batch, accumulated_size) = match fetch_admitted_chunk(
2140 &*self.statement_store,
2141 &self.recently_received_statements,
2142 &self.pending_statements_peers,
2143 &peer_id,
2144 peer_data,
2145 entry.get().cursor,
2146 entry.get().watermark,
2147 max_size,
2148 ) {
2149 Ok(r) => r,
2150 Err(e) => {
2151 log::warn!(
2154 target: LOG_TARGET,
2155 "Failed to fetch statements for initial sync of {peer_id}, will retry: {e:?}",
2156 );
2157 self.initial_sync_peer_queue.push_back(peer_id);
2158 return;
2159 },
2160 };
2161
2162 if accumulated_size > max_size {
2165 log::warn!(target: LOG_TARGET, "Statement too large, skipping");
2166 self.metrics.as_ref().map(|metrics| {
2167 metrics.skipped_oversized_statements.inc();
2168 });
2169 entry.get_mut().cursor = batch.cursor;
2170 self.initial_sync_peer_queue.push_back(peer_id);
2171 return;
2172 }
2173
2174 if batch.statements.is_empty() {
2175 entry.get_mut().cursor = batch.cursor;
2178 self.initial_sync_peer_queue.push_back(peer_id);
2179 return;
2180 }
2181
2182 let next_cursor = batch.cursor;
2183
2184 let statement_count = batch.statements.len();
2185 let send_stmts: Vec<_> = batch.statements.iter().map(|(_, stmt)| stmt).collect();
2186 let encoded = match peer_version {
2187 PeerProtocolVersion::V1 => send_stmts.encode(),
2188 PeerProtocolVersion::V2 => StatementMessage::encode_statement_refs(&send_stmts),
2189 };
2190 let bytes_to_send = encoded.len() as u64;
2191 let Some(message_sink) = self.notification_service.message_sink(&peer_id) else {
2192 log::debug!(
2195 target: LOG_TARGET,
2196 "Failed to get message sink for peer {peer_id}, its initial sync will retry",
2197 );
2198 self.record_send_failure(send_failure::NO_SINK);
2199 self.initial_sync_peer_queue.push_back(peer_id);
2200 return;
2201 };
2202 let sent_latency =
2203 self.metrics.as_ref().map(|metrics| metrics.sent_latency_seconds.clone());
2204 let in_flight = self.initial_sync_in_flight_bytes.saturating_add(bytes_to_send);
2205 self.set_initial_sync_in_flight_bytes(in_flight);
2206 let chunk_id = self.occupy_send_slot(peer_id);
2207 self.pending_sends.push(Box::pin(async move {
2208 let sent_latency_timer = sent_latency.map(|metric| metric.start_timer());
2209 let result = send_with_timeout(message_sink.send_async_notification(encoded)).await;
2210 drop(sent_latency_timer);
2211 PendingSendResult {
2212 peer: peer_id,
2213 statement_count,
2214 bytes_sent: bytes_to_send,
2215 result,
2216 kind: SendKind::InitialSync { sync_id, next_cursor },
2217 chunk_id,
2218 }
2219 }));
2220 }
2221}
2222
2223#[cfg(test)]
2224mod tests {
2225
2226 use super::*;
2227 use governor::clock::FakeRelativeClock;
2228 use std::{
2229 sync::{
2230 atomic::{AtomicBool, AtomicUsize, Ordering},
2231 Mutex,
2232 },
2233 time::Duration,
2234 };
2235
2236 const BLOOM_SEED: u128 = 0x5EED_5EED_5EED_5EED;
2238
2239 fn new_live_statement() -> Statement {
2240 let mut statement = sp_statement_store::Statement::new();
2241 statement.set_expiry_from_parts(u32::MAX, 0);
2242 statement
2243 }
2244
2245 #[derive(Clone)]
2246 struct TestNetwork {
2247 reported_peers: Arc<Mutex<Vec<(PeerId, sc_network::ReputationChange)>>>,
2248 disconnected_peers: Arc<Mutex<Vec<PeerId>>>,
2249 default_role: sc_network::ObservedRole,
2251 added_reserved: Arc<Mutex<Vec<HashSet<sc_network::Multiaddr>>>>,
2252 removed_reserved: Arc<Mutex<Vec<Vec<PeerId>>>>,
2253 }
2254
2255 impl TestNetwork {
2256 fn new() -> Self {
2257 Self {
2258 reported_peers: Arc::new(Mutex::new(Vec::new())),
2259 disconnected_peers: Arc::new(Mutex::new(Vec::new())),
2260 default_role: sc_network::ObservedRole::Full,
2261 added_reserved: Arc::new(Mutex::new(Vec::new())),
2262 removed_reserved: Arc::new(Mutex::new(Vec::new())),
2263 }
2264 }
2265
2266 fn new_light() -> Self {
2267 Self {
2268 reported_peers: Arc::new(Mutex::new(Vec::new())),
2269 disconnected_peers: Arc::new(Mutex::new(Vec::new())),
2270 default_role: sc_network::ObservedRole::Light,
2271 added_reserved: Arc::new(Mutex::new(Vec::new())),
2272 removed_reserved: Arc::new(Mutex::new(Vec::new())),
2273 }
2274 }
2275
2276 fn get_reports(&self) -> Vec<(PeerId, sc_network::ReputationChange)> {
2277 self.reported_peers.lock().unwrap().clone()
2278 }
2279
2280 fn get_disconnected_peers(&self) -> Vec<PeerId> {
2281 self.disconnected_peers.lock().unwrap().clone()
2282 }
2283
2284 fn get_added_reserved(&self) -> Vec<HashSet<sc_network::Multiaddr>> {
2285 self.added_reserved.lock().unwrap().clone()
2286 }
2287
2288 fn get_removed_reserved(&self) -> Vec<Vec<PeerId>> {
2289 self.removed_reserved.lock().unwrap().clone()
2290 }
2291 }
2292
2293 #[async_trait::async_trait]
2294 impl NetworkPeers for TestNetwork {
2295 fn set_authorized_peers(&self, _: std::collections::HashSet<PeerId>) {
2296 unimplemented!()
2297 }
2298
2299 fn set_authorized_only(&self, _: bool) {
2300 unimplemented!()
2301 }
2302
2303 fn add_known_address(&self, _: PeerId, _: sc_network::Multiaddr) {
2304 unimplemented!()
2305 }
2306
2307 fn report_peer(&self, peer_id: PeerId, cost_benefit: sc_network::ReputationChange) {
2308 self.reported_peers.lock().unwrap().push((peer_id, cost_benefit));
2309 }
2310
2311 fn peer_reputation(&self, _: &PeerId) -> i32 {
2312 unimplemented!()
2313 }
2314
2315 fn disconnect_peer(&self, peer: PeerId, _: sc_network::ProtocolName) {
2316 self.disconnected_peers.lock().unwrap().push(peer);
2317 }
2318
2319 fn accept_unreserved_peers(&self) {
2320 unimplemented!()
2321 }
2322
2323 fn deny_unreserved_peers(&self) {
2324 unimplemented!()
2325 }
2326
2327 fn add_reserved_peer(
2328 &self,
2329 _: sc_network::config::MultiaddrWithPeerId,
2330 ) -> Result<(), String> {
2331 unimplemented!()
2332 }
2333
2334 fn remove_reserved_peer(&self, _: PeerId) {
2335 unimplemented!()
2336 }
2337
2338 fn set_reserved_peers(
2339 &self,
2340 _: sc_network::ProtocolName,
2341 _: std::collections::HashSet<sc_network::Multiaddr>,
2342 ) -> Result<(), String> {
2343 unimplemented!()
2344 }
2345
2346 fn add_peers_to_reserved_set(
2347 &self,
2348 _: sc_network::ProtocolName,
2349 addrs: std::collections::HashSet<sc_network::Multiaddr>,
2350 ) -> Result<(), String> {
2351 self.added_reserved.lock().unwrap().push(addrs);
2352 Ok(())
2353 }
2354
2355 fn remove_peers_from_reserved_set(
2356 &self,
2357 _: sc_network::ProtocolName,
2358 peers: Vec<PeerId>,
2359 ) -> Result<(), String> {
2360 self.removed_reserved.lock().unwrap().push(peers);
2361 Ok(())
2362 }
2363
2364 fn sync_num_connected(&self) -> usize {
2365 unimplemented!()
2366 }
2367
2368 fn peer_role(&self, _: PeerId, _: Vec<u8>) -> Option<sc_network::ObservedRole> {
2369 Some(self.default_role)
2370 }
2371
2372 async fn reserved_peers(&self) -> Result<Vec<PeerId>, ()> {
2373 unimplemented!();
2374 }
2375 }
2376
2377 #[derive(Clone)]
2378 struct TestSync {
2379 major_syncing: Arc<AtomicBool>,
2380 }
2381
2382 impl TestSync {
2383 fn new() -> Self {
2384 Self { major_syncing: Arc::new(AtomicBool::new(false)) }
2385 }
2386
2387 fn with_syncing(initial: bool) -> (Self, Arc<AtomicBool>) {
2388 let flag = Arc::new(AtomicBool::new(initial));
2389 (Self { major_syncing: flag.clone() }, flag)
2390 }
2391 }
2392
2393 impl SyncEventStream for TestSync {
2394 fn event_stream(
2395 &self,
2396 _name: &'static str,
2397 ) -> Pin<Box<dyn Stream<Item = sc_network_sync::types::SyncEvent> + Send>> {
2398 Box::pin(futures::stream::pending())
2399 }
2400 }
2401
2402 impl sp_consensus::SyncOracle for TestSync {
2403 fn is_major_syncing(&self) -> bool {
2404 self.major_syncing.load(Ordering::Relaxed)
2405 }
2406
2407 fn is_offline(&self) -> bool {
2408 unimplemented!()
2409 }
2410 }
2411
2412 impl NetworkEventStream for TestNetwork {
2413 fn event_stream(
2414 &self,
2415 _name: &'static str,
2416 ) -> Pin<Box<dyn Stream<Item = sc_network::Event> + Send>> {
2417 unimplemented!()
2418 }
2419 }
2420
2421 #[derive(Debug, Clone)]
2422 struct TestNotificationService {
2423 sent_notifications: Arc<Mutex<Vec<(PeerId, Vec<u8>)>>>,
2424 block_sends: Arc<AtomicBool>,
2425 fail_sends: Arc<AtomicBool>,
2426 sinks_available: Arc<AtomicUsize>,
2427 }
2428
2429 impl TestNotificationService {
2430 fn new() -> Self {
2431 Self {
2432 sent_notifications: Arc::new(Mutex::new(Vec::new())),
2433 block_sends: Arc::new(AtomicBool::new(false)),
2434 fail_sends: Arc::new(AtomicBool::new(false)),
2435 sinks_available: Arc::new(AtomicUsize::new(usize::MAX)),
2436 }
2437 }
2438
2439 fn get_sent_notifications(&self) -> Vec<(PeerId, Vec<u8>)> {
2440 self.sent_notifications.lock().unwrap().clone()
2441 }
2442
2443 fn clear_sent_notifications(&self) {
2444 self.sent_notifications.lock().unwrap().clear();
2445 }
2446
2447 fn block_sends(&self) {
2448 self.block_sends.store(true, Ordering::Relaxed);
2449 }
2450
2451 fn fail_sends(&self) {
2452 self.fail_sends.store(true, Ordering::Relaxed);
2453 }
2454
2455 fn allow_sends(&self) {
2456 self.fail_sends.store(false, Ordering::Relaxed);
2457 }
2458
2459 fn serve_sinks(&self, count: usize) {
2460 self.sinks_available.store(count, Ordering::Relaxed);
2461 }
2462 }
2463
2464 struct TestMessageSink {
2465 peer: PeerId,
2466 sent_notifications: Arc<Mutex<Vec<(PeerId, Vec<u8>)>>>,
2467 block_sends: Arc<AtomicBool>,
2468 fail_sends: Arc<AtomicBool>,
2469 }
2470
2471 #[async_trait::async_trait]
2472 impl sc_network::service::traits::MessageSink for TestMessageSink {
2473 fn send_sync_notification(&self, notification: Vec<u8>) {
2474 self.sent_notifications.lock().unwrap().push((self.peer, notification));
2475 }
2476
2477 async fn send_async_notification(
2478 &self,
2479 notification: Vec<u8>,
2480 ) -> Result<(), sc_network::error::Error> {
2481 if self.block_sends.load(Ordering::Relaxed) {
2482 futures::future::pending::<()>().await;
2483 }
2484 if self.fail_sends.load(Ordering::Relaxed) {
2485 return Err(sc_network::error::Error::ConnectionClosed);
2486 }
2487 self.sent_notifications.lock().unwrap().push((self.peer, notification));
2488 Ok(())
2489 }
2490 }
2491
2492 #[async_trait::async_trait]
2493 impl NotificationService for TestNotificationService {
2494 async fn open_substream(&mut self, _peer: PeerId) -> Result<(), ()> {
2495 unimplemented!()
2496 }
2497
2498 async fn close_substream(&mut self, _peer: PeerId) -> Result<(), ()> {
2499 unimplemented!()
2500 }
2501
2502 fn send_sync_notification(&mut self, peer: &PeerId, notification: Vec<u8>) {
2503 self.sent_notifications.lock().unwrap().push((*peer, notification));
2504 }
2505
2506 async fn send_async_notification(
2507 &mut self,
2508 peer: &PeerId,
2509 notification: Vec<u8>,
2510 ) -> Result<(), sc_network::error::Error> {
2511 if self.fail_sends.load(Ordering::Relaxed) {
2512 return Err(sc_network::error::Error::ConnectionClosed);
2513 }
2514 self.sent_notifications.lock().unwrap().push((*peer, notification));
2515 Ok(())
2516 }
2517
2518 async fn set_handshake(&mut self, _handshake: Vec<u8>) -> Result<(), ()> {
2519 unimplemented!()
2520 }
2521
2522 fn try_set_handshake(&mut self, _handshake: Vec<u8>) -> Result<(), ()> {
2523 unimplemented!()
2524 }
2525
2526 async fn next_event(&mut self) -> Option<sc_network::service::traits::NotificationEvent> {
2527 None
2528 }
2529
2530 fn clone(&mut self) -> Result<Box<dyn NotificationService>, ()> {
2531 unimplemented!()
2532 }
2533
2534 fn protocol(&self) -> &sc_network::types::ProtocolName {
2535 unimplemented!()
2536 }
2537
2538 fn message_sink(
2539 &self,
2540 peer: &PeerId,
2541 ) -> Option<Box<dyn sc_network::service::traits::MessageSink>> {
2542 self.sinks_available
2543 .fetch_update(Ordering::Relaxed, Ordering::Relaxed, |n| n.checked_sub(1))
2544 .ok()?;
2545 Some(Box::new(TestMessageSink {
2546 peer: *peer,
2547 sent_notifications: self.sent_notifications.clone(),
2548 block_sends: self.block_sends.clone(),
2549 fail_sends: self.fail_sends.clone(),
2550 }))
2551 }
2552 }
2553
2554 #[derive(Clone)]
2555 struct TestStatementStore {
2556 statements: Arc<Mutex<HashMap<sp_statement_store::Hash, sp_statement_store::Statement>>>,
2557 recent_statements:
2558 Arc<Mutex<HashMap<sp_statement_store::Hash, sp_statement_store::Statement>>>,
2559 admissions: Arc<Mutex<Vec<sp_statement_store::Hash>>>,
2561 fail_fetches: Arc<AtomicBool>,
2562 }
2563
2564 impl TestStatementStore {
2565 fn new() -> Self {
2566 Self {
2567 statements: Default::default(),
2568 recent_statements: Default::default(),
2569 admissions: Default::default(),
2570 fail_fetches: Arc::new(AtomicBool::new(false)),
2571 }
2572 }
2573
2574 fn insert(&self, statement: sp_statement_store::Statement) {
2576 let hash = statement.hash();
2577 self.statements.lock().unwrap().insert(hash, statement);
2578 self.admit(hash);
2579 }
2580
2581 fn admit(&self, hash: sp_statement_store::Hash) -> u64 {
2584 let mut admissions = self.admissions.lock().unwrap();
2585 if let Some(seq) = admissions.iter().position(|admitted| *admitted == hash) {
2586 return seq as u64;
2587 }
2588 admissions.push(hash);
2589 (admissions.len() - 1) as u64
2590 }
2591 }
2592
2593 impl StatementStore for TestStatementStore {
2594 fn statements(
2595 &self,
2596 ) -> sp_statement_store::Result<
2597 Vec<(sp_statement_store::Hash, sp_statement_store::Statement)>,
2598 > {
2599 Ok(self.statements.lock().unwrap().iter().map(|(h, s)| (*h, s.clone())).collect())
2600 }
2601
2602 fn take_recent_statements(
2603 &self,
2604 ) -> sp_statement_store::Result<
2605 Vec<(u64, sp_statement_store::Hash, sp_statement_store::Statement)>,
2606 > {
2607 let drained: Vec<_> = self.recent_statements.lock().unwrap().drain().collect();
2610 let mut statements = self.statements.lock().unwrap();
2611 for (hash, statement) in &drained {
2612 statements.insert(*hash, statement.clone());
2613 }
2614 drop(statements);
2615 let mut result: Vec<_> = drained
2616 .into_iter()
2617 .map(|(hash, statement)| (self.admit(hash), hash, statement))
2618 .collect();
2619 result.sort_unstable_by_key(|(seq, ..)| *seq);
2620 Ok(result)
2621 }
2622
2623 fn statement(
2624 &self,
2625 _hash: &sp_statement_store::Hash,
2626 ) -> sp_statement_store::Result<Option<sp_statement_store::Statement>> {
2627 unimplemented!()
2628 }
2629
2630 fn has_statement(&self, hash: &sp_statement_store::Hash) -> bool {
2631 self.statements.lock().unwrap().contains_key(hash)
2632 }
2633
2634 fn statements_by_hashes(
2635 &self,
2636 hashes: &[sp_statement_store::Hash],
2637 filter: &mut dyn FnMut(
2638 &sp_statement_store::Hash,
2639 &[u8],
2640 &sp_statement_store::Statement,
2641 ) -> FilterDecision,
2642 ) -> sp_statement_store::Result<(
2643 Vec<(sp_statement_store::Hash, sp_statement_store::Statement)>,
2644 usize,
2645 )> {
2646 if self.fail_fetches.load(Ordering::Relaxed) {
2647 return Err(sp_statement_store::Error::Db("fetch failed".into()));
2648 }
2649 let statements = self.statements.lock().unwrap();
2650 let mut result = Vec::new();
2651 let mut processed = 0;
2652 for hash in hashes {
2653 let Some(stmt) = statements.get(hash) else {
2654 processed += 1;
2655 continue;
2656 };
2657 let encoded = stmt.encode();
2658 match filter(hash, &encoded, stmt) {
2659 FilterDecision::Skip => {
2660 processed += 1;
2661 },
2662 FilterDecision::Take => {
2663 processed += 1;
2664 result.push((*hash, stmt.clone()));
2665 },
2666 FilterDecision::Abort => break,
2667 }
2668 }
2669 Ok((result, processed))
2670 }
2671
2672 fn admission_watermark(&self) -> sp_statement_store::Result<u64> {
2673 Ok(self.admissions.lock().unwrap().len() as u64)
2674 }
2675
2676 fn admitted_statements(
2677 &self,
2678 mut cursor: u64,
2679 watermark: u64,
2680 scan_limit: usize,
2681 filter: &mut dyn FnMut(
2682 &sp_statement_store::Hash,
2683 &[u8],
2684 &sp_statement_store::Statement,
2685 ) -> FilterDecision,
2686 ) -> sp_statement_store::Result<sp_statement_store::AdmittedBatch> {
2687 if self.fail_fetches.load(Ordering::Relaxed) {
2688 return Err(sp_statement_store::Error::Db("fetch failed".into()));
2689 }
2690 let admissions = self.admissions.lock().unwrap();
2691 let statements = self.statements.lock().unwrap();
2692 let mut result = Vec::new();
2693 let mut aborted = false;
2694 let mut scanned = 0usize;
2695 while cursor < watermark {
2696 if scanned == scan_limit {
2697 aborted = true;
2698 break;
2699 }
2700 scanned += 1;
2701 let Some(hash) = admissions.get(cursor as usize) else { break };
2702 let Some(statement) = statements.get(hash) else {
2704 cursor += 1;
2705 continue;
2706 };
2707 let encoded = statement.encode();
2708 match filter(hash, &encoded, statement) {
2709 FilterDecision::Skip => cursor += 1,
2710 FilterDecision::Take => {
2711 result.push((*hash, statement.clone()));
2712 cursor += 1;
2713 },
2714 FilterDecision::Abort => {
2715 aborted = true;
2716 break;
2717 },
2718 }
2719 }
2720 if !aborted && cursor < watermark {
2721 cursor = watermark;
2722 }
2723 Ok(sp_statement_store::AdmittedBatch {
2724 statements: result,
2725 cursor,
2726 done: cursor >= watermark,
2727 })
2728 }
2729
2730 fn broadcasts(
2731 &self,
2732 _match_all_topics: &[sp_statement_store::Topic],
2733 ) -> sp_statement_store::Result<Vec<Vec<u8>>> {
2734 unimplemented!()
2735 }
2736
2737 fn posted(
2738 &self,
2739 _match_all_topics: &[sp_statement_store::Topic],
2740 _dest: [u8; 32],
2741 ) -> sp_statement_store::Result<Vec<Vec<u8>>> {
2742 unimplemented!()
2743 }
2744
2745 fn posted_clear(
2746 &self,
2747 _match_all_topics: &[sp_statement_store::Topic],
2748 _dest: [u8; 32],
2749 ) -> sp_statement_store::Result<Vec<Vec<u8>>> {
2750 unimplemented!()
2751 }
2752
2753 fn broadcasts_stmt(
2754 &self,
2755 _match_all_topics: &[sp_statement_store::Topic],
2756 ) -> sp_statement_store::Result<Vec<Vec<u8>>> {
2757 unimplemented!()
2758 }
2759
2760 fn posted_stmt(
2761 &self,
2762 _match_all_topics: &[sp_statement_store::Topic],
2763 _dest: [u8; 32],
2764 ) -> sp_statement_store::Result<Vec<Vec<u8>>> {
2765 unimplemented!()
2766 }
2767
2768 fn posted_clear_stmt(
2769 &self,
2770 _match_all_topics: &[sp_statement_store::Topic],
2771 _dest: [u8; 32],
2772 ) -> sp_statement_store::Result<Vec<Vec<u8>>> {
2773 unimplemented!()
2774 }
2775
2776 fn submit(
2777 &self,
2778 _statement: sp_statement_store::Statement,
2779 _source: sp_statement_store::StatementSource,
2780 ) -> sp_statement_store::SubmitResult {
2781 unimplemented!()
2782 }
2783
2784 fn remove(&self, _hash: &sp_statement_store::Hash) -> sp_statement_store::Result<()> {
2785 unimplemented!()
2786 }
2787
2788 fn remove_by(&self, _who: [u8; 32]) -> sp_statement_store::Result<()> {
2789 unimplemented!()
2790 }
2791 }
2792
2793 fn build_handler(
2794 num_peers: usize,
2795 ) -> (
2796 StatementHandler<TestNetwork, TestSync>,
2797 TestStatementStore,
2798 TestNetwork,
2799 TestNotificationService,
2800 async_channel::Receiver<(Statement, oneshot::Sender<SubmitResult>)>,
2801 Vec<PeerId>,
2802 ) {
2803 let statement_store = TestStatementStore::new();
2804 let (queue_sender, queue_receiver) = async_channel::bounded(100);
2805 let network = TestNetwork::new();
2806 let notification_service = TestNotificationService::new();
2807 let mut peers = HashMap::new();
2808 let mut peer_ids = Vec::with_capacity(num_peers);
2809
2810 for _ in 0..num_peers {
2811 let peer_id = PeerId::random();
2812 peer_ids.push(peer_id);
2813 peers.insert(
2814 peer_id,
2815 Peer {
2816 rate_limiter: PeerRateLimiter::new(
2817 NonZeroU32::new(DEFAULT_STATEMENTS_PER_SECOND)
2818 .expect("DEFAULT_STATEMENTS_PER_SECOND is nonzero"),
2819 NonZeroU32::new(
2820 DEFAULT_STATEMENTS_PER_SECOND * config::STATEMENTS_BURST_COEFFICIENT,
2821 )
2822 .expect("burst capacity is nonzero"),
2823 ),
2824 protocol_version: PeerProtocolVersion::V1,
2825 topic_affinity: None,
2826 is_light: false,
2827 pending_topic_affinity: None,
2828 sync_watermark: 0,
2829 },
2830 );
2831 }
2832
2833 let handler = StatementHandler {
2834 protocol_name: format!("/{STATEMENT_PROTOCOL_V1}").into(),
2835 notification_service: Box::new(notification_service.clone()),
2836 propagate_timeout: (Box::pin(futures::stream::pending())
2837 as Pin<Box<dyn Stream<Item = ()> + Send>>)
2838 .fuse(),
2839 pending_statements: FuturesUnordered::new(),
2840 pending_statements_peers: HashMap::new(),
2841 recently_received_statements: HashMap::new(),
2842 network: network.clone(),
2843 sync: TestSync::new(),
2844 sync_event_stream: (Box::pin(futures::stream::pending())
2845 as Pin<Box<dyn Stream<Item = sc_network_sync::types::SyncEvent> + Send>>)
2846 .fuse(),
2847 peers,
2848 statement_store: Arc::new(statement_store.clone()),
2849 queue_sender,
2850 statements_per_second: NonZeroU32::new(DEFAULT_STATEMENTS_PER_SECOND)
2851 .expect("DEFAULT_STATEMENTS_PER_SECOND is nonzero"),
2852 metrics: None,
2853 initial_sync_timeout: Box::pin(futures::future::pending()),
2854 pending_affinities_timeout: Box::pin(futures::future::pending()),
2855 pending_initial_syncs: HashMap::new(),
2856 initial_sync_peer_queue: VecDeque::new(),
2857 next_initial_sync_id: 0,
2858 initial_sync_in_flight_bytes: 0,
2859 propagation_outboxes: HashMap::new(),
2860 in_flight_chunks: HashMap::new(),
2861 next_chunk_id: 0,
2862 propagation_in_flight_bytes: 0,
2863 parked_propagations: VecDeque::new(),
2864 pending_sends: FuturesUnordered::new(),
2865 deferred_peers: HashSet::new(),
2866 dropped_statements_during_sync: false,
2867 sync_recovery_peer: None,
2868 sync_recovery_readd_timeout: Box::pin(futures::future::pending()),
2869 };
2870 (handler, statement_store, network, notification_service, queue_receiver, peer_ids)
2871 }
2872
2873 fn get_peer_hashes(sent: &[(PeerId, Vec<u8>)], peer: PeerId) -> Vec<sp_statement_store::Hash> {
2874 sent.iter()
2875 .filter(|(p, _)| *p == peer)
2876 .flat_map(|(_, notification)| {
2877 <Statements as Decode>::decode(&mut notification.as_slice()).unwrap()
2878 })
2879 .map(|s| s.hash())
2880 .collect()
2881 }
2882
2883 async fn import_queued_statement(
2886 handler: &mut StatementHandler<TestNetwork, TestSync>,
2887 statement_store: &TestStatementStore,
2888 queue_receiver: &async_channel::Receiver<(Statement, oneshot::Sender<SubmitResult>)>,
2889 result: SubmitResult,
2890 ) {
2891 let (statement, completion) = queue_receiver.try_recv().unwrap();
2892 let hash = statement.hash();
2893 statement_store.insert(statement.clone());
2894 statement_store.recent_statements.lock().unwrap().insert(hash, statement);
2895 completion.send(result).unwrap();
2896 let (hash, result) = handler.pending_statements.next().await.unwrap();
2897 handler.on_statement_submit_result(hash, result);
2898 }
2899
2900 #[tokio::test]
2901 async fn statement_is_not_sent_back_to_the_peers_it_came_from() {
2902 let (
2903 mut handler,
2904 statement_store,
2905 _network,
2906 notification_service,
2907 queue_receiver,
2908 peer_ids,
2909 ) = build_handler(3);
2910 let (sender_a, sender_b, receiver) = (peer_ids[0], peer_ids[1], peer_ids[2]);
2911
2912 let mut statement = new_live_statement();
2913 statement.set_plain_data(b"statement from two peers".to_vec());
2914 let hash = statement.hash();
2915
2916 handler.on_statements(sender_a, vec![statement.clone()]);
2918 handler.on_statements(sender_b, vec![statement.clone()]);
2919 import_queued_statement(&mut handler, &statement_store, &queue_receiver, SubmitResult::New)
2920 .await;
2921
2922 handler.propagate_statements().await;
2923 handler.flush_pending_sends().await;
2924
2925 let sent = notification_service.get_sent_notifications();
2926 assert!(get_peer_hashes(&sent, sender_a).is_empty(), "statement returned to sender_a");
2927 assert!(get_peer_hashes(&sent, sender_b).is_empty(), "statement returned to sender_b");
2928 assert_eq!(get_peer_hashes(&sent, receiver), vec![hash]);
2929 assert!(
2930 handler.recently_received_statements.is_empty(),
2931 "recently received statements must be cleared after the propagation pass"
2932 );
2933
2934 handler.on_statements(sender_a, vec![statement]);
2937 assert!(handler.recently_received_statements.is_empty());
2938 }
2939
2940 #[tokio::test]
2941 async fn statement_received_mid_import_is_not_sent_back_to_the_sender() {
2942 let (
2943 mut handler,
2944 statement_store,
2945 _network,
2946 notification_service,
2947 queue_receiver,
2948 peer_ids,
2949 ) = build_handler(3);
2950 let (sender_a, sender_b, receiver) = (peer_ids[0], peer_ids[1], peer_ids[2]);
2951
2952 let mut statement = new_live_statement();
2953 statement.set_plain_data(b"statement received mid-import".to_vec());
2954 let hash = statement.hash();
2955
2956 handler.on_statements(sender_a, vec![statement.clone()]);
2957
2958 let (queued, completion) = queue_receiver.try_recv().unwrap();
2961 statement_store.insert(queued.clone());
2962 statement_store.recent_statements.lock().unwrap().insert(hash, queued);
2963
2964 handler.on_statements(sender_b, vec![statement]);
2966 assert!(handler
2967 .pending_statements_peers
2968 .get(&hash)
2969 .is_some_and(|peers| peers.contains(&sender_b)));
2970
2971 completion.send(SubmitResult::New).unwrap();
2972 let (hash, result) = handler.pending_statements.next().await.unwrap();
2973 handler.on_statement_submit_result(hash, result);
2974
2975 handler.propagate_statements().await;
2976 handler.flush_pending_sends().await;
2977
2978 let sent = notification_service.get_sent_notifications();
2979 assert!(get_peer_hashes(&sent, sender_a).is_empty(), "statement returned to sender_a");
2980 assert!(get_peer_hashes(&sent, sender_b).is_empty(), "statement returned to sender_b");
2981 assert_eq!(get_peer_hashes(&sent, receiver), vec![hash]);
2982 }
2983
2984 #[tokio::test]
2985 async fn statement_forwarded_before_the_tick_is_not_sent_back_to_the_forwarder() {
2986 let (
2987 mut handler,
2988 statement_store,
2989 _network,
2990 notification_service,
2991 queue_receiver,
2992 peer_ids,
2993 ) = build_handler(3);
2994 let (sender, forwarder, receiver) = (peer_ids[0], peer_ids[1], peer_ids[2]);
2995
2996 let mut statement = new_live_statement();
2997 statement.set_plain_data(b"late forwarder".to_vec());
2998 let hash = statement.hash();
2999
3000 handler.on_statements(sender, vec![statement.clone()]);
3001 import_queued_statement(
3004 &mut handler,
3005 &statement_store,
3006 &queue_receiver,
3007 SubmitResult::Known,
3008 )
3009 .await;
3010 handler.on_statements(forwarder, vec![statement]);
3012
3013 handler.propagate_statements().await;
3014 handler.flush_pending_sends().await;
3015
3016 let sent = notification_service.get_sent_notifications();
3017 assert!(
3018 get_peer_hashes(&sent, sender).is_empty(),
3019 "statement returned to the peer that sent it"
3020 );
3021 assert!(get_peer_hashes(&sent, forwarder).is_empty(), "statement returned to forwarder");
3022 assert_eq!(get_peer_hashes(&sent, receiver), vec![hash]);
3023 }
3024
3025 #[tokio::test]
3026 async fn recently_received_statements_survive_a_major_sync_early_return() {
3027 let (
3028 mut handler,
3029 statement_store,
3030 _network,
3031 notification_service,
3032 queue_receiver,
3033 peer_ids,
3034 ) = build_handler(2);
3035 let (sender, receiver) = (peer_ids[0], peer_ids[1]);
3036
3037 let mut statement = new_live_statement();
3038 statement.set_plain_data(b"during major sync".to_vec());
3039 let hash = statement.hash();
3040
3041 handler.on_statements(sender, vec![statement]);
3042 import_queued_statement(&mut handler, &statement_store, &queue_receiver, SubmitResult::New)
3043 .await;
3044
3045 handler.sync.major_syncing.store(true, Ordering::Relaxed);
3048 handler.propagate_statements().await;
3049 assert!(
3050 handler.recently_received_statements.contains_key(&hash),
3051 "entries must survive the major-sync early return"
3052 );
3053
3054 handler.sync.major_syncing.store(false, Ordering::Relaxed);
3055 handler.propagate_statements().await;
3056 handler.flush_pending_sends().await;
3057
3058 let sent = notification_service.get_sent_notifications();
3059 assert!(
3060 get_peer_hashes(&sent, sender).is_empty(),
3061 "statement returned to the peer that sent it"
3062 );
3063 assert_eq!(get_peer_hashes(&sent, receiver), vec![hash]);
3064 assert!(handler.recently_received_statements.is_empty());
3065 }
3066
3067 #[tokio::test]
3068 async fn propagation_does_not_wait_for_pending_send() {
3069 let (mut handler, statement_store, _, notification_service, _, _) = build_handler(1);
3070 let mut statement = new_live_statement();
3071 statement.set_plain_data(b"statement".to_vec());
3072 statement_store
3073 .recent_statements
3074 .lock()
3075 .unwrap()
3076 .insert(statement.hash(), statement);
3077
3078 notification_service.block_sends();
3079 let result =
3080 tokio::time::timeout(Duration::from_secs(1), handler.propagate_statements()).await;
3081
3082 assert!(result.is_ok(), "Propagation waited for a pending send");
3083 assert_eq!(handler.pending_sends.len(), 1);
3084 }
3085
3086 #[tokio::test]
3087 async fn slow_peer_keeps_one_propagation_chunk_in_flight() {
3088 let (mut handler, statement_store, _network, notification_service, _, peer_ids) =
3089 build_handler(1);
3090 let peer_id = peer_ids[0];
3091
3092 for i in 0..25u8 {
3094 let mut statement = new_live_statement();
3095 let mut data = vec![0u8; 100 * 1024];
3096 data[0] = i;
3097 statement.set_plain_data(data);
3098 statement_store
3099 .recent_statements
3100 .lock()
3101 .unwrap()
3102 .insert(statement.hash(), statement);
3103 }
3104
3105 notification_service.block_sends();
3107 handler.propagate_statements().await;
3108
3109 assert_eq!(handler.pending_sends.len(), 1, "only one chunk may be in flight");
3110 let backlog = handler.propagation_outboxes.get(&peer_id).unwrap().len();
3111 assert!(backlog > 0, "the remaining hashes stay in the outbox");
3112
3113 let mut statement = new_live_statement();
3115 statement.set_plain_data(b"second tick".to_vec());
3116 statement_store
3117 .recent_statements
3118 .lock()
3119 .unwrap()
3120 .insert(statement.hash(), statement);
3121 handler.propagate_statements().await;
3122
3123 assert_eq!(handler.pending_sends.len(), 1);
3124 assert_eq!(handler.propagation_outboxes.get(&peer_id).unwrap().len(), backlog + 1);
3125 }
3126
3127 #[tokio::test]
3128 async fn statement_pruned_between_tick_and_send_is_skipped() {
3129 let (mut handler, statement_store, _network, notification_service, _, peer_ids) =
3130 build_handler(1);
3131 let peer_id = peer_ids[0];
3132
3133 let mut kept = new_live_statement();
3134 kept.set_plain_data(b"kept".to_vec());
3135 let kept_hash = kept.hash();
3136 statement_store.insert(kept);
3137
3138 let mut pruned = new_live_statement();
3139 pruned.set_plain_data(b"pruned".to_vec());
3140 let pruned_hash = pruned.hash();
3141
3142 handler
3144 .propagation_outboxes
3145 .insert(peer_id, VecDeque::from(vec![pruned_hash, kept_hash]));
3146 handler.try_send_next_chunk(peer_id);
3147 handler.flush_pending_sends().await;
3148
3149 let sent = get_peer_hashes(¬ification_service.get_sent_notifications(), peer_id);
3150 assert_eq!(sent, vec![kept_hash]);
3151 assert!(
3152 !handler.propagation_outboxes.contains_key(&peer_id),
3153 "the drained outbox must be removed"
3154 );
3155 }
3156
3157 #[tokio::test]
3158 async fn failed_store_fetch_parks_the_peer_and_the_tick_retries() {
3159 let (mut handler, statement_store, _network, notification_service, _, peer_ids) =
3160 build_handler(1);
3161 let peer_id = peer_ids[0];
3162
3163 let mut statement = new_live_statement();
3164 statement.set_plain_data(b"statement".to_vec());
3165 let hash = statement.hash();
3166 statement_store.statements.lock().unwrap().insert(hash, statement);
3167 handler.propagation_outboxes.insert(peer_id, VecDeque::from(vec![hash]));
3168
3169 statement_store.fail_fetches.store(true, Ordering::Relaxed);
3170 handler.try_send_next_chunk(peer_id);
3171
3172 assert!(handler.pending_sends.is_empty(), "a failed fetch must not queue a send");
3173 assert_eq!(
3174 handler.propagation_outboxes.get(&peer_id).unwrap().len(),
3175 1,
3176 "the outbox must be retained"
3177 );
3178 assert_eq!(handler.parked_propagations, VecDeque::from([peer_id]));
3179
3180 statement_store.fail_fetches.store(false, Ordering::Relaxed);
3182 handler.propagate_statements().await;
3183 handler.flush_pending_sends().await;
3184
3185 let sent = get_peer_hashes(¬ification_service.get_sent_notifications(), peer_id);
3186 assert_eq!(sent, vec![hash]);
3187 assert!(!handler.propagation_outboxes.contains_key(&peer_id));
3188 assert!(handler.parked_propagations.is_empty());
3189 }
3190
3191 #[tokio::test]
3192 async fn oversized_statement_in_the_outbox_is_consumed() {
3193 let (mut handler, statement_store, _network, notification_service, _, peer_ids) =
3194 build_handler(1);
3195 handler.metrics = Some(Metrics::register(&Registry::new()).unwrap());
3196 let peer_id = peer_ids[0];
3197
3198 let mut oversized = new_live_statement();
3199 oversized.set_plain_data(vec![1u8; MAX_STATEMENT_NOTIFICATION_SIZE as usize]);
3200 let oversized_hash = oversized.hash();
3201 let mut small = new_live_statement();
3202 small.set_plain_data(b"small".to_vec());
3203 let small_hash = small.hash();
3204 statement_store.insert(oversized);
3205 statement_store.insert(small);
3206
3207 handler
3210 .propagation_outboxes
3211 .insert(peer_id, VecDeque::from(vec![oversized_hash, small_hash]));
3212 handler.try_send_next_chunk(peer_id);
3213 handler.flush_pending_sends().await;
3214
3215 let sent = get_peer_hashes(¬ification_service.get_sent_notifications(), peer_id);
3216 assert_eq!(sent, vec![small_hash]);
3217 assert_eq!(handler.metrics.as_ref().unwrap().skipped_oversized_statements.get(), 1);
3218 }
3219
3220 #[tokio::test]
3221 async fn disconnect_clears_the_outbox_and_send_slot() {
3222 let (mut handler, statement_store, _network, notification_service, _, peer_ids) =
3223 build_handler(1);
3224 let peer_id = peer_ids[0];
3225
3226 for i in 0..25u8 {
3229 let mut statement = new_live_statement();
3230 let mut data = vec![0u8; 100 * 1024];
3231 data[0] = i;
3232 statement.set_plain_data(data);
3233 statement_store
3234 .recent_statements
3235 .lock()
3236 .unwrap()
3237 .insert(statement.hash(), statement);
3238 }
3239 notification_service.block_sends();
3240 handler.propagate_statements().await;
3241 assert!(handler.propagation_outboxes.contains_key(&peer_id));
3242 assert!(handler.in_flight_chunks.contains_key(&peer_id));
3243
3244 handler
3245 .handle_notification_event(NotificationEvent::NotificationStreamClosed {
3246 peer: peer_id,
3247 })
3248 .await;
3249
3250 assert!(!handler.propagation_outboxes.contains_key(&peer_id));
3251 assert!(!handler.in_flight_chunks.contains_key(&peer_id));
3252 }
3253
3254 #[tokio::test]
3255 async fn failed_propagation_send_frees_the_slot_and_the_backlog_keeps_draining() {
3256 let (mut handler, statement_store, _network, notification_service, _, peer_ids) =
3257 build_handler(1);
3258 let peer_id = peer_ids[0];
3259
3260 let mut first = new_live_statement();
3262 first.set_plain_data(vec![1u8; 700 * 1024]);
3263 let first_hash = first.hash();
3264 let mut second = new_live_statement();
3265 second.set_plain_data(vec![2u8; 700 * 1024]);
3266 let second_hash = second.hash();
3267 statement_store.insert(first);
3268 statement_store.insert(second);
3269 handler
3270 .propagation_outboxes
3271 .insert(peer_id, VecDeque::from(vec![first_hash, second_hash]));
3272
3273 notification_service.fail_sends();
3275 handler.try_send_next_chunk(peer_id);
3276 assert!(handler.in_flight_chunks.contains_key(&peer_id));
3277 let result = handler.pending_sends.next().await.unwrap();
3278 notification_service.allow_sends();
3279 handler.handle_send_result(result);
3280
3281 handler.flush_pending_sends().await;
3284 let sent = get_peer_hashes(¬ification_service.get_sent_notifications(), peer_id);
3285 assert_eq!(sent, vec![second_hash]);
3286 assert!(!handler.propagation_outboxes.contains_key(&peer_id));
3287 assert!(!handler.in_flight_chunks.contains_key(&peer_id));
3288 }
3289
3290 #[tokio::test]
3291 async fn overflowing_outbox_drops_the_oldest_hashes() {
3292 let (mut handler, statement_store, _network, _notification_service, _, peer_ids) =
3293 build_handler(1);
3294 handler.metrics = Some(Metrics::register(&Registry::new()).unwrap());
3295 let peer_id = peer_ids[0];
3296
3297 let mut old = new_live_statement();
3298 old.set_plain_data(b"oldest".to_vec());
3299 let old_hash = old.hash();
3300
3301 let fresh_hashes: HashSet<_> = (0..3u8)
3304 .map(|i| {
3305 let mut fresh = new_live_statement();
3306 fresh.set_plain_data(vec![i; 8]);
3307 let hash = fresh.hash();
3308 statement_store.recent_statements.lock().unwrap().insert(hash, fresh);
3309 hash
3310 })
3311 .collect();
3312 handler.in_flight_chunks.insert(peer_id, 0);
3313 handler
3314 .propagation_outboxes
3315 .insert(peer_id, VecDeque::from(vec![old_hash; MAX_PROPAGATION_OUTBOX_LEN]));
3316
3317 handler.propagate_statements().await;
3318
3319 let outbox = handler.propagation_outboxes.get(&peer_id).unwrap();
3320 assert_eq!(outbox.len(), MAX_PROPAGATION_OUTBOX_LEN);
3321 let tail: HashSet<_> =
3322 outbox.iter().skip(MAX_PROPAGATION_OUTBOX_LEN - 3).copied().collect();
3323 assert_eq!(tail, fresh_hashes, "the freshest hashes must survive the overflow");
3324
3325 let metrics = handler.metrics.as_ref().unwrap();
3326 assert_eq!(
3327 metrics
3328 .undelivered_statements
3329 .with_label_values(&[send_failure::OUTBOX_FULL])
3330 .get(),
3331 3,
3332 "each dropped hash counts as undelivered"
3333 );
3334 assert_eq!(
3335 metrics.send_failures.with_label_values(&[send_failure::OUTBOX_FULL]).get(),
3336 1,
3337 "one overflow event"
3338 );
3339 }
3340
3341 #[tokio::test]
3342 async fn statement_received_while_queued_in_the_outbox_is_not_echoed() {
3343 let (mut handler, statement_store, _network, notification_service, _, peer_ids) =
3344 build_handler(1);
3345 let peer_id = peer_ids[0];
3346
3347 let mut statement = new_live_statement();
3348 statement.set_plain_data(b"received after append".to_vec());
3349 let hash = statement.hash();
3350 statement_store.insert(statement);
3351
3352 handler.propagation_outboxes.insert(peer_id, VecDeque::from(vec![hash]));
3356 handler.recently_received_statements.insert(hash, HashSet::from_iter([peer_id]));
3357
3358 handler.try_send_next_chunk(peer_id);
3359 handler.flush_pending_sends().await;
3360
3361 assert!(
3362 get_peer_hashes(¬ification_service.get_sent_notifications(), peer_id).is_empty(),
3363 "statement returned to the peer that sent it"
3364 );
3365 }
3366
3367 async fn dispatch_disconnects(
3370 handler: &mut StatementHandler<TestNetwork, TestSync>,
3371 network: &TestNetwork,
3372 ) {
3373 for peer in network.get_disconnected_peers() {
3374 handler
3375 .handle_notification_event(NotificationEvent::NotificationStreamClosed { peer })
3376 .await;
3377 }
3378 }
3379
3380 #[tokio::test]
3381 async fn test_skips_processing_statements_that_already_in_store() {
3382 let (mut handler, statement_store, _network, _notification_service, queue_receiver, _) =
3383 build_handler(1);
3384
3385 let mut statement1 = new_live_statement();
3386 statement1.set_plain_data(b"statement1".to_vec());
3387
3388 statement_store.insert(statement1.clone());
3389
3390 let mut statement2 = new_live_statement();
3391 statement2.set_plain_data(b"statement2".to_vec());
3392 let hash2 = statement2.hash();
3393
3394 let peer_id = *handler.peers.keys().next().unwrap();
3395
3396 handler.on_statements(peer_id, vec![statement1, statement2]);
3397
3398 let to_submit = queue_receiver.try_recv();
3399 assert_eq!(to_submit.unwrap().0.hash(), hash2, "Expected only statement2 to be queued");
3400
3401 let no_more = queue_receiver.try_recv();
3402 assert!(no_more.is_err(), "Expected only one statement to be queued");
3403 }
3404
3405 #[tokio::test]
3406 async fn test_reports_for_duplicate_statements() {
3407 let (mut handler, statement_store, network, _notification_service, queue_receiver, _) =
3408 build_handler(1);
3409
3410 let peer_id = *handler.peers.keys().next().unwrap();
3411
3412 let mut statement1 = new_live_statement();
3413 statement1.set_plain_data(b"statement1".to_vec());
3414
3415 handler.on_statements(peer_id, vec![statement1.clone()]);
3416 {
3417 let (s, _) = queue_receiver.try_recv().unwrap();
3419 statement_store.insert(s);
3420 handler.network.report_peer(peer_id, rep::ANY_STATEMENT_REFUND);
3421 }
3422
3423 handler.on_statements(peer_id, vec![statement1]);
3424
3425 let reports = network.get_reports();
3426 assert_eq!(
3427 reports,
3428 vec![
3429 (peer_id, rep::ANY_STATEMENT), (peer_id, rep::ANY_STATEMENT_REFUND), (peer_id, rep::DUPLICATE_STATEMENT) ],
3433 "Expected ANY_STATEMENT, ANY_STATEMENT_REFUND, DUPLICATE_STATEMENT reputation change, but got: {:?}",
3434 reports
3435 );
3436 }
3437
3438 #[tokio::test]
3439 async fn test_splits_large_batches_into_smaller_chunks() {
3440 let (mut handler, statement_store, _network, notification_service, _queue_receiver, _) =
3441 build_handler(1);
3442
3443 let num_statements = 30;
3444 let statement_size = 100 * 1024; for i in 0..num_statements {
3446 let mut statement = new_live_statement();
3447 let mut data = vec![0u8; statement_size];
3448 data[0] = i as u8;
3449 statement.set_plain_data(data);
3450 let hash = statement.hash();
3451 statement_store.recent_statements.lock().unwrap().insert(hash, statement);
3452 }
3453
3454 handler.propagate_statements().await;
3455 handler.flush_pending_sends().await;
3456
3457 let sent = notification_service.get_sent_notifications();
3458 let mut total_statements_sent = 0;
3459 assert!(
3460 sent.len() == 3,
3461 "Expected batch to be split into 3 chunks, but got {} chunks",
3462 sent.len()
3463 );
3464 for (_peer, notification) in sent.iter() {
3465 assert!(
3466 notification.len() <= MAX_STATEMENT_NOTIFICATION_SIZE as usize,
3467 "Notification size {} exceeds limit {}",
3468 notification.len(),
3469 MAX_STATEMENT_NOTIFICATION_SIZE
3470 );
3471 if let Ok(stmts) = <Statements as Decode>::decode(&mut notification.as_slice()) {
3472 total_statements_sent += stmts.len();
3473 }
3474 }
3475
3476 assert_eq!(
3477 total_statements_sent, num_statements,
3478 "Expected all {} statements to be sent, but only {} were sent",
3479 num_statements, total_statements_sent
3480 );
3481 }
3482
3483 #[tokio::test]
3484 async fn test_skips_only_oversized_statements() {
3485 let (mut handler, statement_store, _network, notification_service, _queue_receiver, _) =
3486 build_handler(1);
3487
3488 let mut statement1 = new_live_statement();
3489 statement1.set_plain_data(vec![1u8; 100]);
3490 let hash1 = statement1.hash();
3491 statement_store
3492 .recent_statements
3493 .lock()
3494 .unwrap()
3495 .insert(hash1, statement1.clone());
3496
3497 let mut oversized1 = new_live_statement();
3498 oversized1.set_plain_data(vec![2u8; MAX_STATEMENT_NOTIFICATION_SIZE as usize * 100]);
3499 let hash_oversized1 = oversized1.hash();
3500 statement_store
3501 .recent_statements
3502 .lock()
3503 .unwrap()
3504 .insert(hash_oversized1, oversized1);
3505
3506 let mut statement2 = new_live_statement();
3507 statement2.set_plain_data(vec![3u8; 100]);
3508 let hash2 = statement2.hash();
3509 statement_store
3510 .recent_statements
3511 .lock()
3512 .unwrap()
3513 .insert(hash2, statement2.clone());
3514
3515 let mut oversized2 = new_live_statement();
3516 oversized2.set_plain_data(vec![4u8; MAX_STATEMENT_NOTIFICATION_SIZE as usize]);
3517 let hash_oversized2 = oversized2.hash();
3518 statement_store
3519 .recent_statements
3520 .lock()
3521 .unwrap()
3522 .insert(hash_oversized2, oversized2);
3523
3524 let mut statement3 = new_live_statement();
3525 statement3.set_plain_data(vec![5u8; 100]);
3526 let hash3 = statement3.hash();
3527 statement_store
3528 .recent_statements
3529 .lock()
3530 .unwrap()
3531 .insert(hash3, statement3.clone());
3532
3533 handler.propagate_statements().await;
3534 handler.flush_pending_sends().await;
3535
3536 let sent = notification_service.get_sent_notifications();
3537
3538 let mut sent_hashes = sent
3539 .iter()
3540 .flat_map(|(_peer, notification)| {
3541 <Statements as Decode>::decode(&mut notification.as_slice()).unwrap()
3542 })
3543 .map(|s| s.hash())
3544 .collect::<Vec<_>>();
3545 sent_hashes.sort();
3546 let mut expected_hashes = vec![hash1, hash2, hash3];
3547 expected_hashes.sort();
3548 assert_eq!(sent_hashes, expected_hashes, "Only small statements should be sent");
3549 }
3550
3551 fn build_handler_no_peers() -> (
3552 StatementHandler<TestNetwork, TestSync>,
3553 TestStatementStore,
3554 TestNetwork,
3555 TestNotificationService,
3556 ) {
3557 let statement_store = TestStatementStore::new();
3558 let (queue_sender, _queue_receiver) = async_channel::bounded(2);
3559 let network = TestNetwork::new();
3560 let notification_service = TestNotificationService::new();
3561
3562 let handler = StatementHandler {
3563 protocol_name: format!("/{STATEMENT_PROTOCOL_V1}").into(),
3564 notification_service: Box::new(notification_service.clone()),
3565 propagate_timeout: (Box::pin(futures::stream::pending())
3566 as Pin<Box<dyn Stream<Item = ()> + Send>>)
3567 .fuse(),
3568 pending_statements: FuturesUnordered::new(),
3569 pending_statements_peers: HashMap::new(),
3570 recently_received_statements: HashMap::new(),
3571 network: network.clone(),
3572 sync: TestSync::new(),
3573 sync_event_stream: (Box::pin(futures::stream::pending())
3574 as Pin<Box<dyn Stream<Item = sc_network_sync::types::SyncEvent> + Send>>)
3575 .fuse(),
3576 peers: HashMap::new(),
3577 statement_store: Arc::new(statement_store.clone()),
3578 queue_sender,
3579 statements_per_second: NonZeroU32::new(DEFAULT_STATEMENTS_PER_SECOND)
3580 .expect("DEFAULT_STATEMENTS_PER_SECOND is nonzero"),
3581 metrics: None,
3582 initial_sync_timeout: Box::pin(futures::future::pending()),
3583 pending_affinities_timeout: Box::pin(futures::future::pending()),
3584 pending_initial_syncs: HashMap::new(),
3585 initial_sync_peer_queue: VecDeque::new(),
3586 next_initial_sync_id: 0,
3587 initial_sync_in_flight_bytes: 0,
3588 propagation_outboxes: HashMap::new(),
3589 in_flight_chunks: HashMap::new(),
3590 next_chunk_id: 0,
3591 propagation_in_flight_bytes: 0,
3592 parked_propagations: VecDeque::new(),
3593 pending_sends: FuturesUnordered::new(),
3594 deferred_peers: HashSet::new(),
3595 dropped_statements_during_sync: false,
3596 sync_recovery_peer: None,
3597 sync_recovery_readd_timeout: Box::pin(futures::future::pending()),
3598 };
3599 (handler, statement_store, network, notification_service)
3600 }
3601
3602 fn build_handler_no_peers_light() -> (
3604 StatementHandler<TestNetwork, TestSync>,
3605 TestStatementStore,
3606 TestNetwork,
3607 TestNotificationService,
3608 ) {
3609 let statement_store = TestStatementStore::new();
3610 let (queue_sender, _queue_receiver) = async_channel::bounded(2);
3611 let network = TestNetwork::new_light();
3612 let notification_service = TestNotificationService::new();
3613
3614 let handler = StatementHandler {
3615 protocol_name: format!("/{STATEMENT_PROTOCOL_V1}").into(),
3616 notification_service: Box::new(notification_service.clone()),
3617 propagate_timeout: (Box::pin(futures::stream::pending())
3618 as Pin<Box<dyn Stream<Item = ()> + Send>>)
3619 .fuse(),
3620 pending_statements: FuturesUnordered::new(),
3621 pending_statements_peers: HashMap::new(),
3622 recently_received_statements: HashMap::new(),
3623 network: network.clone(),
3624 sync: TestSync::new(),
3625 sync_event_stream: (Box::pin(futures::stream::pending())
3626 as Pin<Box<dyn Stream<Item = sc_network_sync::types::SyncEvent> + Send>>)
3627 .fuse(),
3628 peers: HashMap::new(),
3629 statement_store: Arc::new(statement_store.clone()),
3630 queue_sender,
3631 statements_per_second: NonZeroU32::new(DEFAULT_STATEMENTS_PER_SECOND)
3632 .expect("DEFAULT_STATEMENTS_PER_SECOND is nonzero"),
3633 metrics: None,
3634 initial_sync_timeout: Box::pin(futures::future::pending()),
3635 pending_affinities_timeout: Box::pin(futures::future::pending()),
3636 pending_initial_syncs: HashMap::new(),
3637 initial_sync_peer_queue: VecDeque::new(),
3638 next_initial_sync_id: 0,
3639 initial_sync_in_flight_bytes: 0,
3640 propagation_outboxes: HashMap::new(),
3641 in_flight_chunks: HashMap::new(),
3642 next_chunk_id: 0,
3643 propagation_in_flight_bytes: 0,
3644 parked_propagations: VecDeque::new(),
3645 pending_sends: FuturesUnordered::new(),
3646 deferred_peers: HashSet::new(),
3647 dropped_statements_during_sync: false,
3648 sync_recovery_peer: None,
3649 sync_recovery_readd_timeout: Box::pin(futures::future::pending()),
3650 };
3651 (handler, statement_store, network, notification_service)
3652 }
3653
3654 #[tokio::test]
3655 async fn test_initial_sync_burst_single_peer() {
3656 let (mut handler, statement_store, _network, notification_service, _, _) = build_handler(0);
3657
3658 let num_statements = 200;
3661 let statement_size = 100 * 1024; let mut expected_hashes = Vec::new();
3663 for i in 0..num_statements {
3664 let mut statement = new_live_statement();
3665 let mut data = vec![0u8; statement_size];
3666 data[0] = (i % 256) as u8;
3668 data[1] = (i / 256) as u8;
3669 statement.set_plain_data(data);
3670 let hash = statement.hash();
3671 expected_hashes.push(hash);
3672 statement_store.insert(statement);
3673 }
3674
3675 let peer_id = PeerId::random();
3677
3678 handler
3679 .handle_notification_event(NotificationEvent::NotificationStreamOpened {
3680 peer: peer_id,
3681 direction: sc_network::service::traits::Direction::Inbound,
3682 handshake: vec![],
3683 negotiated_fallback: Some(format!("/{STATEMENT_PROTOCOL_V1}").into()),
3684 })
3685 .await;
3686
3687 assert!(handler.peers.contains_key(&peer_id));
3689 assert!(handler.pending_initial_syncs.contains_key(&peer_id));
3690 assert_eq!(handler.initial_sync_peer_queue.len(), 1);
3691
3692 let mut burst_count = 0;
3694 while handler.pending_initial_syncs.contains_key(&peer_id) {
3695 handler.process_initial_sync_burst();
3696 handler.flush_pending_sends().await;
3697 burst_count += 1;
3698 assert!(burst_count <= 300, "Too many bursts, possible infinite loop");
3700 }
3701
3702 assert!(
3705 burst_count >= 10,
3706 "Expected multiple bursts for 200 statements of 100KB each, got {}",
3707 burst_count
3708 );
3709
3710 let sent = notification_service.get_sent_notifications();
3712 let mut sent_hashes: Vec<_> = sent
3713 .iter()
3714 .flat_map(|(peer, notification)| {
3715 assert_eq!(*peer, peer_id);
3716 <Statements as Decode>::decode(&mut notification.as_slice()).unwrap()
3717 })
3718 .map(|s| s.hash())
3719 .collect();
3720 sent_hashes.sort();
3721 expected_hashes.sort();
3722
3723 assert_eq!(
3724 sent_hashes.len(),
3725 expected_hashes.len(),
3726 "Expected {} statements to be sent, got {}",
3727 expected_hashes.len(),
3728 sent_hashes.len()
3729 );
3730 assert_eq!(sent_hashes, expected_hashes, "All statements should be sent");
3731
3732 assert!(!handler.pending_initial_syncs.contains_key(&peer_id));
3734 assert!(handler.initial_sync_peer_queue.is_empty());
3735 }
3736
3737 #[tokio::test]
3738 async fn initial_sync_network_error_leaves_the_sync_to_retry() {
3739 let (mut handler, statement_store, _, notification_service, _, peer_ids) = build_handler(1);
3740 let peer_id = peer_ids[0];
3741 handler.metrics = Some(Metrics::register(&Registry::new()).unwrap());
3742
3743 let mut hashes: Vec<_> = [b"initial-sync-one".to_vec(), b"initial-sync-two".to_vec()]
3746 .into_iter()
3747 .map(|payload| {
3748 let mut statement = new_live_statement();
3749 statement.set_plain_data(payload);
3750 let hash = statement.hash();
3751 statement_store.insert(statement);
3752 hash
3753 })
3754 .collect();
3755
3756 handler.schedule_initial_sync_for_peer(peer_id);
3757 assert!(handler.pending_initial_syncs.contains_key(&peer_id));
3758
3759 notification_service.fail_sends();
3761 handler.process_initial_sync_burst();
3762 handler.flush_pending_sends().await;
3763
3764 assert!(notification_service.get_sent_notifications().is_empty());
3765 let pending = handler.pending_initial_syncs.get(&peer_id).unwrap();
3766 assert_eq!(pending.cursor, 0, "an unconfirmed chunk must not advance the cursor");
3767 assert_eq!(
3768 handler.initial_sync_peer_queue.iter().filter(|p| **p == peer_id).count(),
3769 1,
3770 "the failed send must requeue the peer"
3771 );
3772
3773 let metrics = handler.metrics.as_ref().unwrap();
3774 assert_eq!(metrics.send_failures.with_label_values(&[send_failure::NETWORK]).get(), 1,);
3775 assert_eq!(
3776 metrics.undelivered_statements.with_label_values(&[send_failure::NETWORK]).get(),
3777 0,
3778 "a retried sync chunk is not abandoned, so its statements are not undelivered"
3779 );
3780 assert_eq!(
3781 metrics.send_failures.with_label_values(&[send_failure::TIMEOUT]).get(),
3782 0,
3783 "a network error must not be attributed to a timeout"
3784 );
3785
3786 notification_service.allow_sends();
3788 handler.process_initial_sync_burst();
3789 handler.flush_pending_sends().await;
3790
3791 let mut sent = get_peer_hashes(¬ification_service.get_sent_notifications(), peer_id);
3792 sent.sort();
3793 hashes.sort();
3794 assert_eq!(sent, hashes);
3795
3796 handler.process_initial_sync_burst();
3797 assert!(!handler.pending_initial_syncs.contains_key(&peer_id));
3798 }
3799
3800 #[tokio::test]
3801 async fn rescheduling_a_sync_drops_the_propagation_outbox() {
3802 let (mut handler, statement_store, _network, notification_service, _, peer_ids) =
3803 build_handler(1);
3804 let peer_id = peer_ids[0];
3805
3806 let mut statement = new_live_statement();
3807 statement.set_plain_data(b"queued before re-sync".to_vec());
3808 let hash = statement.hash();
3809 statement_store.insert(statement);
3810 handler.propagation_outboxes.insert(peer_id, VecDeque::from(vec![hash]));
3811
3812 handler.schedule_initial_sync_for_peer(peer_id);
3813
3814 assert!(
3815 !handler.propagation_outboxes.contains_key(&peer_id),
3816 "hashes queued before the sync are covered by its cursor"
3817 );
3818
3819 handler.process_initial_sync_burst();
3821 handler.flush_pending_sends().await;
3822 let sent = get_peer_hashes(¬ification_service.get_sent_notifications(), peer_id);
3823 assert_eq!(sent, vec![hash]);
3824 }
3825
3826 #[tokio::test]
3827 async fn missing_sink_leaves_the_initial_sync_to_retry() {
3828 let (mut handler, statement_store, _network, notification_service, _, peer_ids) =
3829 build_handler(1);
3830 let peer_id = peer_ids[0];
3831
3832 let mut statement = new_live_statement();
3833 statement.set_plain_data(b"initial-sync statement".to_vec());
3834 let hash = statement.hash();
3835 statement_store.insert(statement);
3836
3837 handler.schedule_initial_sync_for_peer(peer_id);
3838
3839 notification_service.serve_sinks(0);
3841 handler.process_initial_sync_burst();
3842
3843 assert!(handler.pending_sends.is_empty());
3844 assert_eq!(
3845 handler.pending_initial_syncs.get(&peer_id).unwrap().cursor,
3846 0,
3847 "a chunk that found no sink must not advance the cursor"
3848 );
3849 assert_eq!(
3850 handler.initial_sync_peer_queue.iter().filter(|p| **p == peer_id).count(),
3851 1,
3852 "the peer must be requeued for a retry"
3853 );
3854
3855 notification_service.serve_sinks(usize::MAX);
3857 handler.process_initial_sync_burst();
3858 handler.flush_pending_sends().await;
3859
3860 let sent = get_peer_hashes(¬ification_service.get_sent_notifications(), peer_id);
3861 assert_eq!(sent, vec![hash]);
3862 }
3863
3864 #[tokio::test]
3865 async fn failed_store_fetch_retains_the_sync_and_a_later_burst_retries() {
3866 let (mut handler, statement_store, _network, notification_service, _, peer_ids) =
3867 build_handler(1);
3868 let peer_id = peer_ids[0];
3869
3870 let mut statement = new_live_statement();
3871 statement.set_plain_data(b"initial-sync statement".to_vec());
3872 let hash = statement.hash();
3873 statement_store.insert(statement);
3874
3875 handler.schedule_initial_sync_for_peer(peer_id);
3876 let watermark = handler.pending_initial_syncs.get(&peer_id).unwrap().watermark;
3877
3878 statement_store.fail_fetches.store(true, Ordering::Relaxed);
3879 handler.process_initial_sync_burst();
3880
3881 assert!(handler.pending_sends.is_empty(), "a failed fetch must not queue a send");
3882 let pending = handler.pending_initial_syncs.get(&peer_id).unwrap();
3883 assert_eq!(pending.cursor, 0, "the cursor must stay at the failed position");
3884 assert_eq!(pending.watermark, watermark);
3885 assert_eq!(
3886 handler.initial_sync_peer_queue.iter().filter(|p| **p == peer_id).count(),
3887 1,
3888 "the peer must be requeued for a retry"
3889 );
3890
3891 statement_store.fail_fetches.store(false, Ordering::Relaxed);
3893 handler.process_initial_sync_burst();
3894 handler.flush_pending_sends().await;
3895
3896 let sent = get_peer_hashes(¬ification_service.get_sent_notifications(), peer_id);
3897 assert_eq!(sent, vec![hash]);
3898
3899 handler.process_initial_sync_burst();
3901 assert!(!handler.pending_initial_syncs.contains_key(&peer_id));
3902 }
3903
3904 #[tokio::test]
3905 async fn missing_sink_counts_every_abandoned_statement() {
3906 let (mut handler, statement_store, _, notification_service, _, peer_ids) = build_handler(1);
3907 let peer_id = peer_ids[0];
3908 handler.metrics = Some(Metrics::register(&Registry::new()).unwrap());
3909
3910 let total = 25;
3914 for i in 0..total {
3915 let mut statement = new_live_statement();
3916 let mut data = vec![0u8; 100 * 1024];
3917 data[0] = i as u8;
3918 statement.set_plain_data(data);
3919 statement_store
3920 .recent_statements
3921 .lock()
3922 .unwrap()
3923 .insert(statement.hash(), statement);
3924 }
3925
3926 notification_service.serve_sinks(1);
3928 handler.propagate_statements().await;
3929 handler.flush_pending_sends().await;
3930
3931 let delivered =
3932 get_peer_hashes(¬ification_service.get_sent_notifications(), peer_id).len();
3933 assert!(delivered > 0);
3934 assert!(
3935 total - delivered > delivered,
3936 "the abandoned remainder must exceed one full chunk, delivered {delivered} of {total}"
3937 );
3938
3939 let metrics = handler.metrics.as_ref().unwrap();
3940 assert_eq!(metrics.send_failures.with_label_values(&[send_failure::NO_SINK]).get(), 1,);
3941 assert_eq!(
3942 metrics.undelivered_statements.with_label_values(&[send_failure::NO_SINK]).get(),
3943 (total - delivered) as u64,
3944 "every statement after the delivered chunk counts as undelivered"
3945 );
3946 }
3947
3948 #[tokio::test]
3949 async fn burst_with_nothing_to_send_returns_the_peer_to_the_queue() {
3950 let (mut handler, statement_store, _network, notification_service, _, peer_ids) =
3951 build_handler(1);
3952 let peer_id = peer_ids[0];
3953
3954 let mut statement = new_live_statement();
3955 statement.set_plain_data(b"filtered by affinity".to_vec());
3956 statement.set_topic(0, [0xAA; 32].into());
3957 statement_store.insert(statement);
3958
3959 let mut filter = AffinityFilter::new(BLOOM_SEED, 0.01, 100);
3962 filter.insert(&[0xBB; 32]);
3963 {
3964 let peer = handler.peers.get_mut(&peer_id).unwrap();
3965 peer.protocol_version = PeerProtocolVersion::V2;
3966 peer.topic_affinity = Some(filter);
3967 }
3968 handler.schedule_initial_sync_for_peer(peer_id);
3969
3970 handler.process_initial_sync_burst();
3973 assert!(handler.pending_sends.is_empty());
3974 assert!(notification_service.get_sent_notifications().is_empty());
3975 assert_eq!(handler.initial_sync_peer_queue.len(), 1);
3976
3977 handler.process_initial_sync_burst();
3978 assert!(!handler.pending_initial_syncs.contains_key(&peer_id));
3979 assert!(handler.initial_sync_peer_queue.is_empty());
3980 }
3981
3982 #[tokio::test]
3983 async fn initial_sync_keeps_one_chunk_per_peer_in_flight() {
3984 let (mut handler, statement_store, _network, notification_service, _, peer_ids) =
3985 build_handler(1);
3986 let peer_id = peer_ids[0];
3987
3988 for i in 0..25u8 {
3991 let mut statement = new_live_statement();
3992 let mut data = vec![0u8; 100 * 1024];
3993 data[0] = i;
3994 statement.set_plain_data(data);
3995 statement_store.insert(statement);
3996 }
3997
3998 notification_service.block_sends();
4000 handler.schedule_initial_sync_for_peer(peer_id);
4001
4002 for _ in 0..10 {
4003 handler.process_initial_sync_burst();
4004 }
4005
4006 assert_eq!(handler.pending_sends.len(), 1);
4007 assert!(handler.initial_sync_peer_queue.is_empty());
4008 }
4009
4010 #[tokio::test]
4011 async fn superseded_initial_sync_result_is_discarded() {
4012 let (mut handler, statement_store, _network, _notification_service, _, peer_ids) =
4013 build_handler(1);
4014 let peer_id = peer_ids[0];
4015
4016 let mut statement = new_live_statement();
4017 statement.set_plain_data(b"superseded".to_vec());
4018 statement_store.insert(statement);
4019
4020 handler.schedule_initial_sync_for_peer(peer_id);
4021 handler.process_initial_sync_burst();
4022 assert_eq!(handler.pending_sends.len(), 1);
4023 assert!(handler.initial_sync_in_flight_bytes > 0);
4024
4025 handler.schedule_initial_sync_for_peer(peer_id);
4027 handler.flush_pending_sends().await;
4028
4029 assert_eq!(
4030 handler.initial_sync_peer_queue.iter().filter(|peer| **peer == peer_id).count(),
4031 1
4032 );
4033 assert_eq!(handler.initial_sync_in_flight_bytes, 0);
4034 }
4035
4036 #[tokio::test]
4037 async fn failed_superseded_initial_sync_result_does_not_abort_the_new_sync() {
4038 let (mut handler, statement_store, _network, notification_service, _, peer_ids) =
4039 build_handler(1);
4040 let peer_id = peer_ids[0];
4041
4042 let mut statement = new_live_statement();
4043 statement.set_plain_data(b"superseded failure".to_vec());
4044 statement_store.insert(statement);
4045
4046 handler.schedule_initial_sync_for_peer(peer_id);
4048 handler.process_initial_sync_burst();
4049 notification_service.fail_sends();
4050
4051 handler.schedule_initial_sync_for_peer(peer_id);
4053 let sync_id = handler.pending_initial_syncs.get(&peer_id).unwrap().sync_id;
4054 handler.flush_pending_sends().await;
4055
4056 assert_eq!(
4057 handler.pending_initial_syncs.get(&peer_id).map(|pending| pending.sync_id),
4058 Some(sync_id)
4059 );
4060 assert_eq!(
4061 handler.initial_sync_peer_queue.iter().filter(|peer| **peer == peer_id).count(),
4062 1
4063 );
4064 }
4065
4066 #[tokio::test]
4067 async fn initial_sync_respects_the_payload_size_boundary() {
4068 let (mut handler, statement_store, _network, notification_service, _, peer_ids) =
4069 build_handler(1);
4070 handler.metrics = Some(Metrics::register(&Registry::new()).unwrap());
4071 let peer_id = peer_ids[0];
4072 let overhead = handler.peers.get(&peer_id).unwrap().protocol_version.envelope_overhead();
4073 let max_size = max_statement_payload_size(overhead);
4074
4075 let mut data_len = max_size - 32;
4076 let exact = loop {
4077 let mut candidate = new_live_statement();
4078 candidate.set_plain_data(vec![7u8; data_len]);
4079 let size = candidate.encoded_size();
4080 assert!(size <= max_size, "no data length encodes to exactly {max_size}");
4081 if size == max_size {
4082 break candidate;
4083 }
4084 data_len += 1;
4085 };
4086 let exact_hash = exact.hash();
4087 let mut oversized = new_live_statement();
4088 oversized.set_plain_data(vec![2u8; MAX_STATEMENT_NOTIFICATION_SIZE as usize]);
4089 statement_store.insert(exact);
4090 statement_store.insert(oversized);
4091
4092 handler.schedule_initial_sync_for_peer(peer_id);
4093 for _ in 0..10 {
4094 if !handler.pending_initial_syncs.contains_key(&peer_id) {
4095 break;
4096 }
4097 handler.process_initial_sync_burst();
4098 handler.flush_pending_sends().await;
4099 }
4100
4101 assert!(!handler.pending_initial_syncs.contains_key(&peer_id));
4102 assert_eq!(handler.metrics.as_ref().unwrap().skipped_oversized_statements.get(), 1);
4103 let sent = get_peer_hashes(¬ification_service.get_sent_notifications(), peer_id);
4104 assert_eq!(sent, vec![exact_hash], "only the exactly-fitting statement is delivered");
4105 }
4106
4107 #[tokio::test]
4108 async fn initial_sync_in_flight_budget_is_bounded() {
4109 let (mut handler, statement_store, _network, notification_service, _, peer_ids) =
4110 build_handler(20);
4111
4112 for i in 0..11u8 {
4115 let mut statement = new_live_statement();
4116 let mut data = vec![0u8; 100 * 1024];
4117 data[0] = i;
4118 statement.set_plain_data(data);
4119 statement_store.insert(statement);
4120 }
4121
4122 notification_service.block_sends();
4124 for peer_id in &peer_ids {
4125 handler.schedule_initial_sync_for_peer(*peer_id);
4126 }
4127
4128 let mut bursts = 0;
4129 while handler.initial_sync_in_flight_bytes < MAX_SEND_IN_FLIGHT_BYTES {
4130 handler.process_initial_sync_burst();
4131 bursts += 1;
4132 assert!(bursts <= 100, "the budget was never reached after {bursts} bursts");
4133 }
4134
4135 let in_flight = handler.initial_sync_in_flight_bytes;
4136 assert!(in_flight >= MAX_SEND_IN_FLIGHT_BYTES);
4137 assert!(
4138 in_flight < MAX_SEND_IN_FLIGHT_BYTES + MAX_STATEMENT_NOTIFICATION_SIZE,
4139 "the budget may only be overshot by the single chunk that crossed it, got {in_flight}"
4140 );
4141
4142 let queued = handler.initial_sync_peer_queue.clone();
4144 handler.process_initial_sync_burst();
4145 assert_eq!(handler.initial_sync_peer_queue, queued);
4146 assert_eq!(handler.initial_sync_in_flight_bytes, in_flight);
4147 }
4148
4149 #[tokio::test]
4150 async fn saturated_send_budget_defers_propagation() {
4151 let (mut handler, statement_store, _network, notification_service, _, peer_ids) =
4152 build_handler(1);
4153 let peer_id = peer_ids[0];
4154
4155 let mut statement = new_live_statement();
4156 statement.set_plain_data(b"deferred by budget".to_vec());
4157 let hash = statement.hash();
4158 statement_store.recent_statements.lock().unwrap().insert(hash, statement);
4159
4160 handler.initial_sync_in_flight_bytes = MAX_SEND_IN_FLIGHT_BYTES;
4162 handler.propagate_statements().await;
4163
4164 assert!(handler.pending_sends.is_empty(), "no chunk may be queued over the budget");
4165 assert_eq!(
4166 handler.propagation_outboxes.get(&peer_id).unwrap(),
4167 &VecDeque::from(vec![hash])
4168 );
4169 assert_eq!(handler.parked_propagations, VecDeque::from([peer_id]));
4170
4171 handler.handle_send_result(PendingSendResult {
4173 peer: peer_id,
4174 statement_count: 1,
4175 bytes_sent: MAX_SEND_IN_FLIGHT_BYTES,
4176 result: SendOutcome::Sent,
4177 kind: SendKind::InitialSync { sync_id: 0, next_cursor: 0 },
4178 chunk_id: 0,
4179 });
4180 handler.flush_pending_sends().await;
4181
4182 let sent = get_peer_hashes(¬ification_service.get_sent_notifications(), peer_id);
4183 assert_eq!(sent, vec![hash]);
4184 assert!(handler.parked_propagations.is_empty());
4185 }
4186
4187 #[tokio::test]
4188 async fn parked_peers_are_refilled_in_parking_order() {
4189 let (mut handler, statement_store, _network, notification_service, _, _) = build_handler(2);
4190
4191 let mut statement = new_live_statement();
4194 statement.set_plain_data(vec![7u8; 900 * 1024]);
4195 let hash = statement.hash();
4196 statement_store.recent_statements.lock().unwrap().insert(hash, statement);
4197
4198 handler.initial_sync_in_flight_bytes = MAX_SEND_IN_FLIGHT_BYTES;
4199 handler.propagate_statements().await;
4200 assert_eq!(handler.parked_propagations.len(), 2);
4201 let first = handler.parked_propagations[0];
4202 let second = handler.parked_propagations[1];
4203
4204 handler.handle_send_result(PendingSendResult {
4206 peer: first,
4207 statement_count: 1,
4208 bytes_sent: 100,
4209 result: SendOutcome::Sent,
4210 kind: SendKind::InitialSync { sync_id: 0, next_cursor: 0 },
4211 chunk_id: 0,
4212 });
4213 assert_eq!(handler.pending_sends.len(), 1);
4214 assert_eq!(handler.parked_propagations, VecDeque::from([second]));
4215
4216 handler.flush_pending_sends().await;
4218 let sent = notification_service.get_sent_notifications();
4219 assert_eq!(
4220 sent.iter().map(|(peer, _)| *peer).collect::<Vec<_>>(),
4221 vec![first, second],
4222 "peers must be served in parking order"
4223 );
4224 assert!(handler.parked_propagations.is_empty());
4225 assert_eq!(handler.propagation_in_flight_bytes, 0);
4226 }
4227
4228 #[tokio::test]
4229 async fn pending_initial_sync_reserves_send_budget_from_propagation() {
4230 let (mut handler, statement_store, _network, notification_service, _, peer_ids) =
4231 build_handler(2);
4232 let propagation_peer = peer_ids[0];
4233 let sync_peer = peer_ids[1];
4234
4235 let mut statement = new_live_statement();
4236 statement.set_plain_data(b"backlog".to_vec());
4237 let hash = statement.hash();
4238 statement_store.insert(statement);
4239 handler.schedule_initial_sync_for_peer(sync_peer);
4240 assert!(handler.pending_initial_syncs.contains_key(&sync_peer));
4241
4242 handler.propagation_in_flight_bytes = MAX_SEND_IN_FLIGHT_BYTES;
4244 handler
4245 .propagation_outboxes
4246 .insert(propagation_peer, VecDeque::from(vec![hash]));
4247 handler.parked_propagations.push_back(propagation_peer);
4248
4249 handler.handle_send_result(PendingSendResult {
4252 peer: propagation_peer,
4253 statement_count: 1,
4254 bytes_sent: INITIAL_SYNC_RESERVED_BYTES,
4255 result: SendOutcome::Sent,
4256 kind: SendKind::Propagation,
4257 chunk_id: 0,
4258 });
4259 assert!(handler.pending_sends.is_empty(), "the reserve must not refill propagation");
4260 assert_eq!(handler.parked_propagations, VecDeque::from([propagation_peer]));
4261
4262 handler.process_initial_sync_burst();
4264 assert_eq!(handler.pending_sends.len(), 1, "the sync burst must use the reserve");
4265 handler.flush_pending_sends().await;
4266 let synced = get_peer_hashes(¬ification_service.get_sent_notifications(), sync_peer);
4267 assert_eq!(synced, vec![hash]);
4268 assert_eq!(
4269 handler.parked_propagations,
4270 VecDeque::from([propagation_peer]),
4271 "the reserve holds while the sync is still pending"
4272 );
4273
4274 handler.process_initial_sync_burst();
4277 assert!(handler.pending_initial_syncs.is_empty());
4278 handler.fill_parked_propagations();
4279 assert!(handler.parked_propagations.is_empty());
4280 handler.flush_pending_sends().await;
4281 let propagated =
4282 get_peer_hashes(¬ification_service.get_sent_notifications(), propagation_peer);
4283 assert_eq!(propagated, vec![hash]);
4284 }
4285
4286 #[tokio::test]
4287 async fn burst_skips_a_peer_with_a_busy_send_slot() {
4288 let (mut handler, statement_store, _network, notification_service, _, peer_ids) =
4289 build_handler(2);
4290
4291 let mut statement = new_live_statement();
4292 statement.set_plain_data(b"burst behind a busy slot".to_vec());
4293 statement_store.insert(statement);
4294
4295 handler.schedule_initial_sync_for_peer(peer_ids[0]);
4296 handler.schedule_initial_sync_for_peer(peer_ids[1]);
4297 handler.in_flight_chunks.insert(peer_ids[0], 42);
4299
4300 handler.process_initial_sync_burst();
4301 handler.flush_pending_sends().await;
4302
4303 let sent = notification_service.get_sent_notifications();
4305 assert_eq!(sent.len(), 1);
4306 assert_eq!(sent[0].0, peer_ids[1]);
4307 assert!(handler.pending_initial_syncs.contains_key(&peer_ids[0]));
4309 assert!(handler.initial_sync_peer_queue.contains(&peer_ids[0]));
4310 }
4311
4312 #[tokio::test]
4313 async fn peer_with_both_kinds_pending_sends_propagation_first() {
4314 let (mut handler, statement_store, _network, notification_service, _, peer_ids) =
4315 build_handler(1);
4316 let peer_id = peer_ids[0];
4317
4318 let mut synced = new_live_statement();
4321 synced.set_plain_data(b"snapshot statement".to_vec());
4322 let synced_hash = synced.hash();
4323 statement_store.insert(synced);
4324 handler.schedule_initial_sync_for_peer(peer_id);
4325 handler.process_initial_sync_burst();
4326 assert!(handler.in_flight_chunks.contains_key(&peer_id));
4327
4328 let mut fresh = new_live_statement();
4329 fresh.set_plain_data(b"fresh gossip".to_vec());
4330 let fresh_hash = fresh.hash();
4331 statement_store.recent_statements.lock().unwrap().insert(fresh_hash, fresh);
4332 handler.propagate_statements().await;
4333 assert_eq!(handler.pending_sends.len(), 1, "the propagation chunk must wait in the outbox");
4334
4335 let result = handler.pending_sends.next().await.unwrap();
4338 handler.handle_send_result(result);
4339 assert_eq!(handler.pending_sends.len(), 1);
4340 assert!(handler.in_flight_chunks.contains_key(&peer_id));
4341 handler.process_initial_sync_burst();
4342 assert_eq!(handler.pending_sends.len(), 1, "a burst must not bypass the busy slot");
4343
4344 handler.flush_pending_sends().await;
4345 let sent = get_peer_hashes(¬ification_service.get_sent_notifications(), peer_id);
4346 assert_eq!(sent, vec![synced_hash, fresh_hash]);
4347 }
4348
4349 #[tokio::test]
4350 async fn initial_sync_and_propagation_share_the_budget() {
4351 let (mut handler, statement_store, _network, _notification_service, _, peer_ids) =
4352 build_handler(1);
4353 let peer_id = peer_ids[0];
4354
4355 let mut statement = new_live_statement();
4356 statement.set_plain_data(b"shared budget".to_vec());
4357 statement_store.insert(statement);
4358 handler.schedule_initial_sync_for_peer(peer_id);
4359
4360 handler.propagation_in_flight_bytes = MAX_SEND_IN_FLIGHT_BYTES;
4362 let queued = handler.initial_sync_peer_queue.clone();
4363 handler.process_initial_sync_burst();
4364
4365 assert!(handler.pending_sends.is_empty());
4366 assert_eq!(
4367 handler.initial_sync_peer_queue, queued,
4368 "a throttled burst must not burn the peer's turn"
4369 );
4370 }
4371
4372 #[tokio::test]
4373 async fn test_initial_sync_burst_multiple_peers_round_robin() {
4374 let (mut handler, statement_store, _network, notification_service, _, _) = build_handler(0);
4375
4376 let num_statements = 200;
4378 let statement_size = 100 * 1024; let mut expected_hashes = Vec::new();
4380 for i in 0..num_statements {
4381 let mut statement = new_live_statement();
4382 let mut data = vec![0u8; statement_size];
4383 data[0] = (i % 256) as u8;
4384 data[1] = (i / 256) as u8;
4385 statement.set_plain_data(data);
4386 let hash = statement.hash();
4387 expected_hashes.push(hash);
4388 statement_store.insert(statement);
4389 }
4390
4391 let peer1 = PeerId::random();
4393 let peer2 = PeerId::random();
4394 let peer3 = PeerId::random();
4395
4396 for peer in [peer1, peer2, peer3] {
4398 handler
4399 .handle_notification_event(NotificationEvent::NotificationStreamOpened {
4400 peer,
4401 direction: sc_network::service::traits::Direction::Inbound,
4402 handshake: vec![],
4403 negotiated_fallback: Some(format!("/{STATEMENT_PROTOCOL_V1}").into()),
4404 })
4405 .await;
4406 }
4407
4408 assert_eq!(handler.peers.len(), 3);
4410 assert_eq!(handler.pending_initial_syncs.len(), 3);
4411 assert_eq!(handler.initial_sync_peer_queue.len(), 3);
4412
4413 let mut peer_burst_order = Vec::new();
4415 let mut burst_count = 0;
4416
4417 while !handler.pending_initial_syncs.is_empty() {
4418 if let Some(&next_peer) = handler.initial_sync_peer_queue.front() {
4420 peer_burst_order.push(next_peer);
4421 }
4422 handler.process_initial_sync_burst();
4423 handler.flush_pending_sends().await;
4424 burst_count += 1;
4425 assert!(burst_count <= 500, "Too many bursts, possible infinite loop");
4427 }
4428
4429 assert!(
4432 burst_count >= 30,
4433 "Expected many bursts for 3 peers with 200 statements each, got {}",
4434 burst_count
4435 );
4436
4437 assert!(peer_burst_order.len() >= 9, "Expected at least 9 bursts");
4439 assert_eq!(peer_burst_order[0], peer1, "First burst should be peer1");
4441 assert_eq!(peer_burst_order[1], peer2, "Second burst should be peer2");
4442 assert_eq!(peer_burst_order[2], peer3, "Third burst should be peer3");
4443 assert_eq!(peer_burst_order[3], peer1, "Fourth burst should be peer1");
4445 assert_eq!(peer_burst_order[4], peer2, "Fifth burst should be peer2");
4446 assert_eq!(peer_burst_order[5], peer3, "Sixth burst should be peer3");
4447
4448 let sent = notification_service.get_sent_notifications();
4450 let mut peer1_hashes = get_peer_hashes(&sent, peer1);
4451 let mut peer2_hashes = get_peer_hashes(&sent, peer2);
4452 let mut peer3_hashes = get_peer_hashes(&sent, peer3);
4453
4454 peer1_hashes.sort();
4455 peer2_hashes.sort();
4456 peer3_hashes.sort();
4457 expected_hashes.sort();
4458
4459 assert_eq!(peer1_hashes, expected_hashes, "Peer1 should receive all statements");
4460 assert_eq!(peer2_hashes, expected_hashes, "Peer2 should receive all statements");
4461 assert_eq!(peer3_hashes, expected_hashes, "Peer3 should receive all statements");
4462
4463 assert!(handler.pending_initial_syncs.is_empty());
4465 assert!(handler.initial_sync_peer_queue.is_empty());
4466 }
4467
4468 #[tokio::test]
4469 async fn test_send_statements_in_chunks_exact_max_size() {
4470 let (mut handler, statement_store, _network, notification_service, _queue_receiver, _) =
4471 build_handler(1);
4472
4473 let max_size = MAX_STATEMENT_NOTIFICATION_SIZE as usize - Compact::<u32>::max_encoded_len();
4487 let num_statements: usize = 100;
4488 let per_statement_overhead = 1 + 1 + 8 + 1 + 2; let total_overhead = per_statement_overhead * num_statements;
4490 let total_data_size = max_size - total_overhead;
4491 let per_statement_data_size = total_data_size / num_statements;
4492 let remainder = total_data_size % num_statements;
4493
4494 let mut expected_hashes = Vec::with_capacity(num_statements);
4495 let mut total_encoded_size = 0;
4496
4497 for i in 0..num_statements {
4498 let mut statement = new_live_statement();
4499 let extra = if i < remainder { 1 } else { 0 };
4501 let mut data = vec![42u8; per_statement_data_size + extra];
4502 data[0] = i as u8;
4504 data[1] = (i >> 8) as u8;
4505 statement.set_plain_data(data);
4506
4507 total_encoded_size += statement.encoded_size();
4508
4509 let hash = statement.hash();
4510 expected_hashes.push(hash);
4511 statement_store.recent_statements.lock().unwrap().insert(hash, statement);
4512 }
4513
4514 assert!(
4516 total_encoded_size == max_size,
4517 "Total encoded size {} should be <= max_size {}",
4518 total_encoded_size,
4519 max_size
4520 );
4521
4522 handler.propagate_statements().await;
4523 handler.flush_pending_sends().await;
4524
4525 let sent = notification_service.get_sent_notifications();
4526
4527 assert_eq!(
4529 sent.len(),
4530 1,
4531 "Expected 1 notification for all {} statements, but got {}",
4532 num_statements,
4533 sent.len()
4534 );
4535
4536 let (_peer, notification) = &sent[0];
4537 assert!(
4538 notification.len() <= MAX_STATEMENT_NOTIFICATION_SIZE as usize,
4539 "Notification size {} exceeds limit {}",
4540 notification.len(),
4541 MAX_STATEMENT_NOTIFICATION_SIZE
4542 );
4543
4544 let decoded = <Statements as Decode>::decode(&mut notification.as_slice()).unwrap();
4545 assert_eq!(
4546 decoded.len(),
4547 num_statements,
4548 "Expected {} statements in the notification",
4549 num_statements
4550 );
4551
4552 let mut received_hashes: Vec<_> = decoded.iter().map(|s| s.hash()).collect();
4554 expected_hashes.sort();
4555 received_hashes.sort();
4556 assert_eq!(expected_hashes, received_hashes, "All statement hashes should match");
4557 }
4558
4559 #[tokio::test]
4560 async fn test_initial_sync_burst_size_limit_consistency() {
4561 let (mut handler, statement_store, _network, notification_service, _, _) = build_handler(0);
4564
4565 let payload_limit = max_statement_payload_size(V1_ENVELOPE_OVERHEAD);
4567
4568 let first_stmt_data_size = payload_limit / 2 + 10;
4570 let mut stmt1 = new_live_statement();
4571 stmt1.set_plain_data(vec![1u8; first_stmt_data_size]);
4572 let stmt1_encoded_size = stmt1.encoded_size();
4573
4574 let remaining = payload_limit.saturating_sub(stmt1_encoded_size);
4577 let target_stmt2_encoded = remaining + 3; let stmt2_data_size = target_stmt2_encoded.saturating_sub(4); let mut stmt2 = new_live_statement();
4580 stmt2.set_plain_data(vec![2u8; stmt2_data_size]);
4581 let stmt2_encoded_size = stmt2.encoded_size();
4582
4583 let total_encoded = stmt1_encoded_size + stmt2_encoded_size;
4584
4585 assert!(
4587 total_encoded > payload_limit,
4588 "Total {} should exceed payload_limit {} so filter rejects second statement",
4589 total_encoded,
4590 payload_limit
4591 );
4592
4593 let hash1 = stmt1.hash();
4594 let hash2 = stmt2.hash();
4595 statement_store.insert(stmt1);
4596 statement_store.insert(stmt2);
4597
4598 let peer_id = PeerId::random();
4600
4601 handler
4602 .handle_notification_event(NotificationEvent::NotificationStreamOpened {
4603 peer: peer_id,
4604 direction: sc_network::service::traits::Direction::Inbound,
4605 handshake: vec![],
4606 negotiated_fallback: Some(format!("/{STATEMENT_PROTOCOL_V1}").into()),
4607 })
4608 .await;
4609
4610 assert!(handler.pending_initial_syncs.contains_key(&peer_id));
4612 let pending = handler.pending_initial_syncs.get(&peer_id).unwrap();
4613 assert_eq!(pending.watermark - pending.cursor, 2);
4614
4615 handler.process_initial_sync_burst();
4617 handler.flush_pending_sends().await;
4618
4619 let sent = notification_service.get_sent_notifications();
4622 assert_eq!(sent.len(), 1, "First burst should send one notification");
4623
4624 let decoded = <Statements as Decode>::decode(&mut sent[0].1.as_slice()).unwrap();
4625 assert_eq!(decoded.len(), 1, "First notification should contain one statement");
4626
4627 let sent_hash = decoded[0].hash();
4629 assert!(
4630 sent_hash == hash1 || sent_hash == hash2,
4631 "Sent statement should be one of the two created"
4632 );
4633
4634 assert!(handler.pending_initial_syncs.contains_key(&peer_id));
4636 let pending = handler.pending_initial_syncs.get(&peer_id).unwrap();
4637 assert_eq!(pending.watermark - pending.cursor, 1);
4638
4639 handler.process_initial_sync_burst();
4641 handler.flush_pending_sends().await;
4642
4643 let sent = notification_service.get_sent_notifications();
4644 assert_eq!(sent.len(), 2, "Second burst should send another notification");
4645
4646 let mut sent_hashes: Vec<_> = sent
4648 .iter()
4649 .flat_map(|(_, notification)| {
4650 <Statements as Decode>::decode(&mut notification.as_slice()).unwrap()
4651 })
4652 .map(|s| s.hash())
4653 .collect();
4654 sent_hashes.sort();
4655 let mut expected_hashes = vec![hash1, hash2];
4656 expected_hashes.sort();
4657 assert_eq!(sent_hashes, expected_hashes, "Both statements should be sent");
4658
4659 handler.process_initial_sync_burst();
4662 assert!(!handler.pending_initial_syncs.contains_key(&peer_id));
4663 }
4664
4665 #[tokio::test]
4666 async fn test_peer_disconnected_on_flooding() {
4667 let (mut handler, _statement_store, network, _notification_service, _queue_receiver, _) =
4668 build_handler(1);
4669
4670 let peer_id = *handler.peers.keys().next().unwrap();
4671
4672 let mut flood_statements = Vec::new();
4673 for i in 0..600_000 {
4674 let mut statement = new_live_statement();
4675 statement.set_plain_data(vec![i as u8, (i >> 8) as u8, (i >> 16) as u8]);
4676 flood_statements.push(statement);
4677 }
4678
4679 handler.on_statements(peer_id, flood_statements);
4680
4681 let reports = network.get_reports();
4682 assert!(
4683 reports
4684 .iter()
4685 .any(|(id, rep)| *id == peer_id && *rep == rep::STATEMENT_FLOODING),
4686 "Expected STATEMENT_FLOODING reputation change, but got: {:?}",
4687 reports
4688 );
4689
4690 let disconnected = network.get_disconnected_peers();
4691 assert!(
4692 disconnected.contains(&peer_id),
4693 "Expected peer {} to be disconnected, but it wasn't. Disconnected peers: {:?}",
4694 peer_id,
4695 disconnected
4696 );
4697
4698 dispatch_disconnects(&mut handler, &network).await;
4699
4700 assert!(!handler.peers.contains_key(&peer_id), "Peer should be removed from peers map");
4702 assert!(
4703 !handler.pending_initial_syncs.contains_key(&peer_id),
4704 "Peer should be removed from pending_initial_syncs"
4705 );
4706 assert!(
4707 !handler.initial_sync_peer_queue.contains(&peer_id),
4708 "Peer should be removed from initial_sync_peer_queue"
4709 );
4710 }
4711
4712 #[tokio::test]
4713 async fn test_legitimate_traffic_not_flagged() {
4714 let (
4715 mut handler,
4716 _statement_store,
4717 network,
4718 _notification_service,
4719 _queue_receiver,
4720 _,
4721 clock,
4722 ) = build_handler_with_fake_clock(1);
4723
4724 let peer_id = *handler.peers.keys().next().unwrap();
4725 let mut counter = 0u32;
4726
4727 for _ in 0..100 {
4730 handler.on_statements(peer_id, statement_batch(5_000, &mut counter));
4731 clock.advance(Duration::from_millis(100));
4732 }
4733
4734 let reports = network.get_reports();
4735 assert!(
4736 !reports
4737 .iter()
4738 .any(|(id, rep)| *id == peer_id && *rep == rep::STATEMENT_FLOODING),
4739 "Legitimate traffic should not trigger flooding detection. Reports: {:?}",
4740 reports
4741 );
4742
4743 let disconnected = network.get_disconnected_peers();
4744 assert!(
4745 !disconnected.contains(&peer_id),
4746 "Legitimate traffic should not cause disconnection. Disconnected peers: {:?}",
4747 disconnected
4748 );
4749
4750 assert!(handler.peers.contains_key(&peer_id), "Peer should still be connected");
4751 }
4752
4753 #[tokio::test]
4754 async fn test_just_over_rate_limit_triggers_flooding() {
4755 let (mut handler, _statement_store, network, _notification_service, _queue_receiver, _) =
4756 build_handler(1);
4757
4758 let peer_id = *handler.peers.keys().next().unwrap();
4759
4760 let mut statements = Vec::new();
4761 for i in 0..260_000 {
4762 let mut statement = new_live_statement();
4763 statement.set_plain_data(vec![
4764 i as u8,
4765 (i >> 8) as u8,
4766 (i >> 16) as u8,
4767 (i >> 24) as u8,
4768 ]);
4769 statements.push(statement);
4770 }
4771
4772 handler.on_statements(peer_id, statements);
4773
4774 let reports = network.get_reports();
4775 let expected_burst = DEFAULT_STATEMENTS_PER_SECOND * config::STATEMENTS_BURST_COEFFICIENT;
4776 assert!(
4777 reports
4778 .iter()
4779 .any(|(id, rep)| *id == peer_id && *rep == rep::STATEMENT_FLOODING),
4780 "Sending 260,000 statements should trigger flooding (burst limit: {}). Reports: {:?}",
4781 expected_burst,
4782 reports
4783 );
4784
4785 let disconnected = network.get_disconnected_peers();
4786 assert!(
4787 disconnected.contains(&peer_id),
4788 "Peer should be disconnected after exceeding rate limit. Disconnected: {:?}",
4789 disconnected
4790 );
4791
4792 dispatch_disconnects(&mut handler, &network).await;
4793
4794 assert!(!handler.peers.contains_key(&peer_id), "Peer should be removed from peers map");
4795 }
4796
4797 #[tokio::test]
4798 async fn test_burst_of_250k_statements_allowed() {
4799 let (mut handler, _statement_store, network, _notification_service, _queue_receiver, _) =
4800 build_handler(1);
4801
4802 let peer_id = *handler.peers.keys().next().unwrap();
4803
4804 let mut statements = Vec::new();
4805 for i in 0..250_000 {
4806 let mut statement = new_live_statement();
4807 statement.set_plain_data(vec![
4808 i as u8,
4809 (i >> 8) as u8,
4810 (i >> 16) as u8,
4811 (i >> 24) as u8,
4812 ]);
4813 statements.push(statement);
4814 }
4815
4816 handler.on_statements(peer_id, statements);
4817
4818 let reports = network.get_reports();
4819 assert!(
4820 !reports
4821 .iter()
4822 .any(|(id, rep)| *id == peer_id && *rep == rep::STATEMENT_FLOODING),
4823 "250k burst should be allowed (burst = rate × 5). Reports: {:?}",
4824 reports
4825 );
4826
4827 assert!(
4828 handler.peers.contains_key(&peer_id),
4829 "Peer should still be connected after 250k burst"
4830 );
4831 }
4832
4833 fn build_handler_with_fake_clock(
4837 num_peers: usize,
4838 ) -> (
4839 StatementHandler<TestNetwork, TestSync>,
4840 TestStatementStore,
4841 TestNetwork,
4842 TestNotificationService,
4843 async_channel::Receiver<(Statement, oneshot::Sender<SubmitResult>)>,
4844 Vec<PeerId>,
4845 FakeRelativeClock,
4846 ) {
4847 let (mut handler, statement_store, network, notification_service, queue_receiver, peer_ids) =
4848 build_handler(num_peers);
4849
4850 let clock = FakeRelativeClock::default();
4851 for peer in handler.peers.values_mut() {
4852 peer.rate_limiter =
4853 PeerRateLimiter::with_clock(statements_per_second(), burst(), &clock);
4854 }
4855
4856 (handler, statement_store, network, notification_service, queue_receiver, peer_ids, clock)
4857 }
4858
4859 fn statements_per_second() -> NonZeroU32 {
4861 NonZeroU32::new(DEFAULT_STATEMENTS_PER_SECOND)
4862 .expect("DEFAULT_STATEMENTS_PER_SECOND is nonzero")
4863 }
4864
4865 fn burst() -> NonZeroU32 {
4866 NonZeroU32::new(DEFAULT_STATEMENTS_PER_SECOND * STATEMENTS_BURST_COEFFICIENT)
4867 .expect("burst capacity is nonzero")
4868 }
4869
4870 fn statement_batch(count: u32, counter: &mut u32) -> Statements {
4871 (0..count)
4872 .map(|i| {
4873 let mut statement = Statement::new();
4874 statement.set_plain_data(vec![
4875 *counter as u8,
4876 (*counter >> 8) as u8,
4877 (*counter >> 16) as u8,
4878 i as u8,
4879 ]);
4880 *counter = counter.wrapping_add(1);
4881 statement
4882 })
4883 .collect()
4884 }
4885
4886 #[tokio::test]
4889 async fn test_sustained_rate_above_limit_triggers_flooding() {
4890 let (
4891 mut handler,
4892 _statement_store,
4893 network,
4894 _notification_service,
4895 _queue_receiver,
4896 _,
4897 clock,
4898 ) = build_handler_with_fake_clock(1);
4899
4900 let peer_id = *handler.peers.keys().next().unwrap();
4901 let mut counter = 0u32;
4902
4903 let flooding_reported = |network: &TestNetwork| {
4904 network
4905 .get_reports()
4906 .iter()
4907 .any(|(id, rep)| *id == peer_id && *rep == rep::STATEMENT_FLOODING)
4908 };
4909
4910 for batch in 0..9 {
4911 handler.on_statements(peer_id, statement_batch(30_000, &mut counter));
4912 assert!(
4913 !flooding_reported(&network),
4914 "batch {batch} is still within the burst and must not be flagged",
4915 );
4916 clock.advance(Duration::from_millis(100));
4917 }
4918
4919 handler.on_statements(peer_id, statement_batch(30_000, &mut counter));
4920 assert!(
4921 flooding_reported(&network),
4922 "the 10th batch overdraws the burst and must be flagged as flooding",
4923 );
4924
4925 let disconnected = network.get_disconnected_peers();
4926 assert!(
4927 disconnected.contains(&peer_id),
4928 "Peer should be disconnected after sustained high rate. Disconnected: {:?}",
4929 disconnected
4930 );
4931
4932 dispatch_disconnects(&mut handler, &network).await;
4933
4934 assert!(!handler.peers.contains_key(&peer_id), "Peer should be removed from peers map");
4935 }
4936
4937 #[tokio::test]
4938 async fn test_v2_peer_detected_when_no_fallback() {
4939 let (mut handler, _statement_store, _network, _notification_service) =
4940 build_handler_no_peers();
4941
4942 let peer_id = PeerId::random();
4943
4944 handler
4946 .handle_notification_event(NotificationEvent::NotificationStreamOpened {
4947 peer: peer_id,
4948 direction: sc_network::service::traits::Direction::Inbound,
4949 handshake: vec![],
4950 negotiated_fallback: None,
4951 })
4952 .await;
4953
4954 assert_eq!(
4955 handler.peers.get(&peer_id).unwrap().protocol_version,
4956 PeerProtocolVersion::V2,
4957 "Peer should be detected as v2 when no fallback is negotiated"
4958 );
4959 }
4960
4961 #[tokio::test]
4962 async fn test_v1_peer_detected_when_fallback_negotiated() {
4963 let (mut handler, _statement_store, _network, _notification_service) =
4964 build_handler_no_peers();
4965
4966 let peer_id = PeerId::random();
4967
4968 handler
4970 .handle_notification_event(NotificationEvent::NotificationStreamOpened {
4971 peer: peer_id,
4972 direction: sc_network::service::traits::Direction::Inbound,
4973 handshake: vec![],
4974 negotiated_fallback: Some(format!("/{STATEMENT_PROTOCOL_V1}").into()),
4975 })
4976 .await;
4977
4978 assert_eq!(
4979 handler.peers.get(&peer_id).unwrap().protocol_version,
4980 PeerProtocolVersion::V1,
4981 "Peer should be detected as v1 when fallback is negotiated"
4982 );
4983 }
4984
4985 #[tokio::test]
4986 async fn test_v1_peer_decodes_raw_statements() {
4987 let (mut handler, _statement_store, _network, _notification_service) =
4988 build_handler_no_peers();
4989
4990 let peer_id = PeerId::random();
4991 let (queue_sender, queue_receiver) = async_channel::bounded(10);
4992 handler.queue_sender = queue_sender;
4993
4994 handler
4996 .handle_notification_event(NotificationEvent::NotificationStreamOpened {
4997 peer: peer_id,
4998 direction: sc_network::service::traits::Direction::Inbound,
4999 handshake: vec![],
5000 negotiated_fallback: Some(format!("/{STATEMENT_PROTOCOL_V1}").into()),
5001 })
5002 .await;
5003
5004 let mut statement = new_live_statement();
5006 statement.set_plain_data(b"v1 statement".to_vec());
5007 let hash = statement.hash();
5008 let raw_encoded = vec![statement].encode();
5009
5010 handler
5011 .handle_notification_event(NotificationEvent::NotificationReceived {
5012 peer: peer_id,
5013 notification: raw_encoded.into(),
5014 })
5015 .await;
5016
5017 let (received, _) = queue_receiver.try_recv().unwrap();
5018 assert_eq!(received.hash(), hash, "V1 peer's raw statement should be decoded correctly");
5019 }
5020
5021 #[tokio::test]
5022 async fn test_v2_peer_decodes_statement_message() {
5023 let (mut handler, _statement_store, _network, _notification_service) =
5024 build_handler_no_peers();
5025
5026 let peer_id = PeerId::random();
5027 let (queue_sender, queue_receiver) = async_channel::bounded(10);
5028 handler.queue_sender = queue_sender;
5029
5030 handler
5032 .handle_notification_event(NotificationEvent::NotificationStreamOpened {
5033 peer: peer_id,
5034 direction: sc_network::service::traits::Direction::Inbound,
5035 handshake: vec![],
5036 negotiated_fallback: None,
5037 })
5038 .await;
5039
5040 let mut statement = new_live_statement();
5042 statement.set_plain_data(b"v2 statement".to_vec());
5043 let hash = statement.hash();
5044 let msg = StatementMessage::Statements(vec![statement].try_into().unwrap());
5045 let encoded = msg.encode();
5046
5047 handler
5048 .handle_notification_event(NotificationEvent::NotificationReceived {
5049 peer: peer_id,
5050 notification: encoded.into(),
5051 })
5052 .await;
5053
5054 let (received, _) = queue_receiver.try_recv().unwrap();
5055 assert_eq!(received.hash(), hash, "V2 peer's StatementMessage should be decoded correctly");
5056 }
5057
5058 #[tokio::test]
5059 async fn test_v2_peer_topic_affinity_stored() {
5060 let (mut handler, _statement_store, _network, _notification_service) =
5061 build_handler_no_peers();
5062
5063 let peer_id = PeerId::random();
5064
5065 handler
5067 .handle_notification_event(NotificationEvent::NotificationStreamOpened {
5068 peer: peer_id,
5069 direction: sc_network::service::traits::Direction::Inbound,
5070 handshake: vec![],
5071 negotiated_fallback: None,
5072 })
5073 .await;
5074
5075 assert!(
5076 handler.peers.get(&peer_id).unwrap().topic_affinity.is_none(),
5077 "Topic affinity should be None initially"
5078 );
5079
5080 let topic: [u8; 32] = [0xAA; 32];
5082 let mut filter = AffinityFilter::new(BLOOM_SEED, 0.01, 100);
5083 filter.insert(&topic);
5084 let msg = StatementMessage::ExplicitTopicAffinity(filter);
5085 let encoded = msg.encode();
5086
5087 handler
5088 .handle_notification_event(NotificationEvent::NotificationReceived {
5089 peer: peer_id,
5090 notification: encoded.into(),
5091 })
5092 .await;
5093
5094 handler.process_pending_affinities();
5096
5097 let peer_data = handler.peers.get(&peer_id).unwrap();
5098 assert!(
5099 peer_data.topic_affinity.is_some(),
5100 "Topic affinity should be set after receiving ExplicitTopicAffinity"
5101 );
5102 assert!(
5104 peer_data.topic_affinity.as_ref().unwrap().contains(&topic),
5105 "Stored affinity filter should match the topic"
5106 );
5107 }
5108
5109 #[tokio::test]
5110 async fn test_topic_affinity_filters_propagation() {
5111 let (mut handler, statement_store, _network, notification_service) =
5112 build_handler_no_peers();
5113
5114 let peer_id = PeerId::random();
5115
5116 handler
5118 .handle_notification_event(NotificationEvent::NotificationStreamOpened {
5119 peer: peer_id,
5120 direction: sc_network::service::traits::Direction::Inbound,
5121 handshake: vec![],
5122 negotiated_fallback: None,
5123 })
5124 .await;
5125
5126 let topic_aa: [u8; 32] = [0xAA; 32];
5128 let topic_bb: [u8; 32] = [0xBB; 32];
5129 let mut filter = AffinityFilter::new(BLOOM_SEED, 0.01, 100);
5130 filter.insert(&topic_aa);
5131 let msg = StatementMessage::ExplicitTopicAffinity(filter);
5132 let encoded = msg.encode();
5133 handler
5134 .handle_notification_event(NotificationEvent::NotificationReceived {
5135 peer: peer_id,
5136 notification: encoded.into(),
5137 })
5138 .await;
5139
5140 handler.process_pending_affinities();
5142
5143 let mut stmt_matching = new_live_statement();
5145 stmt_matching.set_plain_data(b"matching".to_vec());
5146 stmt_matching.set_topic(0, topic_aa.into());
5147 let hash_matching = stmt_matching.hash();
5148
5149 let mut stmt_not_matching = new_live_statement();
5150 stmt_not_matching.set_plain_data(b"not matching".to_vec());
5151 stmt_not_matching.set_topic(0, topic_bb.into());
5152 let hash_not_matching = stmt_not_matching.hash();
5153
5154 let mut stmt_no_topic = new_live_statement();
5155 stmt_no_topic.set_plain_data(b"no topic".to_vec());
5156 let hash_no_topic = stmt_no_topic.hash();
5157
5158 statement_store
5159 .recent_statements
5160 .lock()
5161 .unwrap()
5162 .insert(hash_matching, stmt_matching);
5163 statement_store
5164 .recent_statements
5165 .lock()
5166 .unwrap()
5167 .insert(hash_not_matching, stmt_not_matching);
5168 statement_store
5169 .recent_statements
5170 .lock()
5171 .unwrap()
5172 .insert(hash_no_topic, stmt_no_topic);
5173
5174 handler.propagate_statements().await;
5175 handler.flush_pending_sends().await;
5176
5177 let sent = notification_service.get_sent_notifications();
5178 let mut sent_hashes: Vec<_> = sent
5179 .iter()
5180 .flat_map(|(_, notification)| {
5181 match StatementMessage::decode(&mut notification.as_slice()).unwrap() {
5183 StatementMessage::Statements(stmts) => stmts,
5184 _ => panic!("Expected StatementMessage::Statements"),
5185 }
5186 })
5187 .map(|s| s.hash())
5188 .collect();
5189 sent_hashes.sort();
5190
5191 assert!(
5193 sent_hashes.contains(&hash_matching),
5194 "Statement matching topic affinity should be propagated"
5195 );
5196 assert!(
5197 sent_hashes.contains(&hash_no_topic),
5198 "Statement with no topics should be propagated (broadcast)"
5199 );
5200 assert!(
5201 !sent_hashes.contains(&hash_not_matching),
5202 "Statement NOT matching topic affinity should be filtered out"
5203 );
5204 }
5205
5206 #[tokio::test]
5207 async fn test_v1_peer_no_topic_filtering() {
5208 let (mut handler, statement_store, _network, notification_service) =
5209 build_handler_no_peers();
5210
5211 let peer_id = PeerId::random();
5212
5213 handler
5215 .handle_notification_event(NotificationEvent::NotificationStreamOpened {
5216 peer: peer_id,
5217 direction: sc_network::service::traits::Direction::Inbound,
5218 handshake: vec![],
5219 negotiated_fallback: Some(format!("/{STATEMENT_PROTOCOL_V1}").into()),
5220 })
5221 .await;
5222
5223 let topic_aa: [u8; 32] = [0xAA; 32];
5225 let mut stmt_with_topic = new_live_statement();
5226 stmt_with_topic.set_plain_data(b"with topic".to_vec());
5227 stmt_with_topic.set_topic(0, topic_aa.into());
5228 let hash_with_topic = stmt_with_topic.hash();
5229
5230 let mut stmt_no_topic = new_live_statement();
5231 stmt_no_topic.set_plain_data(b"no topic".to_vec());
5232 let hash_no_topic = stmt_no_topic.hash();
5233
5234 statement_store
5235 .recent_statements
5236 .lock()
5237 .unwrap()
5238 .insert(hash_with_topic, stmt_with_topic);
5239 statement_store
5240 .recent_statements
5241 .lock()
5242 .unwrap()
5243 .insert(hash_no_topic, stmt_no_topic);
5244
5245 handler.propagate_statements().await;
5246 handler.flush_pending_sends().await;
5247
5248 let sent = notification_service.get_sent_notifications();
5249 let sent_hashes: Vec<_> = sent
5250 .iter()
5251 .flat_map(|(_, notification)| {
5252 <Statements as Decode>::decode(&mut notification.as_slice()).unwrap()
5253 })
5254 .map(|s| s.hash())
5255 .collect();
5256
5257 assert_eq!(
5258 sent_hashes.len(),
5259 2,
5260 "V1 peer should receive all statements regardless of topics"
5261 );
5262 assert!(sent_hashes.contains(&hash_with_topic));
5263 assert!(sent_hashes.contains(&hash_no_topic));
5264 }
5265
5266 #[tokio::test]
5267 async fn test_affinity_change_triggers_resync() {
5268 let (mut handler, statement_store, _network, notification_service) =
5269 build_handler_no_peers_light();
5270
5271 let peer_id = PeerId::random();
5272
5273 let topic_aa: [u8; 32] = [0xAA; 32];
5275 let topic_bb: [u8; 32] = [0xBB; 32];
5276
5277 let mut stmt_aa = new_live_statement();
5278 stmt_aa.set_plain_data(b"stmt_aa".to_vec());
5279 stmt_aa.set_topic(0, topic_aa.into());
5280 let hash_aa = stmt_aa.hash();
5281
5282 let mut stmt_bb = new_live_statement();
5283 stmt_bb.set_plain_data(b"stmt_bb".to_vec());
5284 stmt_bb.set_topic(0, topic_bb.into());
5285 let hash_bb = stmt_bb.hash();
5286
5287 let mut stmt_no_topic = new_live_statement();
5288 stmt_no_topic.set_plain_data(b"no topic".to_vec());
5289 let hash_no_topic = stmt_no_topic.hash();
5290
5291 statement_store.insert(stmt_aa);
5292 statement_store.insert(stmt_bb);
5293 statement_store.insert(stmt_no_topic);
5294
5295 handler
5297 .handle_notification_event(NotificationEvent::NotificationStreamOpened {
5298 peer: peer_id,
5299 direction: sc_network::service::traits::Direction::Inbound,
5300 handshake: vec![],
5301 negotiated_fallback: None,
5302 })
5303 .await;
5304
5305 assert!(
5307 !handler.pending_initial_syncs.contains_key(&peer_id),
5308 "Light V2 peer should NOT have initial sync scheduled on connect"
5309 );
5310
5311 let mut filter = AffinityFilter::new(BLOOM_SEED, 0.01, 100);
5313 filter.insert(&topic_aa);
5314 let msg = StatementMessage::ExplicitTopicAffinity(filter);
5315 let encoded = msg.encode();
5316 handler
5317 .handle_notification_event(NotificationEvent::NotificationReceived {
5318 peer: peer_id,
5319 notification: encoded.into(),
5320 })
5321 .await;
5322
5323 handler.process_pending_affinities();
5325
5326 assert!(
5327 handler.pending_initial_syncs.contains_key(&peer_id),
5328 "Initial sync should be scheduled after setting affinity"
5329 );
5330
5331 while handler.pending_initial_syncs.contains_key(&peer_id) {
5333 handler.process_initial_sync_burst();
5334 handler.flush_pending_sends().await;
5335 }
5336
5337 let sent = notification_service.get_sent_notifications();
5338 let sent_hashes: HashSet<_> = sent
5339 .iter()
5340 .flat_map(|(_, notification)| {
5341 match StatementMessage::decode(&mut notification.as_slice()).unwrap() {
5342 StatementMessage::Statements(stmts) => stmts,
5343 _ => panic!("Expected StatementMessage::Statements"),
5344 }
5345 })
5346 .map(|s| s.hash())
5347 .collect();
5348 assert!(sent_hashes.contains(&hash_aa), "stmt_aa should be sent (matches affinity)");
5349 assert!(
5350 sent_hashes.contains(&hash_no_topic),
5351 "stmt_no_topic should be sent (broadcast, no topic)"
5352 );
5353 assert!(!sent_hashes.contains(&hash_bb), "stmt_bb should NOT be sent (filtered)");
5354
5355 let mut filter = AffinityFilter::new(BLOOM_SEED, 0.01, 100);
5357 filter.insert(&topic_bb);
5358 let msg = StatementMessage::ExplicitTopicAffinity(filter);
5359 let encoded = msg.encode();
5360 handler
5361 .handle_notification_event(NotificationEvent::NotificationReceived {
5362 peer: peer_id,
5363 notification: encoded.into(),
5364 })
5365 .await;
5366
5367 handler.process_pending_affinities();
5369
5370 assert!(
5371 handler.pending_initial_syncs.contains_key(&peer_id),
5372 "Initial sync should be re-scheduled after affinity change"
5373 );
5374
5375 notification_service.clear_sent_notifications();
5376 while handler.pending_initial_syncs.contains_key(&peer_id) {
5377 handler.process_initial_sync_burst();
5378 handler.flush_pending_sends().await;
5379 }
5380
5381 let sent_after_bb = notification_service.get_sent_notifications();
5382 let sent_hashes_bb: HashSet<_> = sent_after_bb
5383 .iter()
5384 .flat_map(|(_, notification)| {
5385 match StatementMessage::decode(&mut notification.as_slice()).unwrap() {
5386 StatementMessage::Statements(stmts) => stmts,
5387 _ => panic!("Expected StatementMessage::Statements"),
5388 }
5389 })
5390 .map(|s| s.hash())
5391 .collect();
5392 assert!(
5394 sent_hashes_bb.contains(&hash_bb),
5395 "stmt_bb should now be sent after affinity changed to topic_bb"
5396 );
5397 assert!(
5399 sent_hashes_bb.contains(&hash_no_topic),
5400 "stmt_no_topic should be re-sent (initial sync resends everything matching)"
5401 );
5402 }
5403
5404 #[tokio::test]
5405 async fn test_affinity_change_sends_previously_filtered_statements() {
5406 let (mut handler, statement_store, _network, notification_service) =
5411 build_handler_no_peers_light();
5412
5413 let peer_id = PeerId::random();
5414
5415 let topic_aa: [u8; 32] = [0xAA; 32];
5416 let topic_bb: [u8; 32] = [0xBB; 32];
5417
5418 let mut stmt_aa = new_live_statement();
5419 stmt_aa.set_plain_data(b"stmt_aa".to_vec());
5420 stmt_aa.set_topic(0, topic_aa.into());
5421 let hash_aa = stmt_aa.hash();
5422
5423 let mut stmt_bb = new_live_statement();
5424 stmt_bb.set_plain_data(b"stmt_bb".to_vec());
5425 stmt_bb.set_topic(0, topic_bb.into());
5426 let hash_bb = stmt_bb.hash();
5427
5428 statement_store.insert(stmt_aa.clone());
5429 statement_store.insert(stmt_bb.clone());
5430
5431 statement_store.recent_statements.lock().unwrap().insert(hash_aa, stmt_aa);
5433 statement_store.recent_statements.lock().unwrap().insert(hash_bb, stmt_bb);
5434
5435 handler
5437 .handle_notification_event(NotificationEvent::NotificationStreamOpened {
5438 peer: peer_id,
5439 direction: sc_network::service::traits::Direction::Inbound,
5440 handshake: vec![],
5441 negotiated_fallback: None,
5442 })
5443 .await;
5444
5445 let mut filter = AffinityFilter::new(BLOOM_SEED, 0.01, 100);
5447 filter.insert(&topic_aa);
5448 let msg = StatementMessage::ExplicitTopicAffinity(filter);
5449 let encoded = msg.encode();
5450 handler
5451 .handle_notification_event(NotificationEvent::NotificationReceived {
5452 peer: peer_id,
5453 notification: encoded.into(),
5454 })
5455 .await;
5456
5457 handler.process_pending_affinities();
5459
5460 while handler.pending_initial_syncs.contains_key(&peer_id) {
5462 handler.process_initial_sync_burst();
5463 handler.flush_pending_sends().await;
5464 }
5465
5466 let sent = notification_service.get_sent_notifications();
5467 let sent_hashes: HashSet<_> = sent
5468 .iter()
5469 .flat_map(|(_, notification)| {
5470 match StatementMessage::decode(&mut notification.as_slice()).unwrap() {
5471 StatementMessage::Statements(stmts) => stmts,
5472 _ => panic!("Expected StatementMessage::Statements"),
5473 }
5474 })
5475 .map(|s| s.hash())
5476 .collect();
5477 assert!(sent_hashes.contains(&hash_aa), "stmt_aa should be sent (matches affinity)");
5478 assert!(
5479 !sent_hashes.contains(&hash_bb),
5480 "stmt_bb should NOT be sent (filtered by affinity)"
5481 );
5482
5483 let mut stmt_aa2 = new_live_statement();
5486 stmt_aa2.set_plain_data(b"stmt_aa2".to_vec());
5487 stmt_aa2.set_topic(0, topic_aa.into());
5488 let hash_aa2 = stmt_aa2.hash();
5489 let mut stmt_bb2 = new_live_statement();
5490 stmt_bb2.set_plain_data(b"stmt_bb2".to_vec());
5491 stmt_bb2.set_topic(0, topic_bb.into());
5492 let hash_bb2 = stmt_bb2.hash();
5493 statement_store.recent_statements.lock().unwrap().insert(hash_aa2, stmt_aa2);
5494 statement_store.recent_statements.lock().unwrap().insert(hash_bb2, stmt_bb2);
5495
5496 notification_service.clear_sent_notifications();
5497 handler.propagate_statements().await;
5498 handler.flush_pending_sends().await;
5499
5500 let sent = notification_service.get_sent_notifications();
5501 let sent_hashes: HashSet<_> = sent
5502 .iter()
5503 .flat_map(|(_, notification)| {
5504 match StatementMessage::decode(&mut notification.as_slice()).unwrap() {
5505 StatementMessage::Statements(stmts) => stmts,
5506 _ => panic!("Expected StatementMessage::Statements"),
5507 }
5508 })
5509 .map(|s| s.hash())
5510 .collect();
5511 assert!(
5512 sent_hashes.contains(&hash_aa2),
5513 "stmt_aa2 should be propagated (matches affinity)"
5514 );
5515 assert!(
5516 !sent_hashes.contains(&hash_bb2),
5517 "stmt_bb2 should NOT be propagated (filtered by affinity)"
5518 );
5519 assert!(
5520 !sent_hashes.contains(&hash_aa),
5521 "a statement below the sync watermark is delivered by the cursor, not propagation"
5522 );
5523
5524 let mut filter = AffinityFilter::new(BLOOM_SEED, 0.01, 100);
5526 filter.insert(&topic_aa);
5527 filter.insert(&topic_bb);
5528 let msg = StatementMessage::ExplicitTopicAffinity(filter);
5529 let encoded = msg.encode();
5530
5531 notification_service.clear_sent_notifications();
5532 handler
5533 .handle_notification_event(NotificationEvent::NotificationReceived {
5534 peer: peer_id,
5535 notification: encoded.into(),
5536 })
5537 .await;
5538
5539 handler.process_pending_affinities();
5541
5542 while handler.pending_initial_syncs.contains_key(&peer_id) {
5544 handler.process_initial_sync_burst();
5545 handler.flush_pending_sends().await;
5546 }
5547
5548 let sent = notification_service.get_sent_notifications();
5549 let sent_hashes: HashSet<_> = sent
5550 .iter()
5551 .flat_map(|(_, notification)| {
5552 match StatementMessage::decode(&mut notification.as_slice()).unwrap() {
5553 StatementMessage::Statements(stmts) => stmts,
5554 _ => panic!("Expected StatementMessage::Statements"),
5555 }
5556 })
5557 .map(|s| s.hash())
5558 .collect();
5559 assert!(
5560 sent_hashes.contains(&hash_bb),
5561 "stmt_bb should now be sent after affinity expanded to include topic_bb"
5562 );
5563 assert!(
5565 sent_hashes.contains(&hash_aa),
5566 "stmt_aa should be re-sent (initial sync resends everything matching)"
5567 );
5568 }
5569
5570 #[test]
5571 fn test_encode_statement_refs_matches_derive_encoding() {
5572 let mut stmt1 = new_live_statement();
5573 stmt1.set_plain_data(b"first".to_vec());
5574 let mut stmt2 = new_live_statement();
5575 stmt2.set_plain_data(b"second".to_vec());
5576
5577 let refs: Vec<&Statement> = vec![&stmt1, &stmt2];
5578 let hand_rolled = StatementMessage::encode_statement_refs(&refs);
5579 let derive_encoded =
5580 StatementMessage::Statements(vec![stmt1, stmt2].try_into().unwrap()).encode();
5581
5582 assert_eq!(
5583 hand_rolled, derive_encoded,
5584 "encode_statement_refs must produce identical bytes to derive Encode"
5585 );
5586 }
5587
5588 #[test]
5589 fn test_encode_statement_refs_empty() {
5590 let refs: Vec<&Statement> = vec![];
5591 let hand_rolled = StatementMessage::encode_statement_refs(&refs);
5592 let derive_encoded = StatementMessage::Statements(vec![].try_into().unwrap()).encode();
5593
5594 assert_eq!(hand_rolled, derive_encoded);
5595 }
5596
5597 #[test]
5598 fn test_can_receive_all_combinations() {
5599 let make_peer = |is_light: bool, version: PeerProtocolVersion, has_affinity: bool| {
5600 let topic_affinity = has_affinity.then(|| AffinityFilter::new(BLOOM_SEED, 0.01, 10));
5601 Peer {
5602 rate_limiter: PeerRateLimiter::new(
5603 NonZeroU32::new(DEFAULT_STATEMENTS_PER_SECOND).expect("nonzero"),
5604 NonZeroU32::new(
5605 DEFAULT_STATEMENTS_PER_SECOND * config::STATEMENTS_BURST_COEFFICIENT,
5606 )
5607 .expect("nonzero"),
5608 ),
5609 protocol_version: version,
5610 topic_affinity,
5611 is_light,
5612 pending_topic_affinity: None,
5613 sync_watermark: 0,
5614 }
5615 };
5616
5617 assert!(make_peer(false, PeerProtocolVersion::V1, false).can_receive());
5619 assert!(make_peer(false, PeerProtocolVersion::V2, false).can_receive());
5621 assert!(make_peer(true, PeerProtocolVersion::V1, false).can_receive());
5623 assert!(!make_peer(true, PeerProtocolVersion::V2, false).can_receive());
5625 assert!(make_peer(true, PeerProtocolVersion::V2, true).can_receive());
5627 assert!(make_peer(false, PeerProtocolVersion::V2, true).can_receive());
5629 }
5630
5631 #[tokio::test]
5632 async fn test_send_chunk_v1_vs_v2_encoding() {
5633 let (mut handler, statement_store, _network, notification_service) =
5634 build_handler_no_peers();
5635
5636 let v1_peer = PeerId::random();
5637 let v2_peer = PeerId::random();
5638
5639 handler
5641 .handle_notification_event(NotificationEvent::NotificationStreamOpened {
5642 peer: v1_peer,
5643 direction: sc_network::service::traits::Direction::Inbound,
5644 handshake: vec![],
5645 negotiated_fallback: Some(format!("/{STATEMENT_PROTOCOL_V1}").into()),
5646 })
5647 .await;
5648
5649 handler
5651 .handle_notification_event(NotificationEvent::NotificationStreamOpened {
5652 peer: v2_peer,
5653 direction: sc_network::service::traits::Direction::Inbound,
5654 handshake: vec![],
5655 negotiated_fallback: None,
5656 })
5657 .await;
5658
5659 let mut stmt = new_live_statement();
5660 stmt.set_plain_data(b"encoding test".to_vec());
5661 statement_store.insert(stmt);
5662
5663 notification_service.clear_sent_notifications();
5665 handler.schedule_initial_sync_for_peer(v1_peer);
5666 while handler.pending_initial_syncs.contains_key(&v1_peer) {
5667 handler.process_initial_sync_burst();
5668 handler.flush_pending_sends().await;
5669 }
5670 let v1_sent = notification_service.get_sent_notifications();
5671 assert_eq!(v1_sent.len(), 1);
5672 let v1_bytes = &v1_sent[0].1;
5673 let decoded_v1 = <Statements as Decode>::decode(&mut v1_bytes.as_slice())
5675 .expect("V1 peer should receive raw Vec<Statement> encoding");
5676 assert_eq!(decoded_v1.len(), 1);
5677
5678 notification_service.clear_sent_notifications();
5680 handler.schedule_initial_sync_for_peer(v2_peer);
5681 while handler.pending_initial_syncs.contains_key(&v2_peer) {
5682 handler.process_initial_sync_burst();
5683 handler.flush_pending_sends().await;
5684 }
5685 let v2_sent = notification_service.get_sent_notifications();
5686 assert_eq!(v2_sent.len(), 1);
5687 let v2_bytes = &v2_sent[0].1;
5688 let decoded_v2 = StatementMessage::decode(&mut v2_bytes.as_slice())
5690 .expect("V2 peer should receive StatementMessage encoding");
5691 match decoded_v2 {
5692 StatementMessage::Statements(stmts) => assert_eq!(stmts.len(), 1),
5693 _ => panic!("Expected StatementMessage::Statements for V2 peer"),
5694 }
5695
5696 assert_ne!(v1_bytes, v2_bytes, "V1 and V2 encodings should differ");
5698 }
5699
5700 #[tokio::test]
5701 async fn test_schedule_initial_sync_replaces_existing() {
5702 let (mut handler, statement_store, _network, _notification_service) =
5703 build_handler_no_peers();
5704
5705 let peer_id = PeerId::random();
5706
5707 let mut stmt1 = new_live_statement();
5709 stmt1.set_plain_data(b"stmt1".to_vec());
5710 statement_store.insert(stmt1);
5711
5712 handler
5714 .handle_notification_event(NotificationEvent::NotificationStreamOpened {
5715 peer: peer_id,
5716 direction: sc_network::service::traits::Direction::Inbound,
5717 handshake: vec![],
5718 negotiated_fallback: Some(format!("/{STATEMENT_PROTOCOL_V1}").into()),
5719 })
5720 .await;
5721
5722 assert!(handler.pending_initial_syncs.contains_key(&peer_id));
5724 assert_eq!(
5725 handler.initial_sync_peer_queue.iter().filter(|p| **p == peer_id).count(),
5726 1,
5727 "Peer should appear exactly once in the queue"
5728 );
5729
5730 let mut stmt2 = new_live_statement();
5732 stmt2.set_plain_data(b"stmt2".to_vec());
5733 statement_store.insert(stmt2);
5734
5735 handler.schedule_initial_sync_for_peer(peer_id);
5736
5737 assert_eq!(
5739 handler.initial_sync_peer_queue.iter().filter(|p| **p == peer_id).count(),
5740 1,
5741 "Peer should NOT be duplicated in the queue after re-schedule"
5742 );
5743 let pending = handler.pending_initial_syncs.get(&peer_id).unwrap();
5745 assert_eq!(pending.cursor, 0);
5746 assert_eq!(pending.watermark, 2);
5747 }
5748
5749 #[tokio::test]
5750 async fn test_initial_sync_queued_during_major_sync_processed_after() {
5751 let statement_store = TestStatementStore::new();
5752 let (queue_sender, _queue_receiver) = async_channel::bounded(2);
5753 let network = TestNetwork::new();
5754 let notification_service = TestNotificationService::new();
5755 let sync = TestSync::new();
5756 sync.major_syncing.store(true, Ordering::Relaxed);
5758
5759 let mut handler = StatementHandler {
5760 protocol_name: format!("/{STATEMENT_PROTOCOL_V1}").into(),
5761 notification_service: Box::new(notification_service.clone()),
5762 propagate_timeout: (Box::pin(futures::stream::pending())
5763 as Pin<Box<dyn Stream<Item = ()> + Send>>)
5764 .fuse(),
5765 pending_statements: FuturesUnordered::new(),
5766 pending_statements_peers: HashMap::new(),
5767 recently_received_statements: HashMap::new(),
5768 network: network.clone(),
5769 sync: sync.clone(),
5770 sync_event_stream: (Box::pin(futures::stream::pending())
5771 as Pin<Box<dyn Stream<Item = sc_network_sync::types::SyncEvent> + Send>>)
5772 .fuse(),
5773 peers: HashMap::new(),
5774 statement_store: Arc::new(statement_store.clone()),
5775 queue_sender,
5776 statements_per_second: NonZeroU32::new(DEFAULT_STATEMENTS_PER_SECOND)
5777 .expect("DEFAULT_STATEMENTS_PER_SECOND is nonzero"),
5778 metrics: None,
5779 initial_sync_timeout: Box::pin(futures::future::pending()),
5780 pending_affinities_timeout: Box::pin(futures::future::pending()),
5781 pending_initial_syncs: HashMap::new(),
5782 initial_sync_peer_queue: VecDeque::new(),
5783 next_initial_sync_id: 0,
5784 initial_sync_in_flight_bytes: 0,
5785 propagation_outboxes: HashMap::new(),
5786 in_flight_chunks: HashMap::new(),
5787 next_chunk_id: 0,
5788 propagation_in_flight_bytes: 0,
5789 parked_propagations: VecDeque::new(),
5790 pending_sends: FuturesUnordered::new(),
5791 deferred_peers: HashSet::new(),
5792 dropped_statements_during_sync: false,
5793 sync_recovery_peer: None,
5794 sync_recovery_readd_timeout: Box::pin(futures::future::pending()),
5795 };
5796
5797 let mut stmt = new_live_statement();
5799 stmt.set_plain_data(b"during major sync".to_vec());
5800 statement_store.insert(stmt);
5801
5802 let peer_id = PeerId::random();
5804 handler.peers.insert(
5805 peer_id,
5806 Peer::new_for_testing(
5807 NonZeroU32::new(DEFAULT_STATEMENTS_PER_SECOND).unwrap(),
5808 NonZeroU32::new(
5809 DEFAULT_STATEMENTS_PER_SECOND * config::STATEMENTS_BURST_COEFFICIENT,
5810 )
5811 .unwrap(),
5812 ),
5813 );
5814
5815 handler.schedule_initial_sync_for_peer(peer_id);
5817
5818 assert!(
5819 handler.pending_initial_syncs.contains_key(&peer_id),
5820 "Initial sync should be queued even during major sync"
5821 );
5822 assert_eq!(handler.initial_sync_peer_queue.len(), 1);
5823
5824 handler.process_initial_sync_burst();
5826 handler.flush_pending_sends().await;
5827 assert!(
5828 handler.pending_initial_syncs.contains_key(&peer_id),
5829 "Pending sync should remain untouched during major sync"
5830 );
5831
5832 sync.major_syncing.store(false, Ordering::Relaxed);
5834 while handler.pending_initial_syncs.contains_key(&peer_id) {
5835 handler.process_initial_sync_burst();
5836 handler.flush_pending_sends().await;
5837 }
5838 assert!(
5839 handler.initial_sync_peer_queue.is_empty(),
5840 "Peer should have been processed after major sync ended"
5841 );
5842 }
5843
5844 #[tokio::test]
5845 async fn test_schedule_initial_sync_resends_all_matching() {
5846 let (mut handler, statement_store, _network, _notification_service) =
5847 build_handler_no_peers();
5848
5849 let peer_id = PeerId::random();
5850
5851 let mut stmt1 = new_live_statement();
5853 stmt1.set_plain_data(b"delivered before".to_vec());
5854 let mut stmt2 = new_live_statement();
5855 stmt2.set_plain_data(b"never delivered".to_vec());
5856
5857 statement_store.insert(stmt1);
5858 statement_store.insert(stmt2);
5859
5860 handler.peers.insert(
5861 peer_id,
5862 Peer {
5863 rate_limiter: PeerRateLimiter::new(
5864 NonZeroU32::new(DEFAULT_STATEMENTS_PER_SECOND).unwrap(),
5865 NonZeroU32::new(
5866 DEFAULT_STATEMENTS_PER_SECOND * config::STATEMENTS_BURST_COEFFICIENT,
5867 )
5868 .unwrap(),
5869 ),
5870 protocol_version: PeerProtocolVersion::V1,
5871 topic_affinity: None,
5872 is_light: false,
5873 pending_topic_affinity: None,
5874 sync_watermark: 0,
5875 },
5876 );
5877
5878 handler.schedule_initial_sync_for_peer(peer_id);
5879
5880 let pending = handler.pending_initial_syncs.get(&peer_id).unwrap();
5881 assert_eq!(pending.cursor, 0, "The sync must start from the oldest admission");
5883 assert_eq!(pending.watermark, 2, "Both admissions must sit below the watermark");
5884 }
5885
5886 #[tokio::test]
5887 async fn statement_below_the_watermark_reaches_a_syncing_peer_exactly_once() {
5888 let (mut handler, statement_store, _network, notification_service) =
5889 build_handler_no_peers();
5890 let peer_id = PeerId::random();
5891
5892 let mut statement = new_live_statement();
5893 statement.set_plain_data(b"pre-watermark".to_vec());
5894 let hash = statement.hash();
5895 statement_store.insert(statement.clone());
5896 statement_store.recent_statements.lock().unwrap().insert(hash, statement);
5898
5899 handler
5900 .handle_notification_event(NotificationEvent::NotificationStreamOpened {
5901 peer: peer_id,
5902 direction: sc_network::service::traits::Direction::Inbound,
5903 handshake: vec![],
5904 negotiated_fallback: Some(format!("/{STATEMENT_PROTOCOL_V1}").into()),
5905 })
5906 .await;
5907 assert!(handler.pending_initial_syncs.contains_key(&peer_id));
5908
5909 handler.propagate_statements().await;
5912 assert!(!handler.propagation_outboxes.contains_key(&peer_id));
5913
5914 while handler.pending_initial_syncs.contains_key(&peer_id) {
5915 handler.process_initial_sync_burst();
5916 handler.flush_pending_sends().await;
5917 }
5918
5919 let sent = get_peer_hashes(¬ification_service.get_sent_notifications(), peer_id);
5920 assert_eq!(sent, vec![hash], "the sync cursor is the only delivery path");
5921 }
5922
5923 #[tokio::test]
5924 async fn sync_watermark_keeps_filtering_propagation_after_the_sync_completes() {
5925 let (mut handler, statement_store, _network, notification_service) =
5926 build_handler_no_peers();
5927 let peer_id = PeerId::random();
5928
5929 let mut pre = new_live_statement();
5930 pre.set_plain_data(b"pre-watermark".to_vec());
5931 let pre_hash = pre.hash();
5932 statement_store.insert(pre.clone());
5933
5934 handler
5935 .handle_notification_event(NotificationEvent::NotificationStreamOpened {
5936 peer: peer_id,
5937 direction: sc_network::service::traits::Direction::Inbound,
5938 handshake: vec![],
5939 negotiated_fallback: Some(format!("/{STATEMENT_PROTOCOL_V1}").into()),
5940 })
5941 .await;
5942 while handler.pending_initial_syncs.contains_key(&peer_id) {
5943 handler.process_initial_sync_burst();
5944 handler.flush_pending_sends().await;
5945 }
5946 notification_service.clear_sent_notifications();
5947
5948 let mut fresh = new_live_statement();
5950 fresh.set_plain_data(b"post-watermark".to_vec());
5951 let fresh_hash = fresh.hash();
5952 statement_store.recent_statements.lock().unwrap().insert(pre_hash, pre);
5953 statement_store.recent_statements.lock().unwrap().insert(fresh_hash, fresh);
5954
5955 handler.propagate_statements().await;
5956 handler.flush_pending_sends().await;
5957
5958 let sent = get_peer_hashes(¬ification_service.get_sent_notifications(), peer_id);
5959 assert_eq!(
5960 sent,
5961 vec![fresh_hash],
5962 "only admissions at or above the watermark are propagated"
5963 );
5964 }
5965
5966 #[tokio::test]
5967 async fn test_malformed_v2_message_does_not_panic() {
5968 let (mut handler, _statement_store, _network, _notification_service) =
5969 build_handler_no_peers();
5970
5971 let peer_id = PeerId::random();
5972
5973 handler
5975 .handle_notification_event(NotificationEvent::NotificationStreamOpened {
5976 peer: peer_id,
5977 direction: sc_network::service::traits::Direction::Inbound,
5978 handshake: vec![],
5979 negotiated_fallback: None,
5980 })
5981 .await;
5982
5983 handler
5985 .handle_notification_event(NotificationEvent::NotificationReceived {
5986 peer: peer_id,
5987 notification: vec![0xFF, 0xFE, 0xFD].into(),
5988 })
5989 .await;
5990
5991 let mut stmt = new_live_statement();
5993 stmt.set_plain_data(b"v1 encoded".to_vec());
5994 let v1_encoded = vec![stmt].encode();
5995 handler
5996 .handle_notification_event(NotificationEvent::NotificationReceived {
5997 peer: peer_id,
5998 notification: v1_encoded.into(),
5999 })
6000 .await;
6001
6002 assert!(handler.peers.contains_key(&peer_id), "Peer should still be connected");
6004 }
6005
6006 fn oversized_statement_batch() -> Vec<u8> {
6009 let count = MAX_STATEMENTS_PER_NOTIFICATION + 1;
6010 let mut batch = Compact(count as u32).encode();
6011 batch.extend(iter::repeat_n(0u8, count));
6012 batch
6013 }
6014
6015 #[tokio::test]
6016 async fn test_v1_oversized_batch_reported_as_bad_message() {
6017 let (mut handler, _statement_store, network, _notification_service) =
6018 build_handler_no_peers();
6019
6020 let peer_id = PeerId::random();
6021
6022 handler
6023 .handle_notification_event(NotificationEvent::NotificationStreamOpened {
6024 peer: peer_id,
6025 direction: sc_network::service::traits::Direction::Inbound,
6026 handshake: vec![],
6027 negotiated_fallback: Some(format!("/{STATEMENT_PROTOCOL_V1}").into()),
6028 })
6029 .await;
6030
6031 handler
6032 .handle_notification_event(NotificationEvent::NotificationReceived {
6033 peer: peer_id,
6034 notification: oversized_statement_batch().into(),
6035 })
6036 .await;
6037
6038 let reports = network.get_reports();
6039 assert!(
6040 reports.iter().any(|(id, rep)| *id == peer_id && *rep == rep::BAD_MESSAGE),
6041 "Expected BAD_MESSAGE reputation change, but got: {reports:?}"
6042 );
6043 assert!(
6044 !network.get_disconnected_peers().contains(&peer_id),
6045 "Expected oversized-batch peer to stay connected"
6046 );
6047 }
6048
6049 #[tokio::test]
6050 async fn test_v2_oversized_batch_reported_as_bad_message() {
6051 let (mut handler, _statement_store, network, _notification_service) =
6052 build_handler_no_peers();
6053
6054 let peer_id = PeerId::random();
6055
6056 handler
6057 .handle_notification_event(NotificationEvent::NotificationStreamOpened {
6058 peer: peer_id,
6059 direction: sc_network::service::traits::Direction::Inbound,
6060 handshake: vec![],
6061 negotiated_fallback: None,
6062 })
6063 .await;
6064
6065 let mut notification = vec![STATEMENTS_VARIANT_INDEX];
6067 notification.extend(oversized_statement_batch());
6068
6069 handler
6070 .handle_notification_event(NotificationEvent::NotificationReceived {
6071 peer: peer_id,
6072 notification: notification.into(),
6073 })
6074 .await;
6075
6076 let reports = network.get_reports();
6077 assert!(
6078 reports.iter().any(|(id, rep)| *id == peer_id && *rep == rep::BAD_MESSAGE),
6079 "Expected BAD_MESSAGE reputation change, but got: {reports:?}"
6080 );
6081 assert!(
6082 !network.get_disconnected_peers().contains(&peer_id),
6083 "Expected oversized-batch peer to stay connected"
6084 );
6085 }
6086
6087 #[test]
6088 fn test_max_statement_payload_size_v2_overhead() {
6089 let v1_max = max_statement_payload_size(V1_ENVELOPE_OVERHEAD);
6090 let v2_max = max_statement_payload_size(V2_ENVELOPE_OVERHEAD);
6091
6092 assert!(
6094 v2_max < v1_max,
6095 "V2 payload capacity ({v2_max}) should be less than V1 ({v1_max})"
6096 );
6097 assert_eq!(v1_max - v2_max, 1, "V2 overhead is exactly 1 byte more than V1");
6098 }
6099
6100 #[tokio::test]
6101 async fn test_full_node_v2_gets_initial_sync_immediately() {
6102 let (mut handler, statement_store, _network, _notification_service) =
6103 build_handler_no_peers();
6104
6105 let mut stmt = new_live_statement();
6107 stmt.set_plain_data(b"full node v2".to_vec());
6108 statement_store.insert(stmt);
6109
6110 let peer_id = PeerId::random();
6111
6112 handler
6114 .handle_notification_event(NotificationEvent::NotificationStreamOpened {
6115 peer: peer_id,
6116 direction: sc_network::service::traits::Direction::Inbound,
6117 handshake: vec![],
6118 negotiated_fallback: None,
6119 })
6120 .await;
6121
6122 assert!(
6124 handler.pending_initial_syncs.contains_key(&peer_id),
6125 "Full-node V2 peer should have initial sync scheduled immediately"
6126 );
6127 assert_eq!(handler.peers.get(&peer_id).unwrap().protocol_version, PeerProtocolVersion::V2);
6128 assert!(!handler.peers.get(&peer_id).unwrap().is_light);
6129 }
6130
6131 #[tokio::test]
6132 async fn test_propagation_reaches_all_connected_peers() {
6133 let (
6134 mut handler,
6135 statement_store,
6136 _network,
6137 notification_service,
6138 _queue_receiver,
6139 peer_ids,
6140 ) = build_handler(5);
6141
6142 let mut expected_hashes = Vec::new();
6144 for i in 0..3u8 {
6145 let mut statement = new_live_statement();
6146 statement.set_plain_data(vec![i; 100]);
6147 let hash = statement.hash();
6148 expected_hashes.push(hash);
6149 statement_store.recent_statements.lock().unwrap().insert(hash, statement);
6150 }
6151 expected_hashes.sort();
6152
6153 handler.propagate_statements().await;
6154 handler.flush_pending_sends().await;
6155
6156 let sent = notification_service.get_sent_notifications();
6157
6158 for peer_id in &peer_ids {
6160 let mut received_hashes = get_peer_hashes(&sent, *peer_id);
6161 received_hashes.sort();
6162
6163 assert_eq!(
6164 received_hashes, expected_hashes,
6165 "Peer {peer_id} should have received all 3 statements"
6166 );
6167 }
6168
6169 assert!(statement_store.recent_statements.lock().unwrap().is_empty());
6171 }
6172
6173 #[tokio::test]
6174 async fn test_received_statement_filtering_per_peer() {
6175 let (
6176 mut handler,
6177 statement_store,
6178 _network,
6179 notification_service,
6180 _queue_receiver,
6181 peer_ids,
6182 ) = build_handler(3);
6183
6184 let peer_a = peer_ids[0];
6185 let peer_b = peer_ids[1];
6186 let peer_c = peer_ids[2];
6187
6188 let mut hashes = Vec::new();
6190 for i in 0..5u8 {
6191 let mut statement = new_live_statement();
6192 statement.set_plain_data(vec![i; 100]);
6193 let hash = statement.hash();
6194 hashes.push(hash);
6195 statement_store.recent_statements.lock().unwrap().insert(hash, statement);
6196 }
6197
6198 handler
6200 .recently_received_statements
6201 .insert(hashes[0], HashSet::from_iter([peer_a]));
6202 handler
6203 .recently_received_statements
6204 .insert(hashes[1], HashSet::from_iter([peer_a]));
6205 handler
6206 .recently_received_statements
6207 .insert(hashes[2], HashSet::from_iter([peer_b]));
6208
6209 handler.propagate_statements().await;
6210 handler.flush_pending_sends().await;
6211
6212 let sent = notification_service.get_sent_notifications();
6213
6214 let peer_a_hashes = get_peer_hashes(&sent, peer_a);
6215 let peer_b_hashes = get_peer_hashes(&sent, peer_b);
6216 let peer_c_hashes = get_peer_hashes(&sent, peer_c);
6217
6218 assert_eq!(peer_a_hashes.len(), 3, "peer_a should get 3 statements");
6220 assert!(!peer_a_hashes.contains(&hashes[0]), "peer_a sent s1");
6221 assert!(!peer_a_hashes.contains(&hashes[1]), "peer_a sent s2");
6222 assert!(peer_a_hashes.contains(&hashes[2]));
6223 assert!(peer_a_hashes.contains(&hashes[3]));
6224 assert!(peer_a_hashes.contains(&hashes[4]));
6225
6226 assert_eq!(peer_b_hashes.len(), 4, "peer_b should get 4 statements");
6228 assert!(!peer_b_hashes.contains(&hashes[2]), "peer_b sent s3");
6229 assert!(peer_b_hashes.contains(&hashes[0]));
6230 assert!(peer_b_hashes.contains(&hashes[1]));
6231 assert!(peer_b_hashes.contains(&hashes[3]));
6232 assert!(peer_b_hashes.contains(&hashes[4]));
6233
6234 let mut sorted_peer_c: Vec<_> = peer_c_hashes.into_iter().collect();
6236 sorted_peer_c.sort();
6237 let mut all_hashes = hashes.clone();
6238 all_hashes.sort();
6239 assert_eq!(sorted_peer_c, all_hashes, "peer_c should get all 5 statements");
6240 }
6241
6242 #[test]
6245 fn major_sync_defers_peers_and_handles_disconnect() {
6246 let (sync, _flag) = TestSync::with_syncing(true);
6247 let network = TestNetwork::new();
6248 let notification_service = TestNotificationService::new();
6249 let statement_store = TestStatementStore::new();
6250 let (queue_sender, _queue_receiver) = async_channel::bounded(100);
6251
6252 let mut handler = StatementHandler {
6253 protocol_name: "/statement/1".into(),
6254 notification_service: Box::new(notification_service),
6255 propagate_timeout: (Box::pin(futures::stream::pending())
6256 as Pin<Box<dyn Stream<Item = ()> + Send>>)
6257 .fuse(),
6258 pending_statements: FuturesUnordered::new(),
6259 pending_statements_peers: HashMap::new(),
6260 recently_received_statements: HashMap::new(),
6261 network: network.clone(),
6262 sync,
6263 sync_event_stream: (Box::pin(futures::stream::pending())
6264 as Pin<Box<dyn Stream<Item = sc_network_sync::types::SyncEvent> + Send>>)
6265 .fuse(),
6266 peers: HashMap::new(),
6267 statement_store: Arc::new(statement_store),
6268 queue_sender,
6269 statements_per_second: NonZeroU32::new(DEFAULT_STATEMENTS_PER_SECOND)
6270 .expect("DEFAULT_STATEMENTS_PER_SECOND is nonzero"),
6271 metrics: None,
6272 initial_sync_timeout: Box::pin(futures::future::pending()),
6273 pending_affinities_timeout: Box::pin(futures::future::pending()),
6274 pending_initial_syncs: HashMap::new(),
6275 initial_sync_peer_queue: VecDeque::new(),
6276 next_initial_sync_id: 0,
6277 initial_sync_in_flight_bytes: 0,
6278 propagation_outboxes: HashMap::new(),
6279 in_flight_chunks: HashMap::new(),
6280 next_chunk_id: 0,
6281 propagation_in_flight_bytes: 0,
6282 parked_propagations: VecDeque::new(),
6283 pending_sends: FuturesUnordered::new(),
6284 deferred_peers: HashSet::new(),
6285 dropped_statements_during_sync: false,
6286 sync_recovery_peer: None,
6287 sync_recovery_readd_timeout: Box::pin(pending().fuse()),
6288 };
6289
6290 let peer1 = PeerId::random();
6291 let peer2 = PeerId::random();
6292 let peer3 = PeerId::random();
6293
6294 handler.handle_sync_event(SyncEvent::PeerConnected {
6295 peer_id: peer1,
6296 roles: sc_network::Roles::FULL,
6297 });
6298 handler.handle_sync_event(SyncEvent::PeerConnected {
6299 peer_id: peer2,
6300 roles: sc_network::Roles::FULL,
6301 });
6302 handler.handle_sync_event(SyncEvent::PeerConnected {
6303 peer_id: peer3,
6304 roles: sc_network::Roles::FULL,
6305 });
6306
6307 assert!(network.get_added_reserved().is_empty());
6309 assert!(network.get_removed_reserved().is_empty());
6310 assert_eq!(handler.deferred_peers.len(), 3);
6311
6312 handler.handle_sync_event(SyncEvent::PeerDisconnected(peer1));
6314 assert_eq!(handler.deferred_peers.len(), 2);
6315 assert!(!handler.deferred_peers.contains(&peer1), "disconnected peer must leave buffer");
6316 assert!(handler.deferred_peers.contains(&peer2));
6317 assert!(handler.deferred_peers.contains(&peer3));
6318 assert!(network.get_removed_reserved().is_empty(), "no remove call for buffered peer");
6319 }
6320
6321 #[test]
6322 fn deferred_peers_flushed_on_sync_end_without_remove() {
6323 let (sync, flag) = TestSync::with_syncing(true);
6324 let network = TestNetwork::new();
6325 let notification_service = TestNotificationService::new();
6326 let statement_store = TestStatementStore::new();
6327 let (queue_sender, _queue_receiver) = async_channel::bounded(100);
6328
6329 let peer1 = PeerId::random();
6330 let peer2 = PeerId::random();
6331 let mut deferred = HashSet::new();
6332 deferred.insert(peer1);
6333 deferred.insert(peer2);
6334
6335 let mut handler = StatementHandler {
6336 protocol_name: "/statement/1".into(),
6337 notification_service: Box::new(notification_service),
6338 propagate_timeout: (Box::pin(futures::stream::pending())
6339 as Pin<Box<dyn Stream<Item = ()> + Send>>)
6340 .fuse(),
6341 pending_statements: FuturesUnordered::new(),
6342 pending_statements_peers: HashMap::new(),
6343 recently_received_statements: HashMap::new(),
6344 network: network.clone(),
6345 sync,
6346 sync_event_stream: (Box::pin(futures::stream::pending())
6347 as Pin<Box<dyn Stream<Item = sc_network_sync::types::SyncEvent> + Send>>)
6348 .fuse(),
6349 peers: HashMap::new(),
6350 statement_store: Arc::new(statement_store),
6351 queue_sender,
6352 statements_per_second: NonZeroU32::new(DEFAULT_STATEMENTS_PER_SECOND)
6353 .expect("DEFAULT_STATEMENTS_PER_SECOND is nonzero"),
6354 metrics: None,
6355 initial_sync_timeout: Box::pin(futures::future::pending()),
6356 pending_affinities_timeout: Box::pin(futures::future::pending()),
6357 pending_initial_syncs: HashMap::new(),
6358 initial_sync_peer_queue: VecDeque::new(),
6359 next_initial_sync_id: 0,
6360 initial_sync_in_flight_bytes: 0,
6361 propagation_outboxes: HashMap::new(),
6362 in_flight_chunks: HashMap::new(),
6363 next_chunk_id: 0,
6364 propagation_in_flight_bytes: 0,
6365 parked_propagations: VecDeque::new(),
6366 pending_sends: FuturesUnordered::new(),
6367 deferred_peers: deferred,
6368 dropped_statements_during_sync: false,
6369 sync_recovery_peer: None,
6370 sync_recovery_readd_timeout: Box::pin(pending().fuse()),
6371 };
6372
6373 flag.store(false, std::sync::atomic::Ordering::Relaxed);
6374 handler.drain_deferred_peers();
6375
6376 assert!(handler.deferred_peers.is_empty());
6377
6378 let added = network.get_added_reserved();
6379 assert_eq!(added.len(), 1);
6380 let added_addrs = &added[0];
6381 let expected_addr1: sc_network::Multiaddr =
6382 iter::once(multiaddr::Protocol::P2p(peer1.into())).collect();
6383 let expected_addr2: sc_network::Multiaddr =
6384 iter::once(multiaddr::Protocol::P2p(peer2.into())).collect();
6385 assert!(added_addrs.contains(&expected_addr1), "peer1 must be in added set");
6386 assert!(added_addrs.contains(&expected_addr2), "peer2 must be in added set");
6387
6388 assert!(network.get_removed_reserved().is_empty());
6389 }
6390
6391 #[tokio::test]
6392 async fn sync_recovery_schedules_remove_for_one_connected_peer() {
6393 let network = TestNetwork::new();
6394 let notification_service = TestNotificationService::new();
6395 let (sync, _flag) = TestSync::with_syncing(false);
6396 let (queue_sender, _) = async_channel::bounded(2);
6397 let statement_store = TestStatementStore::new();
6398
6399 let connected_peer = PeerId::random();
6400
6401 let mut peers = HashMap::new();
6402 peers.insert(
6403 connected_peer,
6404 Peer {
6405 rate_limiter: PeerRateLimiter::new(
6406 NonZeroU32::new(DEFAULT_STATEMENTS_PER_SECOND)
6407 .expect("DEFAULT_STATEMENTS_PER_SECOND is nonzero"),
6408 NonZeroU32::new(
6409 DEFAULT_STATEMENTS_PER_SECOND * config::STATEMENTS_BURST_COEFFICIENT,
6410 )
6411 .expect("burst capacity is nonzero"),
6412 ),
6413 protocol_version: PeerProtocolVersion::V1,
6414 topic_affinity: None,
6415 is_light: false,
6416 pending_topic_affinity: None,
6417 sync_watermark: 0,
6418 },
6419 );
6420
6421 let mut handler = StatementHandler {
6422 protocol_name: format!("/{STATEMENT_PROTOCOL_V1}").into(),
6423 notification_service: Box::new(notification_service),
6424 propagate_timeout: (Box::pin(futures::stream::pending())
6425 as Pin<Box<dyn Stream<Item = ()> + Send>>)
6426 .fuse(),
6427 pending_statements: FuturesUnordered::new(),
6428 pending_statements_peers: HashMap::new(),
6429 recently_received_statements: HashMap::new(),
6430 network: network.clone(),
6431 sync,
6432 sync_event_stream: (Box::pin(futures::stream::pending())
6433 as Pin<Box<dyn Stream<Item = sc_network_sync::types::SyncEvent> + Send>>)
6434 .fuse(),
6435 peers,
6436 statement_store: Arc::new(statement_store),
6437 queue_sender,
6438 statements_per_second: NonZeroU32::new(DEFAULT_STATEMENTS_PER_SECOND)
6439 .expect("DEFAULT_STATEMENTS_PER_SECOND is nonzero"),
6440 metrics: None,
6441 initial_sync_timeout: Box::pin(futures::future::pending()),
6442 pending_affinities_timeout: Box::pin(futures::future::pending()),
6443 pending_initial_syncs: HashMap::new(),
6444 initial_sync_peer_queue: VecDeque::new(),
6445 next_initial_sync_id: 0,
6446 initial_sync_in_flight_bytes: 0,
6447 propagation_outboxes: HashMap::new(),
6448 in_flight_chunks: HashMap::new(),
6449 next_chunk_id: 0,
6450 propagation_in_flight_bytes: 0,
6451 parked_propagations: VecDeque::new(),
6452 pending_sends: FuturesUnordered::new(),
6453 deferred_peers: HashSet::new(),
6454 dropped_statements_during_sync: true,
6455 sync_recovery_peer: None,
6456 sync_recovery_readd_timeout: Box::pin(futures::future::pending()),
6457 };
6458
6459 handler.start_sync_recovery();
6460
6461 {
6463 let removed = network.removed_reserved.lock().unwrap();
6464 assert_eq!(
6465 removed.len(),
6466 1,
6467 "Expected exactly one remove_peers_from_reserved_set call"
6468 );
6469 assert!(removed[0].contains(&connected_peer));
6470 }
6471
6472 assert_eq!(handler.sync_recovery_peer, Some(connected_peer));
6474
6475 handler.try_readd_sync_recovery_peer();
6478 assert!(handler.sync_recovery_peer.is_none());
6479 {
6480 let added = network.added_reserved.lock().unwrap();
6481 assert_eq!(added.len(), 1);
6482 let expected_addr: multiaddr::Multiaddr =
6483 iter::once(multiaddr::Protocol::P2p(connected_peer.into())).collect();
6484 assert!(added[0].contains(&expected_addr));
6485 }
6486
6487 {
6490 let peer2 = PeerId::random();
6491 handler.sync_recovery_peer = Some(peer2);
6492 handler.start_sync_recovery();
6493 assert_eq!(
6494 handler.sync_recovery_peer,
6495 Some(peer2),
6496 "Re-entry guard: recovery peer must not change on second call"
6497 );
6498 assert_eq!(
6499 network.removed_reserved.lock().unwrap().len(),
6500 1,
6501 "Re-entry guard: no extra remove call while recovery is in flight"
6502 );
6503 }
6504 }
6505
6506 #[tokio::test]
6507 async fn sync_recovery_gated_by_dropped_statements_flag() {
6508 let make_peer = || Peer {
6509 rate_limiter: PeerRateLimiter::new(
6510 NonZeroU32::new(DEFAULT_STATEMENTS_PER_SECOND)
6511 .expect("DEFAULT_STATEMENTS_PER_SECOND is nonzero"),
6512 NonZeroU32::new(
6513 DEFAULT_STATEMENTS_PER_SECOND * config::STATEMENTS_BURST_COEFFICIENT,
6514 )
6515 .expect("burst capacity is nonzero"),
6516 ),
6517 protocol_version: PeerProtocolVersion::V1,
6518 topic_affinity: None,
6519 is_light: false,
6520 pending_topic_affinity: None,
6521 sync_watermark: 0,
6522 };
6523
6524 let make_handler =
6525 |network: TestNetwork, dropped: bool| -> StatementHandler<TestNetwork, TestSync> {
6526 let (sync, _) = TestSync::with_syncing(false);
6527 let (queue_sender, _) = async_channel::bounded(2);
6528 let mut peers = HashMap::new();
6529 peers.insert(PeerId::random(), make_peer());
6530 StatementHandler {
6531 protocol_name: format!("/{STATEMENT_PROTOCOL_V1}").into(),
6532 notification_service: Box::new(TestNotificationService::new()),
6533 propagate_timeout: (Box::pin(futures::stream::pending())
6534 as Pin<Box<dyn Stream<Item = ()> + Send>>)
6535 .fuse(),
6536 pending_statements: FuturesUnordered::new(),
6537 pending_statements_peers: HashMap::new(),
6538 recently_received_statements: HashMap::new(),
6539 network,
6540 sync,
6541 sync_event_stream: (Box::pin(futures::stream::pending())
6542 as Pin<Box<dyn Stream<Item = sc_network_sync::types::SyncEvent> + Send>>)
6543 .fuse(),
6544 peers,
6545 statement_store: Arc::new(TestStatementStore::new()),
6546 queue_sender,
6547 statements_per_second: NonZeroU32::new(DEFAULT_STATEMENTS_PER_SECOND)
6548 .expect("DEFAULT_STATEMENTS_PER_SECOND is nonzero"),
6549 metrics: None,
6550 initial_sync_timeout: Box::pin(futures::future::pending()),
6551 pending_affinities_timeout: Box::pin(futures::future::pending()),
6552 pending_initial_syncs: HashMap::new(),
6553 initial_sync_peer_queue: VecDeque::new(),
6554 next_initial_sync_id: 0,
6555 initial_sync_in_flight_bytes: 0,
6556 propagation_outboxes: HashMap::new(),
6557 in_flight_chunks: HashMap::new(),
6558 next_chunk_id: 0,
6559 propagation_in_flight_bytes: 0,
6560 parked_propagations: VecDeque::new(),
6561 pending_sends: FuturesUnordered::new(),
6562 deferred_peers: HashSet::new(),
6563 dropped_statements_during_sync: dropped,
6564 sync_recovery_peer: None,
6565 sync_recovery_readd_timeout: Box::pin(pending().fuse()),
6566 }
6567 };
6568
6569 let net = TestNetwork::new();
6571 let mut handler = make_handler(net.clone(), false);
6572 handler.start_sync_recovery();
6573 assert!(handler.sync_recovery_peer.is_none());
6574 assert!(net.get_removed_reserved().is_empty());
6575
6576 let net2 = TestNetwork::new();
6578 let mut handler2 = make_handler(net2.clone(), true);
6579 handler2.start_sync_recovery();
6580 assert!(handler2.sync_recovery_peer.is_some());
6581 assert_eq!(net2.get_removed_reserved().len(), 1);
6582 }
6583
6584 #[test]
6585 fn send_paths_skip_expired_statements() {
6586 let mut live = new_live_statement();
6587 live.set_expiry_from_parts(u32::MAX, 0);
6588 live.set_plain_data(vec![1u8; 16]);
6589 let live_hash = live.hash();
6590
6591 let mut stale = Statement::new();
6592 stale.set_plain_data(vec![2u8; 16]);
6593 let stale_hash = stale.hash();
6594
6595 let store = TestStatementStore::new();
6596 store.insert(stale.clone());
6597 store.insert(live.clone());
6598
6599 let peer = Peer {
6600 rate_limiter: PeerRateLimiter::new(
6601 NonZeroU32::new(DEFAULT_STATEMENTS_PER_SECOND)
6602 .expect("DEFAULT_STATEMENTS_PER_SECOND is nonzero"),
6603 NonZeroU32::new(
6604 DEFAULT_STATEMENTS_PER_SECOND * config::STATEMENTS_BURST_COEFFICIENT,
6605 )
6606 .expect("burst capacity is nonzero"),
6607 ),
6608 protocol_version: PeerProtocolVersion::V1,
6609 topic_affinity: None,
6610 is_light: false,
6611 pending_topic_affinity: None,
6612 sync_watermark: 0,
6613 };
6614 let who = PeerId::random();
6615 let received = HashMap::new();
6616 let pending = HashMap::new();
6617 let max_size = max_statement_payload_size(V1_ENVELOPE_OVERHEAD);
6618
6619 let (statements, _processed, _size) = fetch_statement_chunk(
6620 &store,
6621 &received,
6622 &pending,
6623 &who,
6624 &peer,
6625 &[stale_hash, live_hash],
6626 max_size,
6627 )
6628 .expect("the test store never fails a fetch");
6629 assert_eq!(statements.iter().map(|(hash, _)| *hash).collect::<Vec<_>>(), vec![live_hash],);
6630
6631 let watermark = store.admission_watermark().expect("watermark is readable");
6632 let (batch, _size) =
6633 fetch_admitted_chunk(&store, &received, &pending, &who, &peer, 0, watermark, max_size)
6634 .expect("the test store never fails a walk");
6635 assert_eq!(
6636 batch.statements.iter().map(|(hash, _)| *hash).collect::<Vec<_>>(),
6637 vec![live_hash],
6638 );
6639 }
6640}